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]