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 eb290a881f32a..b8dc02d69269c 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 @@ -1051,11 +1051,17 @@ def _clean(self, event: dict[str, Any], result: dict | None, context: Context) - reraise=True, ) def _write_logs(self, pod: k8s.V1Pod, follow: bool = False, since_time: DateTime | None = None) -> None: - since_seconds = ( - math.ceil((datetime.datetime.now(tz=datetime.timezone.utc) - since_time).total_seconds()) - if since_time - else None - ) + since_seconds = None + if since_time: + try: + since_seconds = math.ceil( + (datetime.datetime.now(tz=datetime.timezone.utc) - since_time).total_seconds() + ) + except TypeError: + self.log.warning( + "Error calculating since_seconds with since_time %s. Using None instead.", + since_time, + ) logs = self.client.read_namespaced_pod_log( name=pod.metadata.name, namespace=pod.metadata.namespace, 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 5a28eb2ff6d3f..66d1e05914d5a 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 @@ -26,6 +26,7 @@ import pendulum import pytest import tenacity +import time_machine from kubernetes.client import ApiClient, V1Pod, V1PodSecurityContext, V1PodStatus, models as k8s from kubernetes.client.exceptions import ApiException @@ -2917,6 +2918,35 @@ def test_write_logs_gives_up_after_max_retries( post_complete_action.assert_called_once() assert "Reading of logs interrupted with error" in caplog.text + @time_machine.travel("2026-01-01 00:00:00", tick=False) + @patch(KUB_OP_PATH.format("client")) + 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) + 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 + assert call_kwargs["since_seconds"] == 30 + + @patch(KUB_OP_PATH.format("client")) + def test_write_logs_with_invalid_since_time_falls_back_to_none(self, mocked_client): + """Test that a TypeError from an invalid since_time is caught, warns, and uses since_seconds=None.""" + pod = k8s.V1Pod(metadata=k8s.V1ObjectMeta(name=TEST_NAME, namespace=TEST_NAMESPACE)) + k = KubernetesPodOperator(task_id="task", get_logs=True) + + with patch.object(k.log, "warning") as mock_warning: + k._write_logs(pod, since_time="not-a-datetime") + + _, call_kwargs = mocked_client.read_namespaced_pod_log.call_args + assert call_kwargs["since_seconds"] is None + mock_warning.assert_called_once_with( + "Error calculating since_seconds with since_time %s. Using None instead.", + "not-a-datetime", + ) + @pytest.mark.parametrize( ("log_pod_spec_on_failure", "expect_match"), [