diff --git a/providers/amazon/src/airflow/providers/amazon/aws/log/cloudwatch_task_handler.py b/providers/amazon/src/airflow/providers/amazon/aws/log/cloudwatch_task_handler.py index 7825a54b4ccc7..d9d0b3f52f424 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/log/cloudwatch_task_handler.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/log/cloudwatch_task_handler.py @@ -25,7 +25,7 @@ import os import shutil from collections.abc import Generator -from datetime import date, datetime, timedelta, timezone +from datetime import UTC, date, datetime, timedelta from functools import cached_property from pathlib import Path from typing import TYPE_CHECKING, Any @@ -307,7 +307,7 @@ def _iter_events() -> Generator[CloudWatchLogEvent, None, None]: # instead of an empty view that looks like remote logging silently failed. if e.response.get("Error", {}).get("Code") != "ResourceNotFoundException": raise - notice_ts = end_time or datetime_to_epoch_utc_ms(datetime.now(tz=timezone.utc)) + notice_ts = end_time or datetime_to_epoch_utc_ms(datetime.now(tz=UTC)) yield { "timestamp": notice_ts, "ingestionTime": notice_ts, @@ -321,7 +321,7 @@ def _iter_events() -> Generator[CloudWatchLogEvent, None, None]: return _iter_events() def _parse_log_event_as_dumped_json(self, event: CloudWatchLogEvent) -> str: - event_dt = datetime.fromtimestamp(event["timestamp"] / 1000.0, tz=timezone.utc).isoformat() + event_dt = datetime.fromtimestamp(event["timestamp"] / 1000.0, tz=UTC).isoformat() event_msg = event["message"] try: message = json.loads(event_msg) @@ -432,7 +432,7 @@ def _read_remote_logs( return messages, logs def _event_to_str(self, event: CloudWatchLogEvent) -> str: - event_dt = datetime.fromtimestamp(event["timestamp"] / 1000.0, tz=timezone.utc) + event_dt = datetime.fromtimestamp(event["timestamp"] / 1000.0, tz=UTC) # Format a datetime object to a string in Zulu time without milliseconds. formatted_event_dt = event_dt.strftime("%Y-%m-%dT%H:%M:%SZ") message = event["message"] diff --git a/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py b/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py index 3b54c32eb16ca..db3e379fe6b25 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py @@ -1189,7 +1189,7 @@ def invoke_defer_method(self, last_log_time=None, context=None) -> None: connection_extras = conn.extra_dejson self.log.info("Successfully resolved connection extras for deferral.") - trigger_start_time = datetime.datetime.now(tz=datetime.timezone.utc) + trigger_start_time = datetime.datetime.now(tz=datetime.UTC) if self.pod is None or self.pod.metadata is None: raise RuntimeError("Pod must be created with metadata before deferring") diff --git a/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py b/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py index 63c5299839781..36b43290fef57 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py @@ -19,7 +19,7 @@ import logging import posixpath import re -from datetime import datetime, timezone +from datetime import UTC, datetime from enum import Enum from importlib import metadata from typing import TYPE_CHECKING, Any @@ -88,7 +88,7 @@ def datetime_to_epoch_ms(date_time: datetime) -> int: def datetime_to_epoch_utc_ms(date_time: datetime) -> int: """Convert a datetime object to an epoch integer (milliseconds) in UTC timezone.""" - return int(date_time.replace(tzinfo=timezone.utc).timestamp() * 1_000) + return int(date_time.replace(tzinfo=UTC).timestamp() * 1_000) def datetime_to_epoch_us(date_time: datetime) -> int: diff --git a/providers/amazon/src/airflow/providers/amazon/aws/utils/eks_get_token.py b/providers/amazon/src/airflow/providers/amazon/aws/utils/eks_get_token.py index 03e2f2513dfee..0dfe77f25615b 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/utils/eks_get_token.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/utils/eks_get_token.py @@ -19,7 +19,7 @@ import argparse import base64 import os -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta import boto3 from botocore.signers import RequestSigner @@ -31,7 +31,7 @@ def get_expiration_time(): - token_expiration = datetime.now(timezone.utc) + timedelta(minutes=TOKEN_EXPIRATION_MINUTES) + token_expiration = datetime.now(UTC) + timedelta(minutes=TOKEN_EXPIRATION_MINUTES) return token_expiration.strftime("%Y-%m-%dT%H:%M:%SZ") diff --git a/providers/amazon/src/airflow/providers/amazon/aws/utils/mixins.py b/providers/amazon/src/airflow/providers/amazon/aws/utils/mixins.py index 436266a83f558..d2c86e3e4e31b 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/utils/mixins.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/utils/mixins.py @@ -28,9 +28,7 @@ from __future__ import annotations from functools import cache, cached_property -from typing import Any, Generic, NamedTuple, TypeVar - -from typing_extensions import final +from typing import Any, Generic, NamedTuple, TypeVar, final from airflow.providers.amazon.aws.hooks.base_aws import AwsGenericHook diff --git a/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py b/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py index 7b4f8563cc70a..a2a7a79f3859e 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py @@ -22,7 +22,7 @@ import re import time from collections.abc import Generator -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from threading import Event, Thread from typing import TYPE_CHECKING @@ -109,7 +109,7 @@ def run(self) -> None: def _forward_log_events(self, continuation_token: AwsLogsHook.ContinuationToken) -> None: prev_timestamp_event = None for log_event in self._get_log_events(continuation_token): - current_timestamp_event = datetime.fromtimestamp(log_event["timestamp"] / 1000.0, tz=timezone.utc) + current_timestamp_event = datetime.fromtimestamp(log_event["timestamp"] / 1000.0, tz=UTC) if current_timestamp_event == prev_timestamp_event: # When multiple events have the same timestamp, somehow, only one event is logged # As a consequence, some logs are missed in the log group (in case they have the same @@ -146,7 +146,7 @@ def _get_log_events(self, skip_token: AwsLogsHook.ContinuationToken | None = Non @staticmethod def event_to_str(event: dict) -> str: - event_dt = datetime.fromtimestamp(event["timestamp"] / 1000.0, tz=timezone.utc) + event_dt = datetime.fromtimestamp(event["timestamp"] / 1000.0, tz=UTC) formatted_event_dt = event_dt.strftime("%Y-%m-%d %H:%M:%S,%f")[:-3] message = event["message"] return f"[{formatted_event_dt}] {message}" diff --git a/providers/amazon/tests/system/amazon/aws/tests/test_lambda_executor_dlq.py b/providers/amazon/tests/system/amazon/aws/tests/test_lambda_executor_dlq.py index 2c5016c66cc4d..65f0b7601144d 100644 --- a/providers/amazon/tests/system/amazon/aws/tests/test_lambda_executor_dlq.py +++ b/providers/amazon/tests/system/amazon/aws/tests/test_lambda_executor_dlq.py @@ -18,7 +18,7 @@ import os import time -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from urllib.parse import urlparse import boto3 @@ -69,7 +69,7 @@ def verify_dlq_activity(dlq_queue_name: str): cloudwatch = boto3.client("cloudwatch") # Try for up to 10 attempts (5 minutes total) for attempt in range(10): - end_time = datetime.now(timezone.utc) + end_time = datetime.now(UTC) start_time = end_time - timedelta(minutes=5) received_response = cloudwatch.get_metric_statistics( Namespace="AWS/SQS", diff --git a/providers/amazon/tests/unit/amazon/aws/executors/batch/test_batch_executor.py b/providers/amazon/tests/unit/amazon/aws/executors/batch/test_batch_executor.py index 05b5a15b00632..3e442850f7f8f 100644 --- a/providers/amazon/tests/unit/amazon/aws/executors/batch/test_batch_executor.py +++ b/providers/amazon/tests/unit/amazon/aws/executors/batch/test_batch_executor.py @@ -1010,7 +1010,7 @@ def test_try_adopt_task_instances(self, mock_executor): task.pool_slots = 1 task.priority_weight = 1 task.context_carrier = {} - task.queued_dttm = dt.datetime(2024, 1, 1, tzinfo=dt.timezone.utc) + task.queued_dttm = dt.datetime(2024, 1, 1, tzinfo=dt.UTC) task.dag_model = mock.Mock() task.dag_model.bundle_name = "test_bundle" task.dag_model.relative_fileloc = "test_dag.py" diff --git a/providers/amazon/tests/unit/amazon/aws/executors/utils/test_exponential_backoff_retry.py b/providers/amazon/tests/unit/amazon/aws/executors/utils/test_exponential_backoff_retry.py index ecb2d3356c497..fcee290e0afa8 100644 --- a/providers/amazon/tests/unit/amazon/aws/executors/utils/test_exponential_backoff_retry.py +++ b/providers/amazon/tests/unit/amazon/aws/executors/utils/test_exponential_backoff_retry.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from unittest import mock import pytest @@ -32,7 +32,7 @@ def test_exponential_backoff_retry_base_case(self, time_machine): time_machine.move_to(datetime(2023, 1, 1, 12, 0, 5)) mock_callable_function = mock.Mock() exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=0, callable_function=mock_callable_function, ) @@ -117,7 +117,7 @@ def test_exponential_backoff_retry_parameterized( mock_callable_function.side_effect = Exception() exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=attempt_number, callable_function=mock_callable_function, ) @@ -129,7 +129,7 @@ def test_exponential_backoff_retry_fail_success(self, time_machine, caplog): mock_callable_function.side_effect = [Exception(), True] time_machine.move_to(datetime(2023, 1, 1, 12, 0, 2)) exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=0, callable_function=mock_callable_function, ) @@ -139,7 +139,7 @@ def test_exponential_backoff_retry_fail_success(self, time_machine, caplog): time_machine.move_to(datetime(2023, 1, 1, 12, 0, 6)) exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=1, callable_function=mock_callable_function, ) @@ -151,7 +151,7 @@ def test_exponential_backoff_retry_max_delay(self, time_machine): mock_callable_function.return_value = Exception() time_machine.move_to(datetime(2023, 1, 1, 12, 4, 15)) exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=4, callable_function=mock_callable_function, max_delay=60 * 5, @@ -159,7 +159,7 @@ def test_exponential_backoff_retry_max_delay(self, time_machine): mock_callable_function.assert_not_called() # delay is 256 seconds; no calls made time_machine.move_to(datetime(2023, 1, 1, 12, 4, 16)) exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=4, callable_function=mock_callable_function, max_delay=60 * 5, @@ -168,7 +168,7 @@ def test_exponential_backoff_retry_max_delay(self, time_machine): time_machine.move_to(datetime(2023, 1, 1, 12, 5, 0)) exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=5, callable_function=mock_callable_function, max_delay=60 * 5, @@ -183,7 +183,7 @@ def test_exponential_backoff_retry_max_attempts(self, time_machine, caplog): time_machine.move_to(datetime(2023, 1, 1, 12, 55, 0)) for i in range(10): exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=i, callable_function=mock_callable_function, max_attempts=3, @@ -270,7 +270,7 @@ def test_exponential_backoff_retry_exponent_base_parameterized( time_machine.move_to(utcnow_value) exponential_backoff_retry( - last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc), + last_attempt_time=datetime(2023, 1, 1, 12, 0, 0, tzinfo=UTC), attempts_since_last_successful=attempt_number, callable_function=mock_callable_function, exponent_base=3, diff --git a/providers/amazon/tests/unit/amazon/aws/hooks/test_base_aws.py b/providers/amazon/tests/unit/amazon/aws/hooks/test_base_aws.py index d91cba0d3e75e..a7119b0ef646a 100644 --- a/providers/amazon/tests/unit/amazon/aws/hooks/test_base_aws.py +++ b/providers/amazon/tests/unit/amazon/aws/hooks/test_base_aws.py @@ -21,7 +21,7 @@ import os from base64 import b64encode from contextlib import nullcontext -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from pathlib import Path from unittest import mock from unittest.mock import MagicMock, PropertyMock, mock_open @@ -321,7 +321,7 @@ def side_effect(): "access_key": "mock-AccessKeyId", "secret_key": "mock-SecretAccessKey", "token": "mock-SessionToken", - "expiry_time": datetime.now(timezone.utc).isoformat(), + "expiry_time": datetime.now(UTC).isoformat(), } mock_refresh.side_effect = side_effect @@ -839,7 +839,7 @@ def test_refreshable_credentials(self, mock_get_connection): expire_on_calls = [] def mock_refresh_credentials(): - expiry_datetime = datetime.now(timezone.utc) + expiry_datetime = datetime.now(UTC) expire_on_call = expire_on_calls.pop() if expire_on_call: expiry_datetime -= timedelta(minutes=1000) 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 5f4da0a8dcc1e..2bc0a03bcfebf 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, timedelta, timezone as std_timezone +from datetime import UTC, datetime as std_datetime, timedelta from pathlib import Path from unittest import mock, mock as async_mock from unittest.mock import AsyncMock, MagicMock, Mock, patch @@ -1171,7 +1171,7 @@ async def test_s3_key_hook_is_keys_unchanged_exception_async(self, mock_list_key @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=std_timezone.utc) + frozen_dt = std_datetime(2026, 1, 1, 12, 0, 5, tzinfo=UTC) time_machine.move_to(frozen_dt, tick=False) mock_list_keys.return_value = ["test"] @@ -1293,7 +1293,7 @@ async def test_s3_key_hook_is_keys_unchanged_pending_async_with_tzinfo(self, moc previous_objects=set(), inactivity_seconds=0, allow_delete=False, - last_activity_time=std_datetime.now(std_timezone.utc), + last_activity_time=std_datetime.now(UTC), ) assert response.get("status") == "pending" diff --git a/providers/amazon/tests/unit/amazon/aws/hooks/test_sagemaker.py b/providers/amazon/tests/unit/amazon/aws/hooks/test_sagemaker.py index 1aa3a34d112a2..4b10fcff20dfe 100644 --- a/providers/amazon/tests/unit/amazon/aws/hooks/test_sagemaker.py +++ b/providers/amazon/tests/unit/amazon/aws/hooks/test_sagemaker.py @@ -18,7 +18,7 @@ from __future__ import annotations import time -from datetime import datetime, timezone +from datetime import UTC, datetime from unittest import mock from unittest.mock import patch @@ -528,7 +528,7 @@ def test_secondary_training_status_changed_false(self): def test_secondary_training_status_message_status_changed(self): now = datetime.now(tzlocal()) SECONDARY_STATUS_DESCRIPTION_1["LastModifiedTime"] = now - expected_time = now.astimezone(tz=timezone.utc).strftime("%Y-%m-%d %H:%M:%S") + expected_time = now.astimezone(tz=UTC).strftime("%Y-%m-%d %H:%M:%S") expected = f"{expected_time} {status} - {message}" assert ( secondary_training_status_message(SECONDARY_STATUS_DESCRIPTION_1, SECONDARY_STATUS_DESCRIPTION_2) diff --git a/providers/amazon/tests/unit/amazon/aws/log/test_cloudwatch_task_handler.py b/providers/amazon/tests/unit/amazon/aws/log/test_cloudwatch_task_handler.py index 2a1fc159049b3..c924658ea2ec0 100644 --- a/providers/amazon/tests/unit/amazon/aws/log/test_cloudwatch_task_handler.py +++ b/providers/amazon/tests/unit/amazon/aws/log/test_cloudwatch_task_handler.py @@ -22,7 +22,7 @@ import os import textwrap import time -from datetime import datetime as dt, timedelta, timezone as std_timezone +from datetime import UTC, datetime as dt, timedelta from pathlib import Path from unittest import mock from unittest.mock import ANY, call @@ -55,7 +55,7 @@ def get_time_str(time_in_milliseconds): - dt_time = dt.fromtimestamp(time_in_milliseconds / 1000.0, tz=std_timezone.utc) + dt_time = dt.fromtimestamp(time_in_milliseconds / 1000.0, tz=UTC) return dt_time.strftime("%Y-%m-%dT%H:%M:%SZ") diff --git a/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py b/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py index bc92c52fb1a43..9adb0859a4ee8 100644 --- a/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py +++ b/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py @@ -369,7 +369,7 @@ async def test_when_there_are_no_fargate_profiles_it_should_only_log_message(sel class TestEksPodTrigger: """Tests for EksPodTrigger.""" - TRIGGER_START_TIME = datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc) + TRIGGER_START_TIME = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) def _create_trigger(self, **overrides): """Create an EksPodTrigger with sensible defaults.""" diff --git a/providers/amazon/tests/unit/amazon/aws/triggers/test_kinesis.py b/providers/amazon/tests/unit/amazon/aws/triggers/test_kinesis.py index f42bdee7aac9c..57f2eeeeeaf33 100644 --- a/providers/amazon/tests/unit/amazon/aws/triggers/test_kinesis.py +++ b/providers/amazon/tests/unit/amazon/aws/triggers/test_kinesis.py @@ -129,9 +129,7 @@ def create_record(sequence_number="1", data=b"message", *, include_timestamp=Tru "Data": data, } if include_timestamp: - record["ApproximateArrivalTimestamp"] = datetime.datetime( - 2026, 8, 4, 12, 30, tzinfo=datetime.timezone.utc - ) + record["ApproximateArrivalTimestamp"] = datetime.datetime(2026, 8, 4, 12, 30, tzinfo=datetime.UTC) return record diff --git a/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py b/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py index 2eea945642b3d..3e7bb087919a8 100644 --- a/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py +++ b/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py @@ -35,7 +35,7 @@ is_resource_in_use_error, ) -DT = datetime.datetime(2000, 1, 1, tzinfo=datetime.timezone.utc) +DT = datetime.datetime(2000, 1, 1, tzinfo=datetime.UTC) EPOCH = 946_684_800 diff --git a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py index 6ab4976110336..0a3a917e13080 100644 --- a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py +++ b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py @@ -26,10 +26,9 @@ import time from collections.abc import Iterable, Mapping from tempfile import NamedTemporaryFile, TemporaryDirectory -from typing import TYPE_CHECKING, Any, Literal +from typing import TYPE_CHECKING, Any, Literal, overload from deprecated import deprecated -from typing_extensions import overload from airflow.exceptions import AirflowOptionalProviderFeatureException, AirflowProviderDeprecationWarning from airflow.providers.common.compat.sdk import ( diff --git a/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py b/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py index 20f6f1f6e449e..fdf1b2065bd5b 100644 --- a/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py +++ b/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py @@ -21,7 +21,7 @@ import logging import os import time -from datetime import datetime, timezone +from datetime import UTC, datetime from fnmatch import fnmatch from functools import lru_cache from typing import TYPE_CHECKING, Any @@ -275,7 +275,7 @@ def _on_delivery(err, _msg) -> None: def _now_iso() -> str: - return datetime.now(timezone.utc).isoformat() + return datetime.now(UTC).isoformat() # DagRun diff --git a/providers/apache/livy/src/airflow/providers/apache/livy/triggers/livy.py b/providers/apache/livy/src/airflow/providers/apache/livy/triggers/livy.py index 2d0cb99a64240..a7bfb90b384ce 100644 --- a/providers/apache/livy/src/airflow/providers/apache/livy/triggers/livy.py +++ b/providers/apache/livy/src/airflow/providers/apache/livy/triggers/livy.py @@ -18,7 +18,7 @@ import asyncio from collections.abc import AsyncIterator -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from typing import Any from airflow.providers.apache.livy.hooks.livy import BatchState, LivyAsyncHook @@ -121,7 +121,7 @@ async def poll_for_termination(self, batch_id: int | str) -> dict[str, Any]: :param batch_id: id of the batch session to monitor. """ if self._execution_timeout is not None: - timeout_datetime = datetime.now(timezone.utc) + self._execution_timeout + timeout_datetime = datetime.now(UTC) + self._execution_timeout else: timeout_datetime = None batch_execution_timed_out = False @@ -130,9 +130,7 @@ async def poll_for_termination(self, batch_id: int | str) -> dict[str, Any]: self.log.info("Batch with id %s is in state: %s", batch_id, state["batch_state"].value) while state["batch_state"] not in hook.TERMINAL_STATES: self.log.info("Batch with id %s is in state: %s", batch_id, state["batch_state"].value) - batch_execution_timed_out = ( - timeout_datetime is not None and datetime.now(timezone.utc) > timeout_datetime - ) + batch_execution_timed_out = timeout_datetime is not None and datetime.now(UTC) > timeout_datetime if batch_execution_timed_out: break self.log.info("Sleeping for %s seconds", self._polling_interval) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py index 5178222c647a7..a6ad43eb5f258 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py @@ -1053,7 +1053,7 @@ def invoke_defer_method( connection_extras = conn.extra_dejson self.log.info("Successfully resolved connection extras for deferral.") - trigger_start_time = datetime.datetime.now(tz=datetime.timezone.utc) + trigger_start_time = datetime.datetime.now(tz=datetime.UTC) # Translate ``execution_timeout`` into an absolute deadline plumbed to # the trigger via ``trigger_kwargs["_execution_deadline"]``. Anchoring @@ -1321,7 +1321,7 @@ def _write_logs(self, pod: k8s.V1Pod, follow: bool = False, since_time: DateTime if isinstance(since_time, str): # against interface spec but accept string as safeguard since_time = pendulum.parse(since_time.replace("Z", "+00:00")) since_seconds = math.ceil( - (datetime.datetime.now(tz=datetime.timezone.utc) - since_time).total_seconds() + (datetime.datetime.now(tz=datetime.UTC) - since_time).total_seconds() ) except (TypeError, ValueError): self.log.warning( @@ -1784,7 +1784,7 @@ def _get_most_recent_pod_index(self, pod_list: list[k8s.V1Pod]) -> int: pod_start_times: list[datetime.datetime] = [ # type: ignore[no-redef] pod.to_dict() .get("metadata", {}) - .get("creation_timestamp", datetime.datetime.now(tz=datetime.timezone.utc)) + .get("creation_timestamp", datetime.datetime.now(tz=datetime.UTC)) for pod in pod_list ] most_recent_start_time = max(pod_start_times) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/spark_kubernetes.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/spark_kubernetes.py index 72119977f661e..062f0418dbf6d 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/spark_kubernetes.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/spark_kubernetes.py @@ -17,7 +17,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from functools import cached_property from pathlib import Path from typing import TYPE_CHECKING, Any, cast @@ -273,7 +273,7 @@ def find_spark_job(self, context, exclude_checked: bool = True): p.metadata.deletion_timestamp is None, p.status.phase == PodPhase.SUCCEEDED, # If the job succeeded while the worker was down. p.status.phase == PodPhase.PENDING, - p.metadata.creation_timestamp or datetime.min.replace(tzinfo=timezone.utc), + p.metadata.creation_timestamp or datetime.min.replace(tzinfo=UTC), p.metadata.name or "", ), ) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/triggers/pod.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/triggers/pod.py index c72a8228ff2dc..1e5d514423af1 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/triggers/pod.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/triggers/pod.py @@ -318,7 +318,7 @@ async def _wait_for_pod_start_within_deadline(self) -> ContainerState: ) try: return await asyncio.wait_for(self._wait_for_pod_start(), timeout=remaining) - except asyncio.TimeoutError as exc: + except TimeoutError as exc: raise PodLaunchTimeoutException( f"Pod {self.pod_namespace}/{self.pod_name} reached the task's " "execution_timeout deadline while waiting for the pod to start." @@ -353,7 +353,7 @@ async def _wait_for_container_completion(self) -> TriggerEvent: Waits until container is no longer in running state. If trigger is configured with a logging period, then will emit an event to resume the task for the purpose of fetching more logs. """ - time_begin = datetime.datetime.now(tz=datetime.timezone.utc) + time_begin = datetime.datetime.now(tz=datetime.UTC) time_get_more_logs = None if self.logging_interval is not None: time_get_more_logs = time_begin + datetime.timedelta(seconds=self.logging_interval) @@ -404,7 +404,7 @@ async def _wait_for_container_completion(self) -> TriggerEvent: } ) self.log.debug("Container is not completed and still working.") - now = datetime.datetime.now(tz=datetime.timezone.utc) + now = datetime.datetime.now(tz=datetime.UTC) if time_get_more_logs and now >= time_get_more_logs: if self.get_logs and self.logging_interval: self.last_log_time = await self.pod_manager.fetch_container_logs_before_current_sec( diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py index ff6fd171fb7ce..5720bf68180be 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py @@ -3218,7 +3218,7 @@ def test_invoke_defer_method_passes_execution_deadline_when_execution_timeout_se k.pod.metadata.namespace = TEST_NAMESPACE ti_mock = MagicMock() - ti_start = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + ti_start = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) ti_mock.start_date = ti_start context = {"ti": ti_mock} @@ -3269,7 +3269,7 @@ def test_invoke_defer_method_pads_defer_timeout_for_slow_poll_interval( k.pod.metadata.namespace = TEST_NAMESPACE ti_mock = MagicMock() - ti_start = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + ti_start = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) ti_mock.start_date = ti_start context = {"ti": ti_mock} @@ -3312,7 +3312,7 @@ def test_invoke_defer_method_floors_defer_timeout_when_deadline_already_past( k.pod.metadata.namespace = TEST_NAMESPACE ti_mock = MagicMock() - ti_start = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + ti_start = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) ti_mock.start_date = ti_start context = {"ti": ti_mock} @@ -3353,7 +3353,7 @@ def test_invoke_defer_method_passes_no_deadline_when_execution_timeout_not_set( k.pod.metadata.namespace = TEST_NAMESPACE ti_mock = MagicMock() - ti_mock.start_date = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + ti_mock.start_date = datetime.datetime(2026, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) context = {"ti": ti_mock} with ( @@ -3667,9 +3667,9 @@ def test_write_logs_gives_up_after_max_retries( def test_write_logs_with_valid_since_time(self, mocked_client): """Test that since_seconds is calculated correctly when since_time is a valid datetime.""" pod = k8s.V1Pod(metadata=k8s.V1ObjectMeta(name=TEST_NAME, namespace=TEST_NAMESPACE)) - since_time = datetime.datetime( - 2026, 1, 1, 0, 0, 0, tzinfo=datetime.timezone.utc - ) - datetime.timedelta(seconds=30) + since_time = datetime.datetime(2026, 1, 1, 0, 0, 0, tzinfo=datetime.UTC) - datetime.timedelta( + seconds=30 + ) k = KubernetesPodOperator(task_id="task", get_logs=True) k._write_logs(pod, since_time=since_time) _, call_kwargs = mocked_client.read_namespaced_pod_log.call_args diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/triggers/test_pod.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/triggers/test_pod.py index b27ad42eba571..71905ed9cb898 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/triggers/test_pod.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/triggers/test_pod.py @@ -52,7 +52,7 @@ GET_LOGS = True STARTUP_TIMEOUT_SECS = 120 STARTUP_CHECK_INTERVAL_SECS = 0.1 -TRIGGER_START_TIME = datetime.datetime.now(tz=datetime.timezone.utc) +TRIGGER_START_TIME = datetime.datetime.now(tz=datetime.UTC) FAILED_RESULT_MSG = "Test message that appears when trigger have failed event." BASE_CONTAINER_NAME = "base" ON_FINISH_ACTION = "delete_pod" @@ -431,7 +431,7 @@ async def test_running_log_interval( """ If log interval given, check that the trigger fetches logs at the right times. """ - fixed_now = datetime.datetime(2022, 1, 1, tzinfo=datetime.timezone.utc) + fixed_now = datetime.datetime(2022, 1, 1, tzinfo=datetime.UTC) mock_datetime.datetime.now.side_effect = [ fixed_now, fixed_now + datetime.timedelta(seconds=1), diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py index 5473ca6c29156..be94e0dae2d8e 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py @@ -17,7 +17,7 @@ from __future__ import annotations import logging -from datetime import datetime, timedelta, timezone as dt_timezone +from datetime import UTC, datetime, timedelta from json.decoder import JSONDecodeError from typing import TYPE_CHECKING, cast from unittest import mock @@ -1132,7 +1132,7 @@ def test_fetch_requested_container_logs_invalid(self, container_running, contain def test_fetch_container_with_valid_since_time(self, logs_available, container_running): """Test that since_seconds is calculated correctly when since_time is a valid datetime.""" mock_pod = MagicMock() - since_time = datetime(2026, 1, 1, 0, 0, 0, tzinfo=dt_timezone.utc) - timedelta(seconds=30) + since_time = datetime(2026, 1, 1, 0, 0, 0, tzinfo=UTC) - timedelta(seconds=30) logs_available.return_value = True container_running.return_value = False self.mock_kube_client.read_namespaced_pod_log.return_value = mock.MagicMock( diff --git a/providers/common/ai/src/airflow/providers/common/ai/batch/base.py b/providers/common/ai/src/airflow/providers/common/ai/batch/base.py index e0d435fcbf77e..44ef109137d87 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/batch/base.py +++ b/providers/common/ai/src/airflow/providers/common/ai/batch/base.py @@ -28,9 +28,9 @@ from abc import ABC, abstractmethod from collections.abc import Iterator, Mapping from dataclasses import dataclass -from typing import TYPE_CHECKING, Any, ClassVar, Literal +from typing import TYPE_CHECKING, Any, ClassVar, Literal, Required -from typing_extensions import Required, TypedDict +from typing_extensions import TypedDict from airflow.providers.common.ai.exceptions import LLMBatchLimitExceededError, LLMBatchModelMismatchError diff --git a/providers/common/ai/src/airflow/providers/common/ai/batch/polling.py b/providers/common/ai/src/airflow/providers/common/ai/batch/polling.py index 233fbb1f60a11..9e0e6eb59ce6a 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/batch/polling.py +++ b/providers/common/ai/src/airflow/providers/common/ai/batch/polling.py @@ -28,7 +28,7 @@ from __future__ import annotations from dataclasses import dataclass -from datetime import datetime, timezone +from datetime import UTC, datetime from typing import Any from airflow.providers.common.ai.batch.base import IN_PROGRESS_STATUSES, TERMINAL_STATUS_MAP, BatchState @@ -75,7 +75,7 @@ def __init__(self, *, batch_id: str, end_time: float, timeout: int, cancel_on_ti @property def deadline_iso(self) -> str: - return datetime.fromtimestamp(self.end_time, tz=timezone.utc).isoformat(timespec="seconds") + return datetime.fromtimestamp(self.end_time, tz=UTC).isoformat(timespec="seconds") def on_state(self, state: BatchState, *, now: float) -> PollOutcome: """Decide after a successful status check.""" diff --git a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py index d9eb2eea76922..7fba7d7b6ac69 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py +++ b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_langchain_tool_agent.py @@ -371,7 +371,7 @@ def get_current_utc_time() -> str: merge freeze active right now", "how recent is the survey data"). LLMs cannot reliably know the wall-clock time on their own. """ - return datetime.datetime.now(datetime.timezone.utc).isoformat(timespec="seconds") + return datetime.datetime.now(datetime.UTC).isoformat(timespec="seconds") return [search_knowledge_base, query_survey_data, search_web, get_current_utc_time] diff --git a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_sandbox_toolset.py b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_sandbox_toolset.py index 456287c23aa89..9e9291f5a329c 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_sandbox_toolset.py +++ b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_sandbox_toolset.py @@ -41,7 +41,7 @@ from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from pydantic import BaseModel @@ -96,7 +96,7 @@ class ColumnMapping(BaseModel): @dag( schedule=None, - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["example", "sandbox"], ) @@ -158,7 +158,7 @@ def example_sandbox_agent_investigation(): # [START howto_sandbox_agent_local] @dag( schedule=None, - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["example", "sandbox"], ) @@ -217,7 +217,7 @@ def example_sandbox_agent_local(): @dag( schedule=None, - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["example", "sandbox"], params={ @@ -296,7 +296,7 @@ def convert(input_uri: str, output_uri: str) -> str: @dag( schedule=None, - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["example", "sandbox"], params={ @@ -389,7 +389,7 @@ def collect(sandbox: str | None, report_uri: str) -> str: @dag( schedule=None, - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["example", "sandbox"], ) diff --git a/providers/common/ai/src/airflow/providers/common/ai/operators/llm_batch.py b/providers/common/ai/src/airflow/providers/common/ai/operators/llm_batch.py index 6f05a216a8750..44b504b6361eb 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/operators/llm_batch.py +++ b/providers/common/ai/src/airflow/providers/common/ai/operators/llm_batch.py @@ -20,7 +20,7 @@ import time from collections.abc import Mapping, Sequence -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from typing import TYPE_CHECKING, Any, Literal from airflow.providers.common.ai.batch import dispatch, results, state @@ -331,13 +331,13 @@ def execute(self, context: Context) -> dict[str, Any]: batch_id=self.batch_id, key=key, spec=spec, - submitted_at=reattach_submitted_at or datetime.now(timezone.utc).isoformat(), + submitted_at=reattach_submitted_at or datetime.now(UTC).isoformat(), reattached=True, ) # Record the intent before the paid call so a crash between "request sent" and # "response recorded" leaves a trace the next attempt can act on. - intent_at = datetime.now(timezone.utc).isoformat() + intent_at = datetime.now(UTC).isoformat() state.write_intent( result_path_osp, key=key, @@ -356,7 +356,7 @@ def execute(self, context: Context) -> dict[str, Any]: request_params=self.request_params, completion_window=self.completion_window, ) - submitted_at = datetime.now(timezone.utc).isoformat() + submitted_at = datetime.now(UTC).isoformat() state.write_submitted( result_path_osp, key=key, @@ -469,7 +469,7 @@ def _resolve_reattach_target( "submit and recording the response); re-attaching instead of resubmitting.", recovered_batch_id, ) - submitted_at = record.intent_at or datetime.now(timezone.utc).isoformat() + submitted_at = record.intent_at or datetime.now(UTC).isoformat() state.write_submitted( result_path_osp, key=key, @@ -771,7 +771,7 @@ def _land( merge_diagnostics=merge_diagnostics, custom_id_prefix=key16, submitted_at=record.submitted_at or "", - completed_at=datetime.now(timezone.utc).isoformat(), + completed_at=datetime.now(UTC).isoformat(), ) def on_kill(self) -> None: diff --git a/providers/common/ai/src/airflow/providers/common/ai/toolsets/mcp.py b/providers/common/ai/src/airflow/providers/common/ai/toolsets/mcp.py index dfac26562c141..1b7cc80992a03 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/toolsets/mcp.py +++ b/providers/common/ai/src/airflow/providers/common/ai/toolsets/mcp.py @@ -18,9 +18,7 @@ from __future__ import annotations -from typing import TYPE_CHECKING, Any - -from typing_extensions import Self +from typing import TYPE_CHECKING, Any, Self from airflow.providers.common.ai.utils.toolset_base import AirflowToolset diff --git a/providers/common/ai/src/airflow/providers/common/ai/toolsets/object_storage.py b/providers/common/ai/src/airflow/providers/common/ai/toolsets/object_storage.py index 6e19b6fe4a07d..c8a971efd10ca 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/toolsets/object_storage.py +++ b/providers/common/ai/src/airflow/providers/common/ai/toolsets/object_storage.py @@ -21,7 +21,7 @@ import lzma import os import zlib -from datetime import datetime, timezone +from datetime import UTC, datetime from pathlib import PurePosixPath from typing import TYPE_CHECKING, Any, Literal @@ -379,7 +379,7 @@ def _as_iso(modified: Any) -> str | None: if isinstance(modified, datetime): return modified.isoformat() if isinstance(modified, (int, float)) and modified > 0: - return datetime.fromtimestamp(modified, tz=timezone.utc).isoformat() + return datetime.fromtimestamp(modified, tz=UTC).isoformat() return None diff --git a/providers/common/ai/src/airflow/providers/common/ai/toolsets/sandbox.py b/providers/common/ai/src/airflow/providers/common/ai/toolsets/sandbox.py index 75762aa1a6a12..2b4561c409816 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/toolsets/sandbox.py +++ b/providers/common/ai/src/airflow/providers/common/ai/toolsets/sandbox.py @@ -27,13 +27,12 @@ import threading import time import uuid -from typing import TYPE_CHECKING, Any, NamedTuple +from typing import TYPE_CHECKING, Any, NamedTuple, Self from fsspec.implementations.local import LocalFileSystem from pydantic_ai.exceptions import ModelRetry from pydantic_ai.tools import ToolDefinition from pydantic_ai.toolsets.abstract import AbstractToolset, ToolsetTool -from typing_extensions import Self from airflow.providers.common.ai.sandbox.base import ( AttachableSandboxBackend, diff --git a/providers/common/ai/src/airflow/providers/common/ai/utils/hitl_review.py b/providers/common/ai/src/airflow/providers/common/ai/utils/hitl_review.py index a97058b22c6f2..7dde07765c598 100644 --- a/providers/common/ai/src/airflow/providers/common/ai/utils/hitl_review.py +++ b/providers/common/ai/src/airflow/providers/common/ai/utils/hitl_review.py @@ -27,7 +27,7 @@ from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from enum import Enum from typing import Any, Literal @@ -85,7 +85,7 @@ class ConversationEntry(BaseModel): role: Literal["assistant", "human"] content: str iteration: int - timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) + timestamp: datetime = Field(default_factory=lambda: datetime.now(UTC)) class AgentSessionData(BaseModel): diff --git a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_modal.py b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_modal.py index 1453e5e9b7f5a..e3a20f158bcae 100644 --- a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_modal.py +++ b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_modal.py @@ -36,7 +36,7 @@ import os import time -from datetime import datetime, timezone +from datetime import UTC, datetime from airflow.providers.common.compat.sdk import dag as airflow_dag, task @@ -51,7 +51,7 @@ @airflow_dag( dag_id=DAG_ID, schedule="@once", - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["common.ai", "sandbox", "modal", "system_test"], ) diff --git a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py index 87d70c00f7506..bc6a33bcee82f 100644 --- a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py +++ b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py @@ -19,7 +19,7 @@ from __future__ import annotations import os -from datetime import datetime, timezone +from datetime import UTC, datetime from airflow.providers.common.compat.sdk import dag as airflow_dag, task @@ -35,7 +35,7 @@ @airflow_dag( dag_id=DAG_ID, schedule="@once", - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["common.ai", "sandbox", "opensandbox", "system_test"], ) diff --git a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_sbx.py b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_sbx.py index c7dee4bbffefc..aa196be9e17f6 100644 --- a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_sbx.py +++ b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_sbx.py @@ -19,7 +19,7 @@ from __future__ import annotations import os -from datetime import datetime, timezone +from datetime import UTC, datetime from airflow.providers.common.compat.sdk import dag as airflow_dag, task @@ -33,7 +33,7 @@ @airflow_dag( dag_id=DAG_ID, schedule="@once", - start_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + start_date=datetime(2024, 1, 1, tzinfo=UTC), catchup=False, tags=["common.ai", "sandbox", "sbx", "system_test"], ) diff --git a/providers/common/ai/tests/unit/common/ai/batch/test_output_schema.py b/providers/common/ai/tests/unit/common/ai/batch/test_output_schema.py index 2a4b982a8a5d0..9c2d7bceb5950 100644 --- a/providers/common/ai/tests/unit/common/ai/batch/test_output_schema.py +++ b/providers/common/ai/tests/unit/common/ai/batch/test_output_schema.py @@ -18,7 +18,7 @@ import json import threading -from datetime import datetime, timezone +from datetime import UTC, datetime from enum import Enum import pytest @@ -150,7 +150,7 @@ def test_datetime_and_enum_fields_serialize_with_mode_json(self): spec = build_output_spec(Event) extracted = ExtractedOutput( kind="json_value", - value={"happened_at": datetime(2026, 9, 11, tzinfo=timezone.utc), "severity": Severity.HIGH}, + value={"happened_at": datetime(2026, 9, 11, tzinfo=UTC), "severity": Severity.HIGH}, ) outcome = validate_extracted_output(extracted, spec) assert outcome.ok is True diff --git a/providers/common/ai/tests/unit/common/ai/durable/test_fingerprint.py b/providers/common/ai/tests/unit/common/ai/durable/test_fingerprint.py index d555336d9556c..3eeac1a5f078b 100644 --- a/providers/common/ai/tests/unit/common/ai/durable/test_fingerprint.py +++ b/providers/common/ai/tests/unit/common/ai/durable/test_fingerprint.py @@ -49,8 +49,8 @@ def make_messages(system: str = "You are a bot.", user: str = "hello", **part_kw class TestModelRequestFingerprint: def test_stable_across_part_timestamps(self): """Part timestamps regenerate on every attempt and must not affect the fingerprint.""" - t1 = datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc) - t2 = datetime.datetime(2026, 1, 2, tzinfo=datetime.timezone.utc) + t1 = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) + t2 = datetime.datetime(2026, 1, 2, tzinfo=datetime.UTC) fp1 = fingerprint_model_request("m", make_messages(timestamp=t1), None, ModelRequestParameters()) fp2 = fingerprint_model_request("m", make_messages(timestamp=t2), None, ModelRequestParameters()) diff --git a/providers/common/ai/tests/unit/common/ai/utils/test_hitl_review.py b/providers/common/ai/tests/unit/common/ai/utils/test_hitl_review.py index 36c749f163cf2..d0ecd7c651949 100644 --- a/providers/common/ai/tests/unit/common/ai/utils/test_hitl_review.py +++ b/providers/common/ai/tests/unit/common/ai/utils/test_hitl_review.py @@ -23,7 +23,7 @@ if not AIRFLOW_V_3_1_PLUS: pytest.skip("Human in the loop is only compatible with Airflow >= 3.1.0", allow_module_level=True) -from datetime import datetime, timezone +from datetime import UTC, datetime from pydantic import ValidationError @@ -61,9 +61,9 @@ def test_from_fields(self): assert isinstance(entry.timestamp, datetime) def test_timestamp_defaults_to_utc_now(self): - before = datetime.now(timezone.utc) + before = datetime.now(UTC) entry = ConversationEntry(role="human", content="Hi", iteration=2) - after = datetime.now(timezone.utc) + after = datetime.now(UTC) assert before <= entry.timestamp <= after diff --git a/providers/databricks/src/airflow/providers/databricks/utils/openlineage.py b/providers/databricks/src/airflow/providers/databricks/utils/openlineage.py index d241db12c9406..d3ae8f9cc36f3 100644 --- a/providers/databricks/src/airflow/providers/databricks/utils/openlineage.py +++ b/providers/databricks/src/airflow/providers/databricks/utils/openlineage.py @@ -116,7 +116,7 @@ def _process_data_from_api(data: list[dict[str, Any]]) -> list[dict[str, Any]]: """Convert timestamp fields to UTC datetime objects.""" for row in data: for key in ("query_start_time_ms", "query_end_time_ms"): - row[key] = datetime.datetime.fromtimestamp(row[key] / 1000, tz=datetime.timezone.utc) + row[key] = datetime.datetime.fromtimestamp(row[key] / 1000, tz=datetime.UTC) return data diff --git a/providers/databricks/tests/unit/databricks/utils/test_openlineage.py b/providers/databricks/tests/unit/databricks/utils/test_openlineage.py index 7517a16037044..e88e50febaa02 100644 --- a/providers/databricks/tests/unit/databricks/utils/test_openlineage.py +++ b/providers/databricks/tests/unit/databricks/utils/test_openlineage.py @@ -134,23 +134,15 @@ def test_process_data_from_api(): { "query_id": "ABC", "status": "FINISHED", - "query_start_time_ms": datetime.datetime( - 2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.timezone.utc - ), - "query_end_time_ms": datetime.datetime( - 2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.timezone.utc - ), + "query_start_time_ms": datetime.datetime(2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.UTC), + "query_end_time_ms": datetime.datetime(2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.UTC), "query_text": "SELECT * FROM table1;", "error_message": "Error occurred", }, { "query_id": "DEF", - "query_start_time_ms": datetime.datetime( - 2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.timezone.utc - ), - "query_end_time_ms": datetime.datetime( - 2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.timezone.utc - ), + "query_start_time_ms": datetime.datetime(2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.UTC), + "query_end_time_ms": datetime.datetime(2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.UTC), }, ] result = _process_data_from_api(data=data) @@ -201,8 +193,8 @@ def test_get_queries_details_from_databricks(mock_api_call): assert details == { "ABC": { "status": "FINISHED", - "start_time": datetime.datetime(2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.timezone.utc), - "end_time": datetime.datetime(2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.timezone.utc), + "start_time": datetime.datetime(2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.UTC), + "end_time": datetime.datetime(2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.UTC), "query_text": "SELECT * FROM table1;", "error_message": "Error occurred", } @@ -286,15 +278,15 @@ def test_emit_openlineage_events_for_databricks_queries(mock_generate_uuid, mock fake_metadata = { "query1": { "status": "FINISHED", - "start_time": datetime.datetime(2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.timezone.utc), - "end_time": datetime.datetime(2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.timezone.utc), + "start_time": datetime.datetime(2020, 7, 21, 18, 44, 46, 200000, tzinfo=datetime.UTC), + "end_time": datetime.datetime(2020, 7, 21, 18, 44, 47, 200000, tzinfo=datetime.UTC), "query_text": "SELECT * FROM table1", # No error for query1 }, "query2": { "status": "CANCELED", - "start_time": datetime.datetime(2020, 7, 21, 18, 44, 48, 200000, tzinfo=datetime.timezone.utc), - "end_time": datetime.datetime(2020, 7, 21, 18, 44, 49, 200000, tzinfo=datetime.timezone.utc), + "start_time": datetime.datetime(2020, 7, 21, 18, 44, 48, 200000, tzinfo=datetime.UTC), + "end_time": datetime.datetime(2020, 7, 21, 18, 44, 49, 200000, tzinfo=datetime.UTC), "query_text": "SELECT * FROM table2", "error_message": "Error occurred", }, diff --git a/providers/databricks/tests/unit/databricks/utils/test_retry.py b/providers/databricks/tests/unit/databricks/utils/test_retry.py index 7ff46a71f7d73..0ff54d483fcfe 100644 --- a/providers/databricks/tests/unit/databricks/utils/test_retry.py +++ b/providers/databricks/tests/unit/databricks/utils/test_retry.py @@ -49,7 +49,7 @@ def test_validate_deferrable_databricks_retry_args_accepts_serde_serializable_va def test_validate_deferrable_databricks_retry_args_accepts_airflow_serde_serializable_values(): - retry_args = {"deadline": datetime.datetime(2026, 5, 29, 12, 30, tzinfo=datetime.timezone.utc)} + retry_args = {"deadline": datetime.datetime(2026, 5, 29, 12, 30, tzinfo=datetime.UTC)} assert validate_deferrable_databricks_retry_args(retry_args, owner="test-owner") is None diff --git a/providers/discord/src/airflow/providers/discord/notifications/embed.py b/providers/discord/src/airflow/providers/discord/notifications/embed.py index 931f72f377ff4..29d903f8f1a10 100644 --- a/providers/discord/src/airflow/providers/discord/notifications/embed.py +++ b/providers/discord/src/airflow/providers/discord/notifications/embed.py @@ -24,9 +24,7 @@ from __future__ import annotations -from typing import Literal, TypedDict - -from typing_extensions import NotRequired, Required +from typing import Literal, NotRequired, Required, TypedDict EmbedType = Literal["rich"] diff --git a/providers/edge3/tests/unit/edge3/models/test_edge_job.py b/providers/edge3/tests/unit/edge3/models/test_edge_job.py index dc4b04135b98d..299a72c23d63c 100644 --- a/providers/edge3/tests/unit/edge3/models/test_edge_job.py +++ b/providers/edge3/tests/unit/edge3/models/test_edge_job.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone as dt_timezone +from datetime import UTC, datetime from typing import TYPE_CHECKING import pytest @@ -90,15 +90,15 @@ def test_build_job_key_keeps_task_key_unless_full_callback_identity(dag_id, run_ assert key == TaskInstanceKey(dag_id, "abc", run_id, try_number, map_index) -@time_machine.travel(datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc), tick=False) +@time_machine.travel(datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC), tick=False) def test_queued_dttm_defaults_to_now(): job = _make_job() - assert job.queued_dttm == datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc) + assert job.queued_dttm == datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC) def test_queued_dttm_explicit_value_is_kept(): - queued = datetime(2025, 6, 1, 8, 30, 0, tzinfo=dt_timezone.utc) + queued = datetime(2025, 6, 1, 8, 30, 0, tzinfo=UTC) job = _make_job(queued_dttm=queued) @@ -106,14 +106,14 @@ def test_queued_dttm_explicit_value_is_kept(): def test_last_update_t_returns_timestamp_of_last_update(): - last_update = datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc) + last_update = datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC) job = _make_job(last_update=last_update) assert job.last_update_t == last_update.timestamp() -@time_machine.travel(datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc), tick=False) +@time_machine.travel(datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC), tick=False) def test_last_update_t_falls_back_to_now_when_unset(): job = _make_job() @@ -128,7 +128,7 @@ def _clean_table(self, session: Session): session.commit() def test_round_trip(self, session: Session): - queued = datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc) + queued = datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC) session.add( _make_job( queued_dttm=queued, diff --git a/providers/edge3/tests/unit/edge3/models/test_edge_logs.py b/providers/edge3/tests/unit/edge3/models/test_edge_logs.py index fc0801b421f5d..9a59d08872e95 100644 --- a/providers/edge3/tests/unit/edge3/models/test_edge_logs.py +++ b/providers/edge3/tests/unit/edge3/models/test_edge_logs.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timedelta, timezone as dt_timezone +from datetime import UTC, datetime, timedelta from typing import TYPE_CHECKING import pytest @@ -27,7 +27,7 @@ if TYPE_CHECKING: from sqlalchemy.orm import Session -CHUNK_TIME = datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc) +CHUNK_TIME = datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC) def _make_log_chunk( diff --git a/providers/edge3/tests/unit/edge3/worker_api/test_datamodels.py b/providers/edge3/tests/unit/edge3/worker_api/test_datamodels.py index 09321b8970ddd..551138b26d7fb 100644 --- a/providers/edge3/tests/unit/edge3/worker_api/test_datamodels.py +++ b/providers/edge3/tests/unit/edge3/worker_api/test_datamodels.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone as dt_timezone +from datetime import UTC, datetime import pytest from pydantic import ValidationError @@ -132,11 +132,11 @@ def test_push_logs_body_parses_iso_timestamp(): log_chunk_data="log line", ) - assert body.log_chunk_time == datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc) + assert body.log_chunk_time == datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC) def test_worker_registration_return_assumes_version_mismatch(): - result = WorkerRegistrationReturn(last_update=datetime(2026, 1, 1, tzinfo=dt_timezone.utc)) + result = WorkerRegistrationReturn(last_update=datetime(2026, 1, 1, tzinfo=UTC)) assert result.versions_match is False diff --git a/providers/edge3/tests/unit/edge3/worker_api/test_datamodels_ui.py b/providers/edge3/tests/unit/edge3/worker_api/test_datamodels_ui.py index 017c34d9592b6..f55aacd657e77 100644 --- a/providers/edge3/tests/unit/edge3/worker_api/test_datamodels_ui.py +++ b/providers/edge3/tests/unit/edge3/worker_api/test_datamodels_ui.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone as dt_timezone +from datetime import UTC, datetime import pytest from pydantic import ValidationError @@ -73,7 +73,7 @@ def test_defaults(self): assert job.last_update is None def test_execution_fields_round_trip(self): - queued = datetime(2026, 1, 1, 12, 0, 0, tzinfo=dt_timezone.utc) + queued = datetime(2026, 1, 1, 12, 0, 0, tzinfo=UTC) job = _make_job(queued_dttm=queued, edge_worker="worker-1") diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/security_manager/override.py b/providers/fab/src/airflow/providers/fab/auth_manager/security_manager/override.py index 1d5ac262714be..bd138d43914d4 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/security_manager/override.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/security_manager/override.py @@ -1673,7 +1673,7 @@ def update_user(self, user: User) -> bool: new_role_ids = {r.id for r in user.roles} new_group_ids = {grp.id for grp in user.groups} if existing_role_ids != new_role_ids or existing_group_ids != new_group_ids: - user.changed_on = datetime.datetime.now(tz=datetime.timezone.utc) + user.changed_on = datetime.datetime.now(tz=datetime.UTC) merged_user = self.session.merge(user) self.session.commit() self._reset_user_permissions_cache(merged_user) diff --git a/providers/fab/tests/unit/fab/auth_manager/security_manager/test_fab_alignment.py b/providers/fab/tests/unit/fab/auth_manager/security_manager/test_fab_alignment.py index e3924bf177e82..37fda8802db4c 100644 --- a/providers/fab/tests/unit/fab/auth_manager/security_manager/test_fab_alignment.py +++ b/providers/fab/tests/unit/fab/auth_manager/security_manager/test_fab_alignment.py @@ -346,7 +346,7 @@ def test_update_user_sets_changed_on_when_roles_change(self, mock_log): assert result is True assert user.changed_on is not None assert isinstance(user.changed_on, datetime.datetime) - assert user.changed_on.tzinfo == datetime.timezone.utc + assert user.changed_on.tzinfo == datetime.UTC mock_session.merge.assert_called_once_with(user) mock_session.commit.assert_called_once() @@ -406,4 +406,4 @@ def test_update_user_sets_changed_on_when_groups_change(self, mock_log): assert result is True assert user.changed_on is not None - assert user.changed_on.tzinfo == datetime.timezone.utc + assert user.changed_on.tzinfo == datetime.UTC diff --git a/providers/fab/tests/unit/fab/www/extensions/test_init_session.py b/providers/fab/tests/unit/fab/www/extensions/test_init_session.py index 4784f286d9c93..0390bd04fcd36 100644 --- a/providers/fab/tests/unit/fab/www/extensions/test_init_session.py +++ b/providers/fab/tests/unit/fab/www/extensions/test_init_session.py @@ -35,7 +35,7 @@ from tests_common.test_utils.config import conf_vars -LOGIN_TIME = datetime.datetime(2026, 9, 9, 12, 0, tzinfo=datetime.timezone.utc) +LOGIN_TIME = datetime.datetime(2026, 9, 9, 12, 0, tzinfo=datetime.UTC) class FakeUser: diff --git a/providers/git/src/airflow/providers/git/hooks/git.py b/providers/git/src/airflow/providers/git/hooks/git.py index 010e889e4d1c1..7b8a4f3b6b000 100644 --- a/providers/git/src/airflow/providers/git/hooks/git.py +++ b/providers/git/src/airflow/providers/git/hooks/git.py @@ -27,7 +27,7 @@ import tempfile import warnings from collections.abc import Generator -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from typing import Any from urllib.parse import unquote @@ -263,7 +263,7 @@ def _ensure_github_app_token(self) -> None: TOKEN_REFRESH_BUFFER = timedelta(minutes=5) if ( self.github_app_token_exp is None - or self.github_app_token_exp < datetime.now(timezone.utc) + TOKEN_REFRESH_BUFFER + or self.github_app_token_exp < datetime.now(UTC) + TOKEN_REFRESH_BUFFER ): log.info( "GitHub App token is missing or near expiry (expires at: %s). Refreshing token.", diff --git a/providers/git/tests/unit/git/hooks/test_git.py b/providers/git/tests/unit/git/hooks/test_git.py index e014c8ae71bb8..b3ec043de8d60 100644 --- a/providers/git/tests/unit/git/hooks/test_git.py +++ b/providers/git/tests/unit/git/hooks/test_git.py @@ -28,6 +28,7 @@ import tempfile import threading import warnings +from datetime import UTC from unittest import mock import pytest @@ -846,9 +847,9 @@ def test_app_auth_with_key_file_reads_file(self, create_connection_without_db, t }, ) ) - from datetime import datetime, timedelta, timezone + from datetime import datetime, timedelta - mock_expiry = datetime.now(timezone.utc) + timedelta(hours=1) + mock_expiry = datetime.now(UTC) + timedelta(hours=1) monkeypatch.setattr( "airflow.providers.git.hooks.git.GitHook._get_github_app_token", lambda self: ("x-access-token", "ghs_test_token", mock_expiry), @@ -877,13 +878,13 @@ def test_app_auth_with_missing_key_file_raises(self, create_connection_without_d def test_app_auth_defers_token_fetch(self, monkeypatch): """GitHub App token is not fetched in __init__, only on configure_hook_env.""" - from datetime import datetime, timedelta, timezone + from datetime import datetime, timedelta mock_called = [] def mock_get_token(self): mock_called.append(True) - return ("x-access-token", "ghs_test_token", datetime.now(timezone.utc) + timedelta(hours=1)) + return ("x-access-token", "ghs_test_token", datetime.now(UTC) + timedelta(hours=1)) monkeypatch.setattr( "airflow.providers.git.hooks.git.GitHook._get_github_app_token", @@ -920,7 +921,7 @@ def test_app_auth_success_stores_app_id_and_installation_id(self): def test_app_id_and_installation_id_are_stored_as_provided( self, app_id, installation_id, create_connection_without_db, monkeypatch ): - from datetime import datetime, timedelta, timezone + from datetime import datetime, timedelta create_connection_without_db( Connection( @@ -936,7 +937,7 @@ def test_app_id_and_installation_id_are_stored_as_provided( ) monkeypatch.setattr( "airflow.providers.git.hooks.git.GitHook._get_github_app_token", - lambda self: ("x-access-token", "token", datetime.now(timezone.utc) + timedelta(hours=1)), + lambda self: ("x-access-token", "token", datetime.now(UTC) + timedelta(hours=1)), ) with pytest.warns(AirflowProviderDeprecationWarning, match="accept-new"): hook = GitHook(git_conn_id="git_app_int_check") @@ -945,14 +946,14 @@ def test_app_id_and_installation_id_are_stored_as_provided( def test_github_app_token_is_scoped_to_the_repository_host(self, monkeypatch): """The installation token goes through the same host-scoped helper as a connection token.""" - from datetime import datetime, timedelta, timezone + from datetime import datetime, timedelta monkeypatch.setattr( "airflow.providers.git.hooks.git.GitHook._get_github_app_token", lambda self: ( "x-access-token", "ghs_installation_token", - datetime.now(timezone.utc) + timedelta(hours=1), + datetime.now(UTC) + timedelta(hours=1), ), ) with pytest.warns(AirflowProviderDeprecationWarning, match="accept-new"): @@ -973,7 +974,7 @@ def test_github_app_token_is_scoped_to_the_repository_host(self, monkeypatch): def test_github_app_token_refresh_near_expiry(self, monkeypatch): """Token is refreshed when near expiry during configure_hook_env.""" - from datetime import datetime, timedelta, timezone + from datetime import datetime, timedelta mock_get_token_call_count = [0] @@ -984,13 +985,13 @@ def mock_get_token(self): return ( "x-access-token", f"token_{mock_get_token_call_count[0]}", - datetime.now(timezone.utc) + timedelta(minutes=3), + datetime.now(UTC) + timedelta(minutes=3), ) # Second call (refresh) returns token expiring in 1 hour return ( "x-access-token", f"token_{mock_get_token_call_count[0]}", - datetime.now(timezone.utc) + timedelta(hours=1), + datetime.now(UTC) + timedelta(hours=1), ) monkeypatch.setattr( @@ -1013,13 +1014,13 @@ def mock_get_token(self): def test_github_app_integration_call_shape(self, monkeypatch): """Verify GithubIntegration is called with correct arguments.""" - from datetime import datetime, timedelta, timezone + from datetime import datetime, timedelta from unittest import mock mock_integration = mock.MagicMock() mock_access_token = mock.MagicMock() mock_access_token.token = "ghs_test_token" - mock_access_token.expires_at = datetime.now(timezone.utc) + timedelta(hours=1) + mock_access_token.expires_at = datetime.now(UTC) + timedelta(hours=1) mock_integration.get_access_token.return_value = mock_access_token import sys diff --git a/providers/google/src/airflow/providers/google/cloud/hooks/managed_kafka.py b/providers/google/src/airflow/providers/google/cloud/hooks/managed_kafka.py index dafc4d56403ed..829cb35130b92 100644 --- a/providers/google/src/airflow/providers/google/cloud/hooks/managed_kafka.py +++ b/providers/google/src/airflow/providers/google/cloud/hooks/managed_kafka.py @@ -67,7 +67,7 @@ def _get_jwt(self, credentials): dict( exp=credentials.expiry.timestamp(), iss="Google", - iat=datetime.datetime.now(datetime.timezone.utc).timestamp(), + iat=datetime.datetime.now(datetime.UTC).timestamp(), scope="kafka", sub=credentials.service_account_email, ) @@ -88,8 +88,8 @@ def _get_kafka_access_token(self, credentials): def confluent_token(self): credentials = self._valid_credentials() - utc_expiry = credentials.expiry.replace(tzinfo=datetime.timezone.utc) - expiry_seconds = (utc_expiry - datetime.datetime.now(datetime.timezone.utc)).total_seconds() + utc_expiry = credentials.expiry.replace(tzinfo=datetime.UTC) + expiry_seconds = (utc_expiry - datetime.datetime.now(datetime.UTC)).total_seconds() return self._get_kafka_access_token(credentials), time.time() + expiry_seconds diff --git a/providers/google/src/airflow/providers/google/cloud/openlineage/mixins.py b/providers/google/src/airflow/providers/google/cloud/openlineage/mixins.py index 3e5520721ed97..5914b6da8db8a 100644 --- a/providers/google/src/airflow/providers/google/cloud/openlineage/mixins.py +++ b/providers/google/src/airflow/providers/google/cloud/openlineage/mixins.py @@ -21,7 +21,7 @@ import json import traceback from collections.abc import Iterable -from datetime import datetime, timezone +from datetime import UTC, datetime from typing import TYPE_CHECKING, cast from airflow.providers.common.compat.openlineage.facet import ( @@ -227,7 +227,7 @@ def _get_bigquery_job_datetime(properties: dict, field_name: str) -> datetime | if value is None: return None try: - return datetime.fromtimestamp(float(value) / 1000, tz=timezone.utc) + return datetime.fromtimestamp(float(value) / 1000, tz=UTC) except (TypeError, ValueError, OverflowError): return None @@ -239,7 +239,7 @@ def _get_child_job_sort_key(cls, properties: dict) -> tuple[datetime, str]: return ( # None is not comparable with datetime, so a missing startTime maps to # datetime.max to sort those children last instead of crashing the sort. - start_time or datetime.max.replace(tzinfo=timezone.utc), + start_time or datetime.max.replace(tzinfo=UTC), get_from_nullable_chain(properties, ["jobReference", "jobId"]) or "", ) diff --git a/providers/google/src/airflow/providers/google/cloud/operators/workflows.py b/providers/google/src/airflow/providers/google/cloud/operators/workflows.py index daf27caaaef5b..51477e641b1fd 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/workflows.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/workflows.py @@ -643,7 +643,7 @@ def __init__( self.workflow_id = workflow_id self.location = location self.start_date_filter = start_date_filter or datetime.datetime.now( - tz=datetime.timezone.utc + tz=datetime.UTC ) - datetime.timedelta(minutes=60) self.project_id = project_id self.retry = retry diff --git a/providers/google/src/airflow/providers/google/cloud/transfers/s3_to_gcs.py b/providers/google/src/airflow/providers/google/cloud/transfers/s3_to_gcs.py index 453026ddaa1e7..fc4a48b28358b 100644 --- a/providers/google/src/airflow/providers/google/cloud/transfers/s3_to_gcs.py +++ b/providers/google/src/airflow/providers/google/cloud/transfers/s3_to_gcs.py @@ -19,7 +19,7 @@ import warnings from collections.abc import Sequence -from datetime import datetime, timezone +from datetime import UTC, datetime from tempfile import NamedTemporaryFile from typing import TYPE_CHECKING, Any @@ -346,7 +346,7 @@ def transfer_files_async(self, files: list[str], gcs_hook: GCSHook, s3_hook: S3H ) def submit_transfer_jobs(self, files: list[str], gcs_hook: GCSHook, s3_hook: S3Hook) -> list[str]: - now = datetime.now(tz=timezone.utc) + now = datetime.now(tz=UTC) one_time_schedule = {"day": now.day, "month": now.month, "year": now.year} gcs_bucket, gcs_prefix = _parse_gcs_url(self.dest_gcs) diff --git a/providers/google/src/airflow/providers/google/common/hooks/base_google.py b/providers/google/src/airflow/providers/google/common/hooks/base_google.py index 649901971694b..ab54f1ad6ee24 100644 --- a/providers/google/src/airflow/providers/google/common/hooks/base_google.py +++ b/providers/google/src/airflow/providers/google/common/hooks/base_google.py @@ -884,7 +884,7 @@ def _now(): # On subsequent calls of `get` it will be used with `datetime.datetime.utcnow()`. # Therefore we have to use an offset-naive datetime. # https://github.com/talkiq/gcloud-aio/blob/f1132b005ba35d8059229a9ca88b90f31f77456d/auth/gcloud/aio/auth/token.py#L204 - return datetime.datetime.now(tz=datetime.timezone.utc).replace(tzinfo=None) + return datetime.datetime.now(tz=datetime.UTC).replace(tzinfo=None) class GoogleBaseAsyncHook(BaseHook): diff --git a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py index aa94eec813469..83947ec360288 100644 --- a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py +++ b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py @@ -23,7 +23,7 @@ import os from copy import deepcopy -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from airflow.models.dag import DAG from airflow.providers.amazon.aws.operators.s3 import S3CreateBucketOperator, S3DeleteBucketOperator @@ -97,7 +97,7 @@ SCHEDULE: { SCHEDULE_START_DATE: datetime(2015, 1, 1).date(), SCHEDULE_END_DATE: datetime(2030, 1, 1).date(), - START_TIME_OF_DAY: (datetime.now(tz=timezone.utc) + timedelta(minutes=1)).time(), + START_TIME_OF_DAY: (datetime.now(tz=UTC) + timedelta(minutes=1)).time(), }, TRANSFER_SPEC: { AWS_S3_DATA_SOURCE: {BUCKET_NAME: BUCKET_SOURCE_AWS}, diff --git a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py index 908b7ae63b71a..b9a761a74be2e 100644 --- a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py +++ b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py @@ -23,7 +23,7 @@ from __future__ import annotations import os -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from pathlib import Path from airflow.models.dag import DAG @@ -89,7 +89,7 @@ SCHEDULE: { SCHEDULE_START_DATE: datetime(2015, 1, 1).date(), SCHEDULE_END_DATE: datetime(2030, 1, 1).date(), - START_TIME_OF_DAY: (datetime.now(tz=timezone.utc) + timedelta(seconds=120)).time(), + START_TIME_OF_DAY: (datetime.now(tz=UTC) + timedelta(seconds=120)).time(), }, TRANSFER_SPEC: { GCS_DATA_SOURCE: {BUCKET_NAME: BUCKET_NAME_SRC}, diff --git a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/commands/cmd_delete.py b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/commands/cmd_delete.py index 65b6cc931afbc..cda559d0858c4 100644 --- a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/commands/cmd_delete.py +++ b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/commands/cmd_delete.py @@ -41,8 +41,8 @@ def _parse_datetime(value: str) -> datetime.datetime: dt = datetime.datetime.fromisoformat(value.replace("Z", "+00:00")) if dt.tzinfo is None: - dt = dt.replace(tzinfo=datetime.timezone.utc) - return dt.astimezone(datetime.timezone.utc) + dt = dt.replace(tzinfo=datetime.UTC) + return dt.astimezone(datetime.UTC) def _parse_create_time(value: Any) -> datetime.datetime | None: @@ -88,7 +88,7 @@ def check_min_age( ) return False - now = now or datetime.datetime.now(datetime.timezone.utc) + now = now or datetime.datetime.now(datetime.UTC) age = now - created_at if age < datetime.timedelta(days=min_age_days): print( diff --git a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/composer.py b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/composer.py index 703d0c192219c..989cbf68e2bd0 100644 --- a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/composer.py +++ b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/composer.py @@ -29,7 +29,7 @@ async def _delete_composer_environment(resource: dict, log_prefix: str): name = get_resource_path(resource) create_time = datetime.datetime.fromisoformat(resource["createTime"].replace("Z", "+00:00")) - age = datetime.datetime.now(datetime.timezone.utc) - create_time + age = datetime.datetime.now(datetime.UTC) - create_time if not age > datetime.timedelta(days=DAYS_PROTECTED): print( f"Composer env with name: {name} was skipped because it is protected for {DAYS_PROTECTED} days." diff --git a/providers/google/tests/unit/google/cloud/openlineage/test_mixins.py b/providers/google/tests/unit/google/cloud/openlineage/test_mixins.py index 0106eb81eb3c8..2a44c648ef973 100644 --- a/providers/google/tests/unit/google/cloud/openlineage/test_mixins.py +++ b/providers/google/tests/unit/google/cloud/openlineage/test_mixins.py @@ -20,7 +20,7 @@ import json import logging import os -from datetime import datetime, timezone +from datetime import UTC, datetime from unittest.mock import MagicMock, patch import pytest @@ -113,7 +113,7 @@ def read_common_json_file(rel: str): def make_task_instance(): - logical_date = datetime(2024, 1, 1, tzinfo=timezone.utc) + logical_date = datetime(2024, 1, 1, tzinfo=UTC) dag_run = MagicMock( logical_date=logical_date, clear_number=0, @@ -497,10 +497,10 @@ def test_get_openlineage_facets_on_complete_script_job(self, mock_emit_query_lin assert mock_emit_query_lineage.call_args.kwargs["task_instance"] is mock_ti assert mock_emit_query_lineage.call_args.kwargs["job_name"] == "dag_id.task_id.query.1" assert mock_emit_query_lineage.call_args.kwargs["start_time"] == datetime.fromtimestamp( - self.query_job_details["statistics"]["startTime"] / 1000, tz=timezone.utc + self.query_job_details["statistics"]["startTime"] / 1000, tz=UTC ) assert mock_emit_query_lineage.call_args.kwargs["end_time"] == datetime.fromtimestamp( - self.query_job_details["statistics"]["endTime"] / 1000, tz=timezone.utc + self.query_job_details["statistics"]["endTime"] / 1000, tz=UTC ) @patch.object( @@ -830,7 +830,7 @@ def test_child_query_lineage_skipped_with_old_openlineage_provider(self): @pytest.mark.parametrize( ("value", "expected"), [ - ("1600000000000", datetime(2020, 9, 13, 12, 26, 40, tzinfo=timezone.utc)), + ("1600000000000", datetime(2020, 9, 13, 12, 26, 40, tzinfo=UTC)), (None, None), ("not-a-timestamp", None), ], diff --git a/providers/google/tests/unit/google/cloud/operators/test_gcs.py b/providers/google/tests/unit/google/cloud/operators/test_gcs.py index fcd17462ead94..63d8dcc1fd34a 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_gcs.py +++ b/providers/google/tests/unit/google/cloud/operators/test_gcs.py @@ -18,7 +18,7 @@ from __future__ import annotations from concurrent.futures import ThreadPoolExecutor -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from pathlib import Path from unittest import mock @@ -484,7 +484,7 @@ def test_get_openlineage_facets_on_start_destination_falls_back_to_rendered_sour class TestGCSTimeSpanFileTransformOperatorDateInterpolation: def test_execute(self): - interp_dt = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=timezone.utc) + interp_dt = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=UTC) assert GCSTimeSpanFileTransformOperator.interpolate_prefix(None, interp_dt) is None @@ -551,7 +551,7 @@ def test_execute(self, mock_hook, mock_subprocess, mock_tempdir): file1 = "file1" file2 = "file2" - timespan_start = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) mock_ti = mock.Mock() @@ -754,7 +754,7 @@ def test_get_openlineage_facets_on_complete( file1 = "file1" file2 = "file2" - timespan_start = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) context = dict( @@ -824,7 +824,7 @@ def test_get_openlineage_facets_on_complete( def test_parallel_download_worker_behavior( self, mock_hook, mock_subprocess, mock_tempdir, mock_executor, workers, should_raise ): - timespan_start = datetime(2015, 2, 1, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) context = { @@ -892,7 +892,7 @@ def test_parallel_download_worker_behavior( def test_parallel_download_failure_behavior( self, mock_hook, mock_subprocess, mock_tempdir, mock_executor, mock_as_completed, continue_on_fail ): - timespan_start = datetime(2015, 2, 1, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) context = { @@ -973,7 +973,7 @@ def test_parallel_download_failure_behavior( def test_parallel_upload_worker_behavior( self, mock_hook, mock_subprocess, mock_tempdir, mock_executor, workers, should_raise ): - timespan_start = datetime(2015, 2, 1, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) context = { @@ -1050,7 +1050,7 @@ def test_parallel_upload_worker_behavior( def test_parallel_upload_failure_behavior( self, mock_hook, mock_subprocess, mock_tempdir, mock_executor, mock_as_completed, continue_on_fail ): - timespan_start = datetime(2015, 2, 1, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) context = { @@ -1133,7 +1133,7 @@ def test_execute_rejects_path_traversal_in_blob_name(self, mock_hook, mock_subpr directory (CWE-22) — exploitable when the source bucket is shared with untrusted writers. """ - timespan_start = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, 15, 16, 17, 345, tzinfo=UTC) timespan_end = timespan_start + timedelta(hours=1) context = dict( logical_date=timespan_start, @@ -1184,7 +1184,7 @@ def test_download_honors_download_num_attempts( succeeds, expected_sleeps, ): - timespan_start = datetime(2015, 2, 1, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, tzinfo=UTC) context = { "logical_date": timespan_start, "data_interval_start": timespan_start, @@ -1246,7 +1246,7 @@ def test_upload_honors_upload_num_attempts( succeeds, expected_sleeps, ): - timespan_start = datetime(2015, 2, 1, tzinfo=timezone.utc) + timespan_start = datetime(2015, 2, 1, tzinfo=UTC) context = { "logical_date": timespan_start, "data_interval_start": timespan_start, diff --git a/providers/google/tests/unit/google/cloud/operators/test_workflows.py b/providers/google/tests/unit/google/cloud/operators/test_workflows.py index 6afd7ae245b6a..32435fd9db4c6 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_workflows.py +++ b/providers/google/tests/unit/google/cloud/operators/test_workflows.py @@ -203,9 +203,7 @@ class TestWorkflowsListWorkflowsOperator: @mock.patch(BASE_PATH.format("WorkflowsHook")) def test_execute(self, mock_hook, mock_object): timestamp = Timestamp() - timestamp.FromDatetime( - datetime.datetime.now(tz=datetime.timezone.utc) + datetime.timedelta(minutes=5) - ) + timestamp.FromDatetime(datetime.datetime.now(tz=datetime.UTC) + datetime.timedelta(minutes=5)) workflow_mock = mock.MagicMock() workflow_mock.start_time = timestamp mock_hook.return_value.list_workflows.return_value = [workflow_mock] @@ -364,7 +362,7 @@ class TestWorkflowExecutionsListExecutionsOperator: @mock.patch(BASE_PATH.format("Execution")) @mock.patch(BASE_PATH.format("WorkflowsHook")) def test_execute(self, mock_hook, mock_object): - start_date_filter = datetime.datetime.now(tz=datetime.timezone.utc) + datetime.timedelta(minutes=5) + start_date_filter = datetime.datetime.now(tz=datetime.UTC) + datetime.timedelta(minutes=5) execution_mock = mock.MagicMock() execution_mock.start_time = start_date_filter mock_hook.return_value.list_executions.return_value = [execution_mock] diff --git a/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py b/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py index aa0c21ee7e67e..ccb83b7fd211e 100644 --- a/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py +++ b/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py @@ -116,7 +116,7 @@ def _assert_uploaded_file_content(self, bucket, obj, tmp_filename, mime_type, gz (datetime.date(1000, 1, 2), "1000-01-02"), (datetime.datetime(1970, 1, 1, 1, 0, tzinfo=None), "1970-01-01T01:00:00"), ( - datetime.datetime(2022, 1, 1, 2, 0, tzinfo=datetime.timezone.utc), + datetime.datetime(2022, 1, 1, 2, 0, tzinfo=datetime.UTC), 1641002400.0, ), (datetime.time(hour=0, minute=0, second=0), "0:00:00"), diff --git a/providers/google/tests/unit/google/cloud/triggers/test_gcs.py b/providers/google/tests/unit/google/cloud/triggers/test_gcs.py index 27c0d3b16c347..4ae9cf78a3697 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_gcs.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_gcs.py @@ -18,7 +18,7 @@ from __future__ import annotations import asyncio -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from typing import Any from unittest import mock from unittest.mock import AsyncMock @@ -41,7 +41,7 @@ TEST_GCP_CONN_ID = "TEST_GCP_CONN_ID" TEST_POLLING_INTERVAL = 3.0 TEST_HOOK_PARAMS: dict[str, Any] = {} -TEST_TS_OBJECT = datetime.now(tz=timezone.utc) +TEST_TS_OBJECT = datetime.now(tz=UTC) TEST_INACTIVITY_PERIOD = 5.0 diff --git a/providers/google/tests/unit/google/cloud/triggers/test_kubernetes_engine.py b/providers/google/tests/unit/google/cloud/triggers/test_kubernetes_engine.py index 34cb30740a90c..05ebb8e30a806 100644 --- a/providers/google/tests/unit/google/cloud/triggers/test_kubernetes_engine.py +++ b/providers/google/tests/unit/google/cloud/triggers/test_kubernetes_engine.py @@ -56,7 +56,7 @@ GET_LOGS = True STARTUP_TIMEOUT_SECS = 120 SCHEDULE_TIMEOUT_SECS = 60 -TRIGGER_START_TIME = datetime.datetime.now(tz=datetime.timezone.utc) +TRIGGER_START_TIME = datetime.datetime.now(tz=datetime.UTC) CLUSTER_URL = "https://test-host" SSL_CA_CERT = "TEST_SSL_CA_CERT_CONTENT" FAILED_RESULT_MSG = "Test message that appears when trigger have failed event." diff --git a/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py b/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py index bc34041d9d459..ed6c337c0685f 100644 --- a/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py +++ b/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py @@ -191,7 +191,7 @@ def test_asset_type_patterns_are_unique(): def test_check_min_age(): - now = datetime.datetime(2026, 5, 20, tzinfo=datetime.timezone.utc) + now = datetime.datetime(2026, 5, 20, tzinfo=datetime.UTC) assert check_min_age({"name": "res1"}, None, now=now) assert not check_min_age({"name": "res1"}, 3, now=now) @@ -201,7 +201,7 @@ def test_check_min_age(): async def test_handle_asset_type_filters_by_min_age(): asset_type = "ai" - now = datetime.datetime.now(datetime.timezone.utc) + now = datetime.datetime.now(datetime.UTC) resources = [ {"name": "old", "createTime": (now - datetime.timedelta(days=4)).isoformat()}, {"name": "new", "createTime": now.isoformat()}, diff --git a/providers/microsoft/azure/tests/unit/microsoft/azure/hooks/test_msgraph.py b/providers/microsoft/azure/tests/unit/microsoft/azure/hooks/test_msgraph.py index f3998f18c620d..4aa0fea5acadc 100644 --- a/providers/microsoft/azure/tests/unit/microsoft/azure/hooks/test_msgraph.py +++ b/providers/microsoft/azure/tests/unit/microsoft/azure/hooks/test_msgraph.py @@ -641,8 +641,8 @@ def test_get_credentials_returns_async_certificate_credential(self): .issuer_name(name) .public_key(private_key.public_key()) .serial_number(x509.random_serial_number()) - .not_valid_before(datetime.datetime.now(datetime.timezone.utc)) - .not_valid_after(datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(days=1)) + .not_valid_before(datetime.datetime.now(datetime.UTC)) + .not_valid_after(datetime.datetime.now(datetime.UTC) + datetime.timedelta(days=1)) .sign(private_key, hashes.SHA256()) ) pem = private_key.private_bytes( diff --git a/providers/openlineage/src/airflow/providers/openlineage/api/datasets.py b/providers/openlineage/src/airflow/providers/openlineage/api/datasets.py index e8e9f5b68ccc7..d2c65043b8995 100644 --- a/providers/openlineage/src/airflow/providers/openlineage/api/datasets.py +++ b/providers/openlineage/src/airflow/providers/openlineage/api/datasets.py @@ -19,7 +19,7 @@ from __future__ import annotations import logging -from datetime import datetime, timezone +from datetime import UTC, datetime from typing import TYPE_CHECKING from openlineage.client.event_v2 import Dataset, Job, Run, RunEvent, RunState @@ -169,7 +169,7 @@ def my_task(): event = RunEvent( eventType=RunState.RUNNING, - eventTime=datetime.now(tz=timezone.utc).isoformat(), + eventTime=datetime.now(tz=UTC).isoformat(), run=Run(runId=task_uuid, facets=run_facets), job=Job( namespace=lineage_job_namespace(), diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_manual_lineage_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_manual_lineage_dag.py index f754dcdfb87d4..7718e99ebf209 100644 --- a/providers/openlineage/tests/system/openlineage/example_openlineage_manual_lineage_dag.py +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_manual_lineage_dag.py @@ -127,8 +127,8 @@ def query_maximal() -> None: ctx = get_current_context() ti = ctx["task_instance"] - start = dt.datetime(2024, 5, 1, 10, 0, 0, tzinfo=dt.timezone.utc) - end = dt.datetime(2024, 5, 1, 10, 0, 5, tzinfo=dt.timezone.utc) + start = dt.datetime(2024, 5, 1, 10, 0, 0, tzinfo=dt.UTC) + end = dt.datetime(2024, 5, 1, 10, 0, 5, tzinfo=dt.UTC) emit_query_lineage( query_id="qid-max-1", diff --git a/providers/openlineage/tests/unit/openlineage/api/test_datasets.py b/providers/openlineage/tests/unit/openlineage/api/test_datasets.py index 3f67c554ea58b..b02cc7fafc42c 100644 --- a/providers/openlineage/tests/unit/openlineage/api/test_datasets.py +++ b/providers/openlineage/tests/unit/openlineage/api/test_datasets.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from unittest import mock import pytest @@ -36,15 +36,15 @@ def _make_task_instance(): task_id="task_id", try_number=1, map_index=-1, - logical_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + logical_date=datetime(2024, 1, 1, tzinfo=UTC), ) dag_run = mock.MagicMock( - logical_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + logical_date=datetime(2024, 1, 1, tzinfo=UTC), clear_number=0, - run_after=datetime(2024, 1, 1, tzinfo=timezone.utc), + run_after=datetime(2024, 1, 1, tzinfo=UTC), conf={}, - data_interval_start=datetime(2024, 1, 1, tzinfo=timezone.utc), - data_interval_end=datetime(2024, 1, 2, tzinfo=timezone.utc), + data_interval_start=datetime(2024, 1, 1, tzinfo=UTC), + data_interval_end=datetime(2024, 1, 2, tzinfo=UTC), ) task = mock.MagicMock( owner="alice,bob", doc=None, doc_md=None, doc_json=None, doc_yaml=None, doc_rst=None diff --git a/providers/openlineage/tests/unit/openlineage/api/test_sql.py b/providers/openlineage/tests/unit/openlineage/api/test_sql.py index 565db5d8440f0..9c04a104eec6d 100644 --- a/providers/openlineage/tests/unit/openlineage/api/test_sql.py +++ b/providers/openlineage/tests/unit/openlineage/api/test_sql.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from unittest import mock import pytest @@ -34,12 +34,12 @@ def _make_task_instance(task_id: str = "task_id", dr_conf: dict | None = None): task_id=task_id, try_number=1, map_index=-1, - logical_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + logical_date=datetime(2024, 1, 1, tzinfo=UTC), ) dag_run = mock.MagicMock( - logical_date=datetime(2024, 1, 1, tzinfo=timezone.utc), + logical_date=datetime(2024, 1, 1, tzinfo=UTC), clear_number=0, - run_after=datetime(2024, 1, 1, tzinfo=timezone.utc), + run_after=datetime(2024, 1, 1, tzinfo=UTC), conf=dr_conf or {}, ) ti.dag_run = dag_run @@ -193,8 +193,8 @@ def test_emits_fail_event_with_error_message(patched_emit): def test_uses_explicit_start_and_end_times(patched_emit): ti = _make_task_instance() - start = datetime(2024, 5, 1, 10, 0, tzinfo=timezone.utc) - end = datetime(2024, 5, 1, 10, 5, tzinfo=timezone.utc) + start = datetime(2024, 5, 1, 10, 0, tzinfo=UTC) + end = datetime(2024, 5, 1, 10, 5, tzinfo=UTC) emit_query_lineage( query_id="qid", query_source_namespace="snowflake://ACCT", diff --git a/providers/openlineage/tests/unit/openlineage/plugins/test_macros.py b/providers/openlineage/tests/unit/openlineage/plugins/test_macros.py index 5a0dcb6e348eb..1bf5f5aa2da89 100644 --- a/providers/openlineage/tests/unit/openlineage/plugins/test_macros.py +++ b/providers/openlineage/tests/unit/openlineage/plugins/test_macros.py @@ -16,7 +16,7 @@ # under the License. from __future__ import annotations -from datetime import datetime, timezone +from datetime import UTC, datetime from unittest import mock import pytest @@ -51,13 +51,13 @@ def test_lineage_job_name(): dag_id="dag_id", task_id="task_id", try_number=1, - **{LOGICAL_DATE_KEY: datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=timezone.utc)}, + **{LOGICAL_DATE_KEY: datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=UTC)}, ) assert lineage_job_name(task_instance) == "dag_id.task_id" def test_lineage_run_id(): - date = datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=timezone.utc) + date = datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=UTC) dag_run = mock.MagicMock(run_id="run_id") dag_run.logical_date = date task_instance = mock.MagicMock( @@ -80,7 +80,7 @@ def test_lineage_run_id(): @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="Test only for Airflow 3.0+") def test_lineage_run_after_airflow_3(): dag_run = mock.MagicMock(run_id="run_id") - dag_run.run_after = datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=timezone.utc) + dag_run.run_after = datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=UTC) dag_run.logical_date = None task_instance = mock.MagicMock( dag_id="dag_id", @@ -105,7 +105,7 @@ def test_lineage_parent_id(mock_run_id): dag_id="dag_id", task_id="task_id", try_number=1, - **{LOGICAL_DATE_KEY: datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=timezone.utc)}, + **{LOGICAL_DATE_KEY: datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=UTC)}, ) actual = lineage_parent_id(task_instance) expected = f"{_DAG_NAMESPACE}/dag_id.task_id/run_id" @@ -347,7 +347,7 @@ def test_lineage_root_macros_use_parent_from_conf_when_root_missing_af2(): ) @pytest.mark.skipif(AIRFLOW_V_3_0_PLUS, reason="Test only for Airflow 2") def test_lineage_root_macros_use_dagrun_info_when_missing_or_invalid_conf_af2(conf): - date = datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=timezone.utc) + date = datetime(2020, 1, 1, 1, 1, 1, 0, tzinfo=UTC) conf = {} dag_run = mock.MagicMock(run_id="run_id", conf=conf) dag_run.logical_date = date diff --git a/providers/openlineage/tests/unit/openlineage/utils/test_utils.py b/providers/openlineage/tests/unit/openlineage/utils/test_utils.py index 0b5d4fe2f8788..2fce4012f6af2 100644 --- a/providers/openlineage/tests/unit/openlineage/utils/test_utils.py +++ b/providers/openlineage/tests/unit/openlineage/utils/test_utils.py @@ -212,21 +212,21 @@ def test_get_airflow_dag_run_facet(): dagrun_mock.conf = {} dagrun_mock.clear_number = 0 dagrun_mock.dag_id = dag.dag_id - dagrun_mock.data_interval_start = datetime.datetime(2024, 6, 1, 1, 2, 3, tzinfo=datetime.timezone.utc) - dagrun_mock.data_interval_end = datetime.datetime(2024, 6, 1, 2, 3, 4, tzinfo=datetime.timezone.utc) + dagrun_mock.data_interval_start = datetime.datetime(2024, 6, 1, 1, 2, 3, tzinfo=datetime.UTC) + dagrun_mock.data_interval_end = datetime.datetime(2024, 6, 1, 2, 3, 4, tzinfo=datetime.UTC) dagrun_mock.external_trigger = True dagrun_mock.run_id = "manual_2024-06-01T00:00:00+00:00" dagrun_mock.run_type = DagRunType.MANUAL - dagrun_mock.execution_date = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.timezone.utc) - dagrun_mock.logical_date = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.timezone.utc) - dagrun_mock.run_after = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.timezone.utc) - dagrun_mock.start_date = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.timezone.utc) - dagrun_mock.end_date = datetime.datetime(2024, 6, 1, 1, 2, 14, 34172, tzinfo=datetime.timezone.utc) + dagrun_mock.execution_date = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.UTC) + dagrun_mock.logical_date = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.UTC) + dagrun_mock.run_after = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.UTC) + dagrun_mock.start_date = datetime.datetime(2024, 6, 1, 1, 2, 4, tzinfo=datetime.UTC) + dagrun_mock.end_date = datetime.datetime(2024, 6, 1, 1, 2, 14, 34172, tzinfo=datetime.UTC) dagrun_mock.triggering_user_name = "user1" dagrun_mock.triggered_by = "something" dagrun_mock.note = "note" dagrun_mock.partition_key = "some_partition_key" - dagrun_mock.partition_date = datetime.datetime(2024, 6, 1, 2, 3, 34, tzinfo=datetime.timezone.utc) + dagrun_mock.partition_date = datetime.datetime(2024, 6, 1, 2, 3, 34, tzinfo=datetime.UTC) dagrun_mock.dag_versions = [ MagicMock( bundle_name="bundle_name", @@ -306,8 +306,8 @@ def test_get_airflow_dag_run_facet(): ({"start_date": "2024-06-01T01:02:04+00:00", "end_date": "2024-06-01T01:02:14.034172+00:00"}, None), ( { - "start_date": datetime.datetime(2025, 1, 1, 6, 1, 1, tzinfo=datetime.timezone.utc), - "end_date": datetime.datetime(2025, 1, 1, 6, 1, 12, 3456, tzinfo=datetime.timezone.utc), + "start_date": datetime.datetime(2025, 1, 1, 6, 1, 1, tzinfo=datetime.UTC), + "end_date": datetime.datetime(2025, 1, 1, 6, 1, 12, 3456, tzinfo=datetime.UTC), }, 11.003456, ), @@ -2883,7 +2883,7 @@ def test_dagrun_with_deadline_and_alert(self): alert.callback_def = {"path": "my_module.on_deadline_missed", "kwargs": {}} deadline = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.UTC) deadline.missed = False deadline.deadline_alert = alert @@ -2920,12 +2920,12 @@ def test_dagrun_with_multiple_deadlines(self): alert2.callback_def = {"path": "mod.cb2", "kwargs": {"notify": True}} d1 = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - d1.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + d1.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.UTC) d1.missed = True d1.deadline_alert = alert1 d2 = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - d2.deadline_time = datetime.datetime(2025, 6, 1, 14, 0, 0, tzinfo=datetime.timezone.utc) + d2.deadline_time = datetime.datetime(2025, 6, 1, 14, 0, 0, tzinfo=datetime.UTC) d2.missed = False d2.deadline_alert = alert2 @@ -2955,7 +2955,7 @@ def test_dagrun_with_multiple_deadlines(self): def test_dagrun_deadline_alert_access_fails(self): """When the alert relationship can't be loaded, execution details still appear.""" deadline = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.UTC) deadline.missed = False type(deadline).deadline_alert = PropertyMock(side_effect=Exception("DB not available")) @@ -2981,7 +2981,7 @@ def test_dagrun_deadline_none_alert_fields_excluded(self): alert.callback_def = None deadline = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.UTC) deadline.missed = True deadline.deadline_alert = alert @@ -3035,7 +3035,7 @@ def test_dagrun_bad_deadline_skipped_others_preserved(self): bad_deadline.deadline_alert = None good_deadline = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - good_deadline.deadline_time = datetime.datetime(2025, 6, 1, 14, 0, 0, tzinfo=datetime.timezone.utc) + good_deadline.deadline_time = datetime.datetime(2025, 6, 1, 14, 0, 0, tzinfo=datetime.UTC) good_deadline.missed = True good_deadline.deadline_alert = None @@ -3058,7 +3058,7 @@ def test_dagrun_alert_attribute_access_raises(self): type(alert).reference = PropertyMock(side_effect=Exception("Column error")) deadline = MagicMock(spec=["deadline_time", "missed", "deadline_alert"]) - deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + deadline.deadline_time = datetime.datetime(2025, 6, 1, 12, 0, 0, tzinfo=datetime.UTC) deadline.missed = False deadline.deadline_alert = alert @@ -3082,7 +3082,7 @@ def test_dagrun_info_af3(mocked_dag_versions): from airflow.models.dag_version import DagVersion from airflow.utils.types import DagRunTriggeredByType - date = datetime.datetime(2024, 6, 1, tzinfo=datetime.timezone.utc) + date = datetime.datetime(2024, 6, 1, tzinfo=datetime.UTC) dv1 = DagVersion() dv2 = DagVersion() dv2.id = "version_id" @@ -3157,7 +3157,7 @@ def test_dagrun_info_af3(mocked_dag_versions): @pytest.mark.skipif(AIRFLOW_V_3_0_PLUS, reason="Airflow 2 test") def test_dagrun_info_af2(): - date = datetime.datetime(2024, 6, 1, tzinfo=datetime.timezone.utc) + date = datetime.datetime(2024, 6, 1, tzinfo=datetime.UTC) dag = DAG( "dag_id", schedule=None, @@ -3260,7 +3260,7 @@ def test_taskinstance_info_af3(): @pytest.mark.skipif(AIRFLOW_V_3_0_PLUS, reason="Airflow 2 test") @patch.object(TaskInstance, "log_url", "some_log_url") # Depends on the host, hard to test exact value def test_taskinstance_info_af2(): - some_date = datetime.datetime(2024, 6, 1, tzinfo=datetime.timezone.utc) + some_date = datetime.datetime(2024, 6, 1, tzinfo=datetime.UTC) task_obj = PythonOperator(task_id="task_id", python_callable=lambda x: x) ti = TaskInstance( task=task_obj, run_id="task_instance_run_id", state=TaskInstanceState.RUNNING, map_index=2 @@ -3748,7 +3748,7 @@ def test_is_dag_run_asset_triggered_af2(): def test_build_task_instance_ol_run_id(): """Test deterministic UUID generation for task instance.""" - logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) run_id = build_task_instance_ol_run_id( dag_id="test_dag", task_id="test_task", @@ -3782,7 +3782,7 @@ def test_build_task_instance_ol_run_id(): def test_build_dag_run_ol_run_id(): """Test deterministic UUID generation for DAG run.""" - logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) run_id = build_dag_run_ol_run_id( dag_id="test_dag", logical_date=logical_date, @@ -3838,7 +3838,7 @@ class TestExtractOlInfoFromAssetEvent: def test_extract_ol_info_from_task_instance(self): """Test extraction from TaskInstance (priority 1).""" - logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) # Mock TaskInstance - using MagicMock without spec to avoid SQLAlchemy mapper inspection ti = MagicMock() @@ -3908,7 +3908,7 @@ def test_extract_ol_info_from_task_instance_no_logical_date(self): @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="Airflow 3 specific test") def test_extract_ol_info_from_task_instance_run_after_fallback(self): """Test extraction from TaskInstance with run_after fallback (AF3).""" - run_after = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + run_after = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) # Mock TaskInstance ti = MagicMock() @@ -4248,7 +4248,7 @@ def test_get_dag_job_dependency_facet_insufficient_info_skipped(self, mock_get_e @patch("airflow.providers.openlineage.utils.utils._get_eagerly_loaded_dagrun_consumed_asset_events") def test_get_dag_job_dependency_facet_with_events(self, mock_get_events): """Test facet generation with asset events - tests full flow.""" - logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) # Create mock asset events with source TaskInstance (priority 1 source) ti1 = MagicMock() @@ -4342,7 +4342,7 @@ def test_get_dag_job_dependency_facet_with_events(self, mock_get_events): @patch("airflow.providers.openlineage.utils.utils._get_eagerly_loaded_dagrun_consumed_asset_events") def test_get_dag_job_dependency_facet_deduplication(self, mock_get_events): """Test that duplicate asset events from same job/run are deduplicated.""" - logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.timezone.utc) + logical_date = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) # Create two events from the same source TI (should be deduplicated) ti = MagicMock() @@ -4465,8 +4465,8 @@ def test_build_task_event_run_facets_composes_all_sections( dag = MagicMock(dag_id="my_dag") dag_run = MagicMock( conf={"k": "v"}, - data_interval_start=datetime.datetime(2024, 1, 1, tzinfo=datetime.timezone.utc), - data_interval_end=datetime.datetime(2024, 1, 2, tzinfo=datetime.timezone.utc), + data_interval_start=datetime.datetime(2024, 1, 1, tzinfo=datetime.UTC), + data_interval_end=datetime.datetime(2024, 1, 2, tzinfo=datetime.UTC), ) facets = build_task_event_run_facets( task_instance=MagicMock(), diff --git a/providers/opsgenie/src/airflow/providers/opsgenie/typing/opsgenie.py b/providers/opsgenie/src/airflow/providers/opsgenie/typing/opsgenie.py index 4e6621785d63d..e02815c0a21b3 100644 --- a/providers/opsgenie/src/airflow/providers/opsgenie/typing/opsgenie.py +++ b/providers/opsgenie/src/airflow/providers/opsgenie/typing/opsgenie.py @@ -16,9 +16,7 @@ # under the License. from __future__ import annotations -from typing import TypedDict - -from typing_extensions import NotRequired, Required # For compat with Python < 3.11 +from typing import NotRequired, Required, TypedDict # For compat with Python < 3.11 class CreateAlertPayload(TypedDict): diff --git a/providers/sftp/tests/unit/sftp/sensors/test_sftp.py b/providers/sftp/tests/unit/sftp/sensors/test_sftp.py index be81f712b280b..5bb4341862083 100644 --- a/providers/sftp/tests/unit/sftp/sensors/test_sftp.py +++ b/providers/sftp/tests/unit/sftp/sensors/test_sftp.py @@ -19,7 +19,7 @@ import importlib import warnings -from datetime import datetime, timezone as stdlib_timezone +from datetime import UTC, datetime from unittest import mock from unittest.mock import Mock, patch @@ -148,7 +148,7 @@ def test_file_not_new_enough(self, sftp_hook_mock): "newer_than", ( datetime(2020, 1, 2), - datetime(2020, 1, 2, tzinfo=stdlib_timezone.utc), + datetime(2020, 1, 2, tzinfo=UTC), "2020-01-02", "2020-01-02 00:00:00+00:00", "2020-01-02 00:00:00.001+00:00", diff --git a/providers/slack/src/airflow/providers/slack/hooks/slack.py b/providers/slack/src/airflow/providers/slack/hooks/slack.py index 2a5f79339a4b6..ee838da830096 100644 --- a/providers/slack/src/airflow/providers/slack/hooks/slack.py +++ b/providers/slack/src/airflow/providers/slack/hooks/slack.py @@ -31,12 +31,11 @@ from collections.abc import Sequence from functools import cached_property from pathlib import Path -from typing import TYPE_CHECKING, Any, TypedDict +from typing import TYPE_CHECKING, Any, NotRequired, TypedDict from slack_sdk import WebClient from slack_sdk.errors import SlackApiError from slack_sdk.web.async_client import AsyncWebClient -from typing_extensions import NotRequired from airflow.providers.common.compat.connection import get_async_connection from airflow.providers.common.compat.sdk import AirflowException, AirflowNotFoundException, BaseHook diff --git a/providers/snowflake/src/airflow/providers/snowflake/utils/sql_api_generate_jwt.py b/providers/snowflake/src/airflow/providers/snowflake/utils/sql_api_generate_jwt.py index 47558f1a1fdc0..75ab553504735 100644 --- a/providers/snowflake/src/airflow/providers/snowflake/utils/sql_api_generate_jwt.py +++ b/providers/snowflake/src/airflow/providers/snowflake/utils/sql_api_generate_jwt.py @@ -19,7 +19,7 @@ import base64 import hashlib import logging -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from typing import Any # This class relies on the PyJWT module (https://pypi.org/project/PyJWT/). @@ -87,7 +87,7 @@ def __init__( self.lifetime = lifetime self.renewal_delay = renewal_delay self.private_key = private_key - self.renew_time = datetime.now(timezone.utc) + self.renew_time = datetime.now(UTC) self.token: str | None = None def prepare_account_name_for_jwt(self, raw_account: str) -> str: @@ -116,7 +116,7 @@ def get_token(self) -> str | None: If a JWT has been already been generated earlier, return the previously generated token unless the specified renewal time has passed. """ - now = datetime.now(timezone.utc) # Fetch the current time + now = datetime.now(UTC) # Fetch the current time # If the token has expired or doesn't exist, regenerate the token. if self.token is None or self.renew_time <= now: diff --git a/providers/snowflake/tests/unit/snowflake/hooks/test_snowflake_sql_api.py b/providers/snowflake/tests/unit/snowflake/hooks/test_snowflake_sql_api.py index eb4fc930c6e3f..14bd79f16511a 100644 --- a/providers/snowflake/tests/unit/snowflake/hooks/test_snowflake_sql_api.py +++ b/providers/snowflake/tests/unit/snowflake/hooks/test_snowflake_sql_api.py @@ -16,7 +16,6 @@ # under the License. from __future__ import annotations -import asyncio import base64 import uuid from collections.abc import Mapping @@ -1726,7 +1725,7 @@ async def test_make_api_call_with_retries_async_retries_on_timeout_error(self, m ) mock_async_request.__aenter__.side_effect = [ - asyncio.TimeoutError(), + TimeoutError(), create_async_request_client_response_success(json=GET_RESPONSE, status_code=200), ] diff --git a/providers/snowflake/tests/unit/snowflake/operators/test_snowpark_containers.py b/providers/snowflake/tests/unit/snowflake/operators/test_snowpark_containers.py index fa888c383f809..1e344d080d0ef 100644 --- a/providers/snowflake/tests/unit/snowflake/operators/test_snowpark_containers.py +++ b/providers/snowflake/tests/unit/snowflake/operators/test_snowpark_containers.py @@ -17,7 +17,7 @@ from __future__ import annotations import itertools -from datetime import datetime, timedelta, timezone +from datetime import UTC, datetime, timedelta from unittest import mock import pytest @@ -371,7 +371,7 @@ def test_execute_defer_without_execution_timeout(self, mock_submit): @mock.patch.object(SnowparkContainerJobOperator, "_submit_job", return_value=JOB_NAME) def test_execute_defer_uses_execution_timeout_for_deadline_and_buffer(self, mock_submit, time_machine): time_machine.move_to(1000, tick=False) - context = {"ti": mock.Mock(start_date=datetime.fromtimestamp(1000, tz=timezone.utc))} + context = {"ti": mock.Mock(start_date=datetime.fromtimestamp(1000, tz=UTC))} op = _make_operator( deferrable=True, timeout=3600, diff --git a/providers/snowflake/tests/unit/snowflake/utils/test_openlineage.py b/providers/snowflake/tests/unit/snowflake/utils/test_openlineage.py index a460aa213e8b5..6b27372e3378e 100644 --- a/providers/snowflake/tests/unit/snowflake/utils/test_openlineage.py +++ b/providers/snowflake/tests/unit/snowflake/utils/test_openlineage.py @@ -159,15 +159,15 @@ def test_process_data_from_api(): { "QUERY_ID": "ABC", "EXECUTION_STATUS": "SUCCESS", - "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.timezone.utc), - "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.timezone.utc), + "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.UTC), + "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.UTC), "QUERY_TEXT": "SELECT * FROM test_table;", "ERROR_CODE": None, "ERROR_MESSAGE": None, }, { - "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.timezone.utc), - "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.timezone.utc), + "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.UTC), + "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.UTC), }, ] result = _process_data_from_api(data=data) @@ -309,8 +309,8 @@ def test_get_queries_details_from_snowflake_single_query_api_hook(mock_run_singl expected_details = { "QUERY_ID": "ABC", "EXECUTION_STATUS": "SUCCESS", - "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.timezone.utc), - "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.timezone.utc), + "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.UTC), + "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.UTC), "QUERY_TEXT": "SELECT * FROM test_table;", "ERROR_CODE": None, "ERROR_MESSAGE": None, @@ -395,8 +395,8 @@ def test_get_queries_details_from_snowflake_multiple_queries_api_hook(mock_run_s { "QUERY_ID": "ABC", "EXECUTION_STATUS": "SUCCESS", - "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.timezone.utc), - "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.timezone.utc), + "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.UTC), + "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.UTC), "QUERY_TEXT": "SELECT * FROM table1;", "ERROR_CODE": None, "ERROR_MESSAGE": None, @@ -404,8 +404,8 @@ def test_get_queries_details_from_snowflake_multiple_queries_api_hook(mock_run_s { "QUERY_ID": "DEF", "EXECUTION_STATUS": "FAILED", - "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.timezone.utc), - "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.timezone.utc), + "START_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 326000, tzinfo=datetime.UTC), + "END_TIME": datetime.datetime(2025, 6, 18, 11, 12, 51, 387000, tzinfo=datetime.UTC), "QUERY_TEXT": "SELECT * FROM table2;", "ERROR_CODE": "123", "ERROR_MESSAGE": "Some error", diff --git a/providers/standard/src/airflow/providers/standard/example_dags/example_sensors.py b/providers/standard/src/airflow/providers/standard/example_dags/example_sensors.py index 73ccaadbe0f02..874a5dbca57d6 100644 --- a/providers/standard/src/airflow/providers/standard/example_dags/example_sensors.py +++ b/providers/standard/src/airflow/providers/standard/example_dags/example_sensors.py @@ -63,20 +63,18 @@ def failure_callable(): # [END example_time_delta_sensor_async] # [START example_time_sensors] - t1 = TimeSensor( - task_id="fire_immediately", target_time=datetime.datetime.now(tz=datetime.timezone.utc).time() - ) + t1 = TimeSensor(task_id="fire_immediately", target_time=datetime.datetime.now(tz=datetime.UTC).time()) t2 = TimeSensor( task_id="timeout_after_second_date_in_the_future", timeout=1, soft_fail=True, - target_time=(datetime.datetime.now(tz=datetime.timezone.utc) + datetime.timedelta(hours=1)).time(), + target_time=(datetime.datetime.now(tz=datetime.UTC) + datetime.timedelta(hours=1)).time(), ) t1a = TimeSensor( task_id="fire_immediately_async", - target_time=datetime.datetime.now(tz=datetime.timezone.utc).time(), + target_time=datetime.datetime.now(tz=datetime.UTC).time(), deferrable=True, ) @@ -84,7 +82,7 @@ def failure_callable(): task_id="timeout_after_second_date_in_the_future_async", timeout=1, soft_fail=True, - target_time=(datetime.datetime.now(tz=datetime.timezone.utc) + datetime.timedelta(hours=1)).time(), + target_time=(datetime.datetime.now(tz=datetime.UTC) + datetime.timedelta(hours=1)).time(), deferrable=True, ) # [END example_time_sensors] diff --git a/providers/standard/tests/unit/standard/decorators/test_external_python.py b/providers/standard/tests/unit/standard/decorators/test_external_python.py index 17994faf82e6f..2f3fa42f5d6e3 100644 --- a/providers/standard/tests/unit/standard/decorators/test_external_python.py +++ b/providers/standard/tests/unit/standard/decorators/test_external_python.py @@ -210,7 +210,7 @@ def f(_): return None with dag_maker(serialized=True): - v = f(datetime.datetime.now(tz=datetime.timezone.utc)) + v = f(datetime.datetime.now(tz=datetime.UTC)) dr = dag_maker.create_dagrun() ti = dr.get_task_instances()[0] diff --git a/providers/standard/tests/unit/standard/decorators/test_python_virtualenv.py b/providers/standard/tests/unit/standard/decorators/test_python_virtualenv.py index cc02d92a3b1a5..0eeb8f16a6625 100644 --- a/providers/standard/tests/unit/standard/decorators/test_python_virtualenv.py +++ b/providers/standard/tests/unit/standard/decorators/test_python_virtualenv.py @@ -314,7 +314,7 @@ def f(_): return None with dag_maker(serialized=True): - f(datetime.datetime.now(tz=datetime.timezone.utc)) + f(datetime.datetime.now(tz=datetime.UTC)) dr = dag_maker.create_dagrun() dag_maker.run_ti("f", dr) diff --git a/providers/standard/tests/unit/standard/operators/test_datetime.py b/providers/standard/tests/unit/standard/operators/test_datetime.py index c9e38e6c9602f..0535423018978 100644 --- a/providers/standard/tests/unit/standard/operators/test_datetime.py +++ b/providers/standard/tests/unit/standard/operators/test_datetime.py @@ -149,8 +149,8 @@ def test_branch_datetime_operator_falls_within_range(self, target_lower, target_ @pytest.mark.parametrize( "date", ( - datetime.datetime(2020, 7, 7, 12, 0, 0, tzinfo=datetime.timezone.utc), - datetime.datetime(2020, 6, 7, 9, 0, 0, tzinfo=datetime.timezone.utc), + datetime.datetime(2020, 7, 7, 12, 0, 0, tzinfo=datetime.UTC), + datetime.datetime(2020, 6, 7, 9, 0, 0, tzinfo=datetime.UTC), ), ) def test_branch_datetime_operator_falls_outside_range(self, date, target_lower, target_upper): diff --git a/providers/standard/tests/unit/standard/operators/test_python.py b/providers/standard/tests/unit/standard/operators/test_python.py index 46f3a77a8e50f..d7b2269662224 100644 --- a/providers/standard/tests/unit/standard/operators/test_python.py +++ b/providers/standard/tests/unit/standard/operators/test_python.py @@ -29,7 +29,7 @@ import warnings from collections import namedtuple from collections.abc import Generator -from datetime import date, datetime, timezone as _timezone +from datetime import UTC, date, datetime from functools import partial from importlib.util import find_spec from pathlib import Path @@ -2023,7 +2023,7 @@ def test_nonimported_as_arg(self): def f(_): return None - self.run_as_task(f, op_args=[datetime.now(tz=_timezone.utc)]) + self.run_as_task(f, op_args=[datetime.now(tz=UTC)]) def test_context(self): def f(templates_dict): diff --git a/providers/standard/tests/unit/standard/utils/test_skipmixin.py b/providers/standard/tests/unit/standard/utils/test_skipmixin.py index c1507577abf59..6f1d799e5777a 100644 --- a/providers/standard/tests/unit/standard/utils/test_skipmixin.py +++ b/providers/standard/tests/unit/standard/utils/test_skipmixin.py @@ -60,7 +60,7 @@ def teardown_method(self): self.clean_db() def test_skip(self, dag_maker, session, time_machine): - now = datetime.datetime.now(tz=datetime.timezone.utc) + now = datetime.datetime.now(tz=datetime.UTC) time_machine.move_to(now, tick=False) with dag_maker("dag"): tasks = [EmptyOperator(task_id="task")]