Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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)
Expand Down Expand Up @@ -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"]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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")

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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}"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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,
)
Expand Down Expand Up @@ -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,
)
Expand All @@ -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,
)
Expand All @@ -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,
)
Expand All @@ -151,15 +151,15 @@ 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,
)
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,
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
6 changes: 3 additions & 3 deletions providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"]

Expand Down Expand Up @@ -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"

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading