myps6415 commented on code in PR #67592:
URL: https://github.com/apache/airflow/pull/67592#discussion_r4162062323


##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py:
##########
@@ -244,6 +251,12 @@ def ti_run(
                 extra=json.dumps({"host_name": ti_run_payload.hostname}) if 
ti_run_payload.hostname else None,
             )
         )
+        # One sample per queue wait, not per try: the scheduler refreshes 
queued_dttm on every
+        # queueing, so a retry and a resume from deferral each waited for a 
slot of their own.
+        # task.scheduled_duration counts per try instead, so the two disagree 
on retries by design.
+        # queued_dttm is None only in rare races and test setups.

Review Comment:
   Used your wording for the scheduled_duration line, and the next line now 
names runs that skip the scheduler's queueing, e.g. dag.test(). 0353c4f



##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py:
##########
@@ -299,6 +324,14 @@ def ti_run(
             or 0
         )
 
+        if emit_queued_duration:
+            # Tags mirror the sibling task.scheduled_duration, which 
emit_state_change_metric sends
+            # as {**ti.stats_tags, "queue": ti.queue}; stats_tags reads the 
team off the transient
+            # _team_name. stats.timing also emits the legacy dotted name from 
the metrics registry.
+            dr._team_name = dr.team_name
+            tags = {**dr.stats_tags, "task_id": ti.task_id, "queue": ti.queue}
+            stats.timing("task.queued_duration", timezone.utcnow() - 
ti.queued_dttm, tags=tags)

Review Comment:
   Moved — the emit now sits just before `return context`, after 
`issue_execution_token`, so only the session commit can fail past it. 
test_ti_run_skips_queued_duration_metric_when_the_response_fails covers it and 
fails if the block goes back where it was. 0353c4f



##########
airflow-core/newsfragments/67592.bugfix.rst:
##########
@@ -0,0 +1 @@
+Restore the ``task.queued_duration`` metric, which stopped being emitted when 
Airflow 3 workers moved to the Task SDK, and record it once per queue wait so 
that a task resuming from deferral also reports the wait for its worker slot.

Review Comment:
   Rewritten. It no longer frames deferral resumes as new, and now names 
retries and reschedule pokes, the roughly 60-versus-1 sample count for a sensor 
poking every minute, and the percentile shift for alerts carried over from 2.x. 
0353c4f



##########
airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py:
##########
@@ -1245,6 +1245,309 @@ def test_ti_run_creates_audit_log(self, client, 
session, create_task_instance, t
         assert logs[0].owner == ti.task.owner
         assert logs[0].extra == '{"host_name": "random-hostname"}'
 
+    @pytest.mark.parametrize(
+        "scenario",
+        ["first_run", "retry", "deferral_resume"],
+    )
+    def test_ti_run_emits_queued_duration_metric(
+        self, client, session, create_task_instance, time_machine, scenario
+    ):
+        """task.queued_duration is emitted once per queue wait.
+
+        The scheduler refreshes queued_dttm every time it queues a task, so a 
retry and a
+        resume from deferral each waited for a worker slot of their own and 
must emit too. A
+        retry still carries the previous attempt's end_date on the row when 
ti_run is reached,
+        so this also asserts the emit does not depend on end_date being unset.
+        """
+        queued_at = timezone.parse("2024-09-30T12:00:00Z")
+        run_at = queued_at.add(seconds=42)
+
+        ti = create_task_instance(
+            task_id=f"test_ti_run_emits_queued_duration_metric_{scenario}",
+            state=State.QUEUED,
+            dagrun_state=DagRunState.RUNNING,
+            session=session,
+            start_date=queued_at,
+            dag_id=str(uuid4()),
+        )
+        ti.queued_dttm = queued_at
+        ti.queue = "default"
+        if scenario == "retry":
+            # A retried TI still has the previous attempt's end_date set on 
the row until
+            # ti_run clears it; the metric must fire regardless.
+            ti.end_date = queued_at.add(seconds=10)
+        elif scenario == "deferral_resume":
+            ti.next_method = "execute_complete"
+        session.commit()
+
+        # The metric has to stay sliceable the same way as its sibling 
task.scheduled_duration,
+        # which emit_state_change_metric sends as {**ti.stats_tags, "queue": 
ti.queue}. Deriving
+        # the expectation from that same expression makes the two drift apart 
only if this fails.
+        expected_tags = {**ti.stats_tags, "queue": ti.queue}
+        assert "run_type" in expected_tags
+
+        time_machine.move_to(run_at, tick=False)
+
+        with 
mock.patch("airflow.api_fastapi.execution_api.routes.task_instances.stats") as 
mock_stats:
+            response = client.patch(
+                f"/execution/task-instances/{ti.id}/run",
+                json={
+                    "state": "running",
+                    "hostname": "random-hostname",
+                    "unixname": "random-unixname",
+                    "pid": 100,
+                    "start_date": run_at.isoformat(),
+                },
+            )
+
+        assert response.status_code == 200
+        mock_stats.timing.assert_called_once_with(
+            "task.queued_duration",
+            run_at - queued_at,
+            tags=expected_tags,
+        )
+
+    def test_ti_run_skips_queued_duration_metric_without_queued_dttm(
+        self, client, session, create_task_instance, time_machine
+    ):
+        """queued_dttm is what the wait is measured from, so a row without one 
(rare race /
+        test setups) has nothing to report."""

Review Comment:
   Same wording fix applied here. 0353c4f



##########
airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py:
##########
@@ -1245,6 +1245,309 @@ def test_ti_run_creates_audit_log(self, client, 
session, create_task_instance, t
         assert logs[0].owner == ti.task.owner
         assert logs[0].extra == '{"host_name": "random-hostname"}'
 
+    @pytest.mark.parametrize(
+        "scenario",
+        ["first_run", "retry", "deferral_resume"],
+    )
+    def test_ti_run_emits_queued_duration_metric(
+        self, client, session, create_task_instance, time_machine, scenario
+    ):
+        """task.queued_duration is emitted once per queue wait.
+
+        The scheduler refreshes queued_dttm every time it queues a task, so a 
retry and a
+        resume from deferral each waited for a worker slot of their own and 
must emit too. A
+        retry still carries the previous attempt's end_date on the row when 
ti_run is reached,
+        so this also asserts the emit does not depend on end_date being unset.
+        """
+        queued_at = timezone.parse("2024-09-30T12:00:00Z")
+        run_at = queued_at.add(seconds=42)
+
+        ti = create_task_instance(
+            task_id=f"test_ti_run_emits_queued_duration_metric_{scenario}",
+            state=State.QUEUED,
+            dagrun_state=DagRunState.RUNNING,
+            session=session,
+            start_date=queued_at,
+            dag_id=str(uuid4()),
+        )
+        ti.queued_dttm = queued_at
+        ti.queue = "default"
+        if scenario == "retry":
+            # A retried TI still has the previous attempt's end_date set on 
the row until
+            # ti_run clears it; the metric must fire regardless.
+            ti.end_date = queued_at.add(seconds=10)
+        elif scenario == "deferral_resume":
+            ti.next_method = "execute_complete"
+        session.commit()
+
+        # The metric has to stay sliceable the same way as its sibling 
task.scheduled_duration,
+        # which emit_state_change_metric sends as {**ti.stats_tags, "queue": 
ti.queue}. Deriving
+        # the expectation from that same expression makes the two drift apart 
only if this fails.
+        expected_tags = {**ti.stats_tags, "queue": ti.queue}
+        assert "run_type" in expected_tags
+
+        time_machine.move_to(run_at, tick=False)
+
+        with 
mock.patch("airflow.api_fastapi.execution_api.routes.task_instances.stats") as 
mock_stats:

Review Comment:
   Added autospec=True to all seven new stats patches. 0353c4f



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to