This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/v3-3-test by this push:
new ef00c94574b [v3-3-test] Fix TaskInstance duration calculation with
SQLite (#68142) (#70734)
ef00c94574b is described below
commit ef00c94574b8a5a582f8d608948774ddb44700aa
Author: Rahul Vats <[email protected]>
AuthorDate: Fri Jul 31 00:56:43 2026 +0530
[v3-3-test] Fix TaskInstance duration calculation with SQLite (#68142)
(#70734)
---
airflow-core/src/airflow/models/taskinstance.py | 5 ++-
.../tests/unit/models/test_taskinstance.py | 37 +++++++++++++++++++++-
2 files changed, 38 insertions(+), 4 deletions(-)
diff --git a/airflow-core/src/airflow/models/taskinstance.py
b/airflow-core/src/airflow/models/taskinstance.py
index f905e4c8e3d..f2bd3f18beb 100644
--- a/airflow-core/src/airflow/models/taskinstance.py
+++ b/airflow-core/src/airflow/models/taskinstance.py
@@ -2281,9 +2281,8 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload):
return query.values(
{
"end_date": end_date,
- "duration": (
- (func.strftime("%s", end_date) - func.strftime("%s",
cls.start_date))
- + func.round((func.strftime("%f", end_date) -
func.strftime("%f", cls.start_date)), 3)
+ "duration": func.round(
+ (func.julianday(end_date) -
func.julianday(cls.start_date)) * 86400, 3
),
}
)
diff --git a/airflow-core/tests/unit/models/test_taskinstance.py
b/airflow-core/tests/unit/models/test_taskinstance.py
index 1221a96d9f8..d02fcc593b1 100644
--- a/airflow-core/tests/unit/models/test_taskinstance.py
+++ b/airflow-core/tests/unit/models/test_taskinstance.py
@@ -34,7 +34,7 @@ import uuid6
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.trace.propagation.tracecontext import
TraceContextTextMapPropagator
-from sqlalchemy import delete, func, inspect as sa_inspect, select
+from sqlalchemy import delete, func, inspect as sa_inspect, select, update
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import load_only
from sqlalchemy.orm.attributes import set_committed_value
@@ -1805,6 +1805,41 @@ class TestTaskInstance:
ti.set_duration()
assert ti.duration is None
+ @pytest.mark.backend("sqlite")
+ @pytest.mark.parametrize(
+ ("start_date", "end_date", "expected_duration"),
+ [
+ (
+ timezone.datetime(2026, 6, 7, 12, 0, 1, 200000),
+ timezone.datetime(2026, 6, 7, 12, 0, 3, 500000),
+ 2.3,
+ ),
+ (
+ timezone.datetime(2026, 6, 7, 12, 0, 59, 200000),
+ timezone.datetime(2026, 6, 7, 12, 1, 0, 500000),
+ 1.3,
+ ),
+ (
+ timezone.datetime(2026, 6, 7, 12),
+ timezone.datetime(2026, 6, 7, 13),
+ 3600.0,
+ ),
+ ],
+ )
+ def test_duration_expression_update_sqlite(
+ self, create_task_instance, session, start_date, end_date,
expected_duration
+ ):
+ ti = create_task_instance()
+ ti.start_date = start_date
+ session.flush()
+
+ query = update(TI).where(TI.id == ti.id)
+ session.execute(TI.duration_expression_update(end_date, query,
session.get_bind()))
+ ti.refresh_from_db(session=session)
+
+ assert ti.end_date == end_date
+ assert ti.duration == expected_duration
+
def test_outlet_asset_extra(self, dag_maker: DagMaker, session: Session):
from airflow.sdk.definitions.asset import Asset