kaxil commented on code in PR #73763:
URL: https://github.com/apache/airflow/pull/73763#discussion_r4158824447


##########
airflow-ctl/src/airflowctl/api/datamodels/generated.py:
##########
@@ -1651,6 +1665,13 @@ class BackfillPostBody(BaseModel):
             title="Run On Latest Version",
         ),
     ] = None
+    drain_dag: Annotated[

Review Comment:
   `BackfillOperations.create` and `create_dry_run` dump with 
`exclude_none=True`, which keeps a `False`, so every `airflowctl backfill 
create` now sends `"drain_dag": false`. The 3.3.2 `BackfillPostBody` is strict 
and rejects that with `extra_forbidden`, the same break #63388 fixed for 
`team_name`. Dumping with `exclude_defaults=True`, as `TasksOperations.clear` 
does, avoids it. `dags trigger` on main already has the same problem through 
`bundle_version: null`, so it's worth fixing both together before the next 
airflowctl release.



##########
airflow-core/src/airflow/models/dag.py:
##########
@@ -509,6 +509,24 @@ def set_scheduling_state(self, state: DagSchedulingState) 
-> None:
         self.is_paused = state == DagSchedulingState.PAUSED
         self.is_draining = state == DagSchedulingState.DRAINING
 
+    @classmethod
+    def start_drain(cls, dag_id: str, *, session: Session) -> None:
+        """
+        Put the Dag into the draining state.
+
+        Call this in the transaction that creates the explicit run the drain 
is started for,
+        before the run is inserted. The row lock keeps 
``_finalize_draining_dags`` from
+        pausing the Dag before that run is committed. Taking it after the 
insert could
+        deadlock on MySQL, where the insert's foreign-key check already holds 
a shared lock
+        on the Dag row that concurrent triggers would both try to upgrade.

Review Comment:
   I don't think the MySQL sentence holds: `dag_run`'s foreign keys go to 
`log_template`, `backfill` and `dag_version`, and nothing inserted on these 
paths references `dag`, so the insert never takes a shared lock on the Dag row.
   
   The lock also covers less than this says. It protects a paused or active 
Dag, because the state change is an UPDATE. On a Dag that is already draining, 
`set_scheduling_state(DRAINING)` emits no UPDATE, so a finalizer whose `NOT 
EXISTS` snapshot predates the trigger's commit can still pause the Dag once the 
lock is released. Any manual run on a draining Dag already hits that race, so a 
re-check in the finalizer can be a follow-up. The docstring just shouldn't 
promise it.



##########
airflow-core/docs/core-concepts/dags.rst:
##########
@@ -813,6 +815,16 @@ the scheduler automatically changes the Dag to paused. A 
backfill started while
 the drain until its Dag runs have been created and finished. While a Dag is 
draining, you can cancel the drain
 to make the Dag active again.
 
+A manual run, backfill or asset materialization on a paused Dag can drain the 
Dag instead of unpausing it.
+Choose **Drain** under **Dag is paused** in the UI form, or set ``drain_dag`` 
to ``true`` in the REST API
+request. This changes the whole Dag, not only the new run: every unfinished 
Dag run proceeds, so Dag runs that

Review Comment:
   Runs of a paused backfill don't proceed here. 
`get_queued_dag_runs_to_set_running` skips them, and the finalizer still counts 
them as unfinished, so the Dag stays draining until that backfill is unpaused. 
The form already leaves them out of the count. This paragraph, the newsfragment 
and both `drain_dag` field descriptions say "every unfinished run proceeds" and 
should name the exception.



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py:
##########
@@ -817,6 +817,10 @@ def trigger_dag_run(
                     f"DAG with dag_id: '{dag_id}' does not support bundle 
versioning",
                 )
 
+        if body.drain_dag:
+            requires_access_dag(method="PUT", param_dag_id=dag_id)(request, 
user)

Review Comment:
   A user who can trigger but can't edit the Dag gets a bare `Forbidden` here 
(same in `assets.py:514` and `backfills.py:308`). `TriggerDAGForm` then 
disables submit on any 403 and doesn't clear the error when the option changes, 
so they can't fall back to Keep paused without reopening the modal. 
`get_auth_manager().is_authorized_dag(...)` with a message naming the missing 
permission, as the run check in `materialize_asset` does, would fix the message 
and drop the `request` params that `create_backfill` and `materialize_asset` 
gained only to feed this call. Resetting the mutation error on option change 
fixes the form.



##########
airflow-core/src/airflow/ui/src/components/TriggerDag/TriggerDAGModal.tsx:
##########
@@ -71,10 +63,13 @@ const TriggerDAGModal = ({
     undefined,
     {
       enabled: open,
+      // The paused state decides whether a drain is requested, so it must not 
come from a stale cache.

Review Comment:
   With react-query v5, `staleTime: 0` still renders cached data right away and 
refetches in the background, so for one round trip the form can show the stale 
state this comment rules out. Gating submit on `isFetching` would make the 
comment true; otherwise I'd soften it. Same in `CreateAssetEventModal.tsx`.



##########
airflow-core/src/airflow/models/backfill.py:
##########
@@ -711,6 +712,9 @@ def _create_backfill(
 
         first_info = dagrun_info_list[0]
         try:
+            if drain_dag:
+                # After the backfill row's commit, so the drain commits or 
rolls back with the runs.
+                DagModel.start_drain(dag_id, session=session)

Review Comment:
   With the lock taken before `_create_runs_*`, the Dag row stays locked until 
every DagRun and TI is inserted. The dag-processor's `find_orm_dags` locks Dag 
rows without `skip_locked`, so a bundle sync that touches this Dag waits for 
it, and so does a `PATCH /dags`. `initializing_backfill_exists` already keeps 
the finalizer away until the runs commit, so `start_drain` could move after run 
creation, still inside the same `try`.



##########
airflow-core/src/airflow/ui/src/components/TriggerDag/PausedDagOptions.tsx:
##########
@@ -0,0 +1,112 @@
+/*!
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+import { Box, HStack, Text } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+
+import { useBackfillServiceListBackfillsUi, useDagRunServiceGetDagRuns } from 
"openapi/queries";
+import type { DagRunType } from "openapi/requests/types.gen";
+
+import { RadioCardItem, RadioCardRoot } from "src/system-components";
+
+import type { PausedDagAction } from "./types";
+
+const PAUSED_DAG_OPTIONS = [
+  { description: "pausedDag.unpauseDescription", label: "pausedDag.unpause", 
value: "unpause" },
+  { description: "pausedDag.drainDescription", label: "pausedDag.drain", 
value: "drain" },
+  { description: "pausedDag.keepPausedDescription", label: 
"pausedDag.keepPaused", value: "keepPaused" },
+] as const satisfies Array<{ description: string; label: string; value: 
PausedDagAction }>;
+
+const NON_BACKFILL_RUN_TYPES: Array<Exclude<DagRunType, "backfill">> = [
+  "scheduled",
+  "manual",
+  "operator_triggered",
+  "asset_triggered",
+  "asset_materialization",
+];
+
+type PausedDagOptionsProps = {
+  readonly dagId: string;
+  readonly onChange: (action: PausedDagAction) => void;
+  readonly value: PausedDagAction;
+};
+
+const PausedDagOptions = ({ dagId, onChange, value }: PausedDagOptionsProps) 
=> {
+  const { t: translate } = useTranslation("components");
+  // Unpausing or draining lets every unfinished run proceed, not only the one 
being created. The cached
+  // count can predate runs that have since finished (e.g. the run of an 
earlier drain), so it is
+  // refetched and only trusted once fetched for this form.
+  const { data: activeBackfills } = useBackfillServiceListBackfillsUi({ 
active: true, dagId }, undefined, {
+    staleTime: 0,
+  });
+  // A paused backfill's runs do not start while it stays paused, whichever 
option is chosen.
+  const isBackfillPaused = activeBackfills?.backfills.some((backfill) => 
backfill.is_paused) ?? false;
+  const { data: unfinishedRuns, isFetchedAfterMount } = 
useDagRunServiceGetDagRuns(

Review Comment:
   `runType` depends on `isBackfillPaused`, which is `false` until the 
backfills query returns, so the runs query fires twice and the first answer can 
count paused-backfill runs. Enabling it only once the backfills query has 
fetched avoids that.
   
   Separately, the title is a `Text` outside `RadioCardRoot`, so the radio 
group has no accessible label. `RunBackfillForm` uses `RadioCardLabel` inside 
the root.



##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_backfills.py:
##########
@@ -461,6 +462,67 @@ def test_create_backfill(self, repro_act, repro_exp, 
session, dag_maker, test_cl
         }
         check_last_log(session, dag_id="TEST_DAG_1", event="create_backfill", 
logical_date=None)
 
+    @pytest.mark.parametrize(
+        ("drain_dag", "expected_state"),
+        [
+            pytest.param(True, DagSchedulingState.DRAINING, id="drain"),
+            pytest.param(False, DagSchedulingState.PAUSED, id="leave-paused"),
+        ],
+    )
+    def test_create_backfill_on_paused_dag_with_drain_dag(
+        self, session, dag_maker, test_client, drain_dag, expected_state
+    ):
+        with dag_maker(session=session, dag_id="TEST_DAG_1", 
schedule="@daily") as dag:
+            EmptyOperator(task_id="mytask")
+        session.execute(update(DagModel).where(DagModel.dag_id == 
dag.dag_id).values(is_paused=True))
+        session.commit()
+
+        response = test_client.post(
+            url="/backfills",
+            json={
+                "dag_id": dag.dag_id,
+                "from_date": to_iso(pendulum.parse("2024-01-01")),
+                "to_date": to_iso(pendulum.parse("2024-01-03")),
+                "drain_dag": drain_dag,
+            },
+        )
+
+        assert response.status_code == 200
+        session.expire_all()
+        assert session.get(DagModel, dag.dag_id).scheduling_state == 
expected_state
+        assert session.scalar(
+            select(func.count()).select_from(DagRun).where(DagRun.backfill_id 
== response.json()["id"])
+        )
+
+    def test_create_backfill_with_drain_dag_requires_dag_edit_access(

Review Comment:
   Only `test_dag_run.py` has a no-drain case where edit access is denied and 
the request still succeeds. If `create_backfill` or `materialize_asset` checked 
PUT unconditionally, this test and its `test_assets.py` sibling would still 
pass.
   
   On the UI side, the passthrough test in `useTrigger.test.tsx` and 
`useCreateBackfill.test.tsx` only sends `true`, so a hard-coded `drain_dag: 
true` stays green. The `it.each` title for the `false` case also reads 
"refreshes ... only when drain_dag=false" when the body asserts it does not 
refresh.



-- 
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