From c3d9758062bee57ae2c3e97805dbcd7968e2c9c2 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 23 Jul 2026 15:42:29 -0700 Subject: [PATCH 1/5] Convert DockerOperator mounts to Mount objects after rendering mounts is a template field, rendered after __init__ runs. The constructor converted dict mounts into Mount objects and tagged each with template_fields so nested values would render. Mount is a dict subclass, so Jinja already renders its values natively without that tagging. Keep the raw input in __init__ and convert to Mount objects at the start of execute(), after rendering. related: #70296 Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../providers/docker/operators/docker.py | 11 ++++++---- .../unit/docker/operators/test_docker.py | 22 +++++++++++++------ .../validate_operators_init_exemptions.txt | 1 - 3 files changed, 22 insertions(+), 12 deletions(-) diff --git a/providers/docker/src/airflow/providers/docker/operators/docker.py b/providers/docker/src/airflow/providers/docker/operators/docker.py index e87f8d08e1664..e0af928e4744c 100644 --- a/providers/docker/src/airflow/providers/docker/operators/docker.py +++ b/providers/docker/src/airflow/providers/docker/operators/docker.py @@ -308,10 +308,10 @@ def __init__( self.mount_tmp_dir = mount_tmp_dir self.tmp_dir = tmp_dir self.user = user - mounts = [mount if isinstance(mount, Mount) else Mount(**mount) for mount in (mounts or [])] - self.mounts: list[Mount] = mounts - for mount in self.mounts: - mount.template_fields = ("Source", "Target", "Type") + # mounts is a template field; keep the raw input (dicts or Mount objects) here so Jinja + # renders it (Mount is a dict subclass, so its values render natively), and convert to + # Mount objects in execute(), after rendering. + self.mounts = mounts or [] self.entrypoint = entrypoint self.working_dir = working_dir self.xcom_all = xcom_all @@ -490,6 +490,9 @@ def _copy_from_docker(self, container_id, src): return lib.load(file) def execute(self, context: Context) -> list[str] | str | None: + # mounts is a template field held as raw input; convert to Mount objects now, after + # Jinja rendering has resolved their values. + self.mounts = [m if isinstance(m, Mount) else Mount(**m) for m in self.mounts] # Pull the docker image if `force_pull` is set or image does not exist locally if self.force_pull or not self.cli.images(name=self.image): self.log.info("::group::Pulling docker image %s", self.image) diff --git a/providers/docker/tests/unit/docker/operators/test_docker.py b/providers/docker/tests/unit/docker/operators/test_docker.py index 81fd757ddf765..c72d9d76ece12 100644 --- a/providers/docker/tests/unit/docker/operators/test_docker.py +++ b/providers/docker/tests/unit/docker/operators/test_docker.py @@ -874,12 +874,18 @@ def test_dict_mounts_are_normalized_to_mount_objects(self): Mount(target="/logs", source="logs", type="volume"), ], ) - assert all(isinstance(m, Mount) for m in op.mounts) - assert op.mounts[0]["Target"] == "/data" - assert op.mounts[0]["Source"] == "workspace" - assert op.mounts[0]["Type"] == "volume" - assert op.mounts[0]["ReadOnly"] is False - assert op.mounts[1]["Target"] == "/logs" + # mounts is a template field, so __init__ keeps the raw input; normalization to Mount + # objects happens in execute(), after rendering. + assert not isinstance(op.mounts[0], Mount) + + op.execute(None) + + passed_mounts = self.client_mock.create_host_config.call_args.kwargs["mounts"] + assert all(isinstance(m, Mount) for m in passed_mounts) + assert passed_mounts[0]["Target"] == "/data" + assert passed_mounts[0]["Source"] == "workspace" + assert passed_mounts[0]["ReadOnly"] is False + assert passed_mounts[1]["Target"] == "/logs" @pytest.mark.db_test def test_dict_mounts_are_templated(self, create_task_instance_of_operator): @@ -893,4 +899,6 @@ def test_dict_mounts_are_templated(self, create_task_instance_of_operator): ], ) rendered = ti.render_templates() - assert rendered.mounts[0]["Target"] == f"/{ti.run_id}" + # mounts stays a raw dict through rendering (Mount is a dict subclass; Jinja renders its + # values natively), so the templated value resolves before execute() converts it. + assert rendered.mounts[0]["target"] == f"/{ti.run_id}" diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt b/scripts/ci/prek/validate_operators_init_exemptions.txt index 29ee3dd263922..0890990686d76 100644 --- a/scripts/ci/prek/validate_operators_init_exemptions.txt +++ b/scripts/ci/prek/validate_operators_init_exemptions.txt @@ -22,7 +22,6 @@ providers/amazon/src/airflow/providers/amazon/aws/transfers/gcs_to_s3.py::GCSToS providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_redshift.py::S3ToRedshiftOperator providers/anthropic/src/airflow/providers/anthropic/operators/agent.py::AnthropicAgentSessionOperator providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py::KubernetesPodOperator -providers/docker/src/airflow/providers/docker/operators/docker.py::DockerOperator providers/google/src/airflow/providers/google/cloud/operators/bigquery.py::BigQueryInsertJobOperator providers/google/src/airflow/providers/google/cloud/operators/cloud_batch.py::CloudBatchSubmitJobOperator providers/google/src/airflow/providers/google/cloud/operators/cloud_build.py::CloudBuildCreateBuildOperator From 6483a5d178d2ccb0ffb5d784aed1b8ea80e3b6eb Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 23 Jul 2026 23:05:09 -0700 Subject: [PATCH 2/5] Tighten DockerOperator mounts comments Reduce the multi-line narration around deferring Mount conversion to a single concise note at each site. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../src/airflow/providers/docker/operators/docker.py | 8 +++----- .../docker/tests/unit/docker/operators/test_docker.py | 6 ++---- 2 files changed, 5 insertions(+), 9 deletions(-) diff --git a/providers/docker/src/airflow/providers/docker/operators/docker.py b/providers/docker/src/airflow/providers/docker/operators/docker.py index e0af928e4744c..b39f72a2d4377 100644 --- a/providers/docker/src/airflow/providers/docker/operators/docker.py +++ b/providers/docker/src/airflow/providers/docker/operators/docker.py @@ -308,9 +308,8 @@ def __init__( self.mount_tmp_dir = mount_tmp_dir self.tmp_dir = tmp_dir self.user = user - # mounts is a template field; keep the raw input (dicts or Mount objects) here so Jinja - # renders it (Mount is a dict subclass, so its values render natively), and convert to - # Mount objects in execute(), after rendering. + # mounts is a template field; keep the raw dicts/Mounts so Jinja renders them (Mount is a + # dict subclass), and convert to Mount objects in execute() after rendering. self.mounts = mounts or [] self.entrypoint = entrypoint self.working_dir = working_dir @@ -490,8 +489,7 @@ def _copy_from_docker(self, container_id, src): return lib.load(file) def execute(self, context: Context) -> list[str] | str | None: - # mounts is a template field held as raw input; convert to Mount objects now, after - # Jinja rendering has resolved their values. + # Convert the rendered mounts (raw dicts or Mounts) to Mount objects, now that rendering ran. self.mounts = [m if isinstance(m, Mount) else Mount(**m) for m in self.mounts] # Pull the docker image if `force_pull` is set or image does not exist locally if self.force_pull or not self.cli.images(name=self.image): diff --git a/providers/docker/tests/unit/docker/operators/test_docker.py b/providers/docker/tests/unit/docker/operators/test_docker.py index c72d9d76ece12..56e6589dd1e8d 100644 --- a/providers/docker/tests/unit/docker/operators/test_docker.py +++ b/providers/docker/tests/unit/docker/operators/test_docker.py @@ -874,8 +874,7 @@ def test_dict_mounts_are_normalized_to_mount_objects(self): Mount(target="/logs", source="logs", type="volume"), ], ) - # mounts is a template field, so __init__ keeps the raw input; normalization to Mount - # objects happens in execute(), after rendering. + # __init__ keeps the raw input; execute() normalizes to Mount objects. assert not isinstance(op.mounts[0], Mount) op.execute(None) @@ -899,6 +898,5 @@ def test_dict_mounts_are_templated(self, create_task_instance_of_operator): ], ) rendered = ti.render_templates() - # mounts stays a raw dict through rendering (Mount is a dict subclass; Jinja renders its - # values natively), so the templated value resolves before execute() converts it. + # Raw dict keeps its lowercase input keys through rendering; execute() converts to Mount. assert rendered.mounts[0]["target"] == f"/{ti.run_id}" From 5bcc5b4bd47c1a8988b4667998f29bfc2e34d4f7 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Fri, 24 Jul 2026 11:26:14 -0700 Subject: [PATCH 3/5] Drop narrating comments from DockerOperator Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../docker/src/airflow/providers/docker/operators/docker.py | 3 --- providers/docker/tests/unit/docker/operators/test_docker.py | 2 -- 2 files changed, 5 deletions(-) diff --git a/providers/docker/src/airflow/providers/docker/operators/docker.py b/providers/docker/src/airflow/providers/docker/operators/docker.py index b39f72a2d4377..fcb3732a0e7df 100644 --- a/providers/docker/src/airflow/providers/docker/operators/docker.py +++ b/providers/docker/src/airflow/providers/docker/operators/docker.py @@ -308,8 +308,6 @@ def __init__( self.mount_tmp_dir = mount_tmp_dir self.tmp_dir = tmp_dir self.user = user - # mounts is a template field; keep the raw dicts/Mounts so Jinja renders them (Mount is a - # dict subclass), and convert to Mount objects in execute() after rendering. self.mounts = mounts or [] self.entrypoint = entrypoint self.working_dir = working_dir @@ -489,7 +487,6 @@ def _copy_from_docker(self, container_id, src): return lib.load(file) def execute(self, context: Context) -> list[str] | str | None: - # Convert the rendered mounts (raw dicts or Mounts) to Mount objects, now that rendering ran. self.mounts = [m if isinstance(m, Mount) else Mount(**m) for m in self.mounts] # Pull the docker image if `force_pull` is set or image does not exist locally if self.force_pull or not self.cli.images(name=self.image): diff --git a/providers/docker/tests/unit/docker/operators/test_docker.py b/providers/docker/tests/unit/docker/operators/test_docker.py index 56e6589dd1e8d..abf53d8277681 100644 --- a/providers/docker/tests/unit/docker/operators/test_docker.py +++ b/providers/docker/tests/unit/docker/operators/test_docker.py @@ -874,7 +874,6 @@ def test_dict_mounts_are_normalized_to_mount_objects(self): Mount(target="/logs", source="logs", type="volume"), ], ) - # __init__ keeps the raw input; execute() normalizes to Mount objects. assert not isinstance(op.mounts[0], Mount) op.execute(None) @@ -898,5 +897,4 @@ def test_dict_mounts_are_templated(self, create_task_instance_of_operator): ], ) rendered = ti.render_templates() - # Raw dict keeps its lowercase input keys through rendering; execute() converts to Mount. assert rendered.mounts[0]["target"] == f"/{ti.run_id}" From 177d9ce79bf3bade61a70d36eaf85f8f0eec8b51 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Fri, 24 Jul 2026 14:28:59 -0700 Subject: [PATCH 4/5] Convert DockerSwarmOperator mounts at execute too DockerSwarmOperator overrides execute() and calls _run_service() without super().execute(), so the mount conversion moved out of __init__ never ran for it. Share the conversion via _normalize_mounts() and call it from both execute paths. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../providers/docker/operators/docker.py | 5 ++++- .../docker/operators/docker_swarm.py | 1 + .../docker/operators/test_docker_swarm.py | 19 +++++++++++++++++++ 3 files changed, 24 insertions(+), 1 deletion(-) diff --git a/providers/docker/src/airflow/providers/docker/operators/docker.py b/providers/docker/src/airflow/providers/docker/operators/docker.py index fcb3732a0e7df..ec4ac81d14aac 100644 --- a/providers/docker/src/airflow/providers/docker/operators/docker.py +++ b/providers/docker/src/airflow/providers/docker/operators/docker.py @@ -486,8 +486,11 @@ def _copy_from_docker(self, container_id, src): lib = getattr(self, "pickling_library", pickle) return lib.load(file) - def execute(self, context: Context) -> list[str] | str | None: + def _normalize_mounts(self) -> None: self.mounts = [m if isinstance(m, Mount) else Mount(**m) for m in self.mounts] + + def execute(self, context: Context) -> list[str] | str | None: + self._normalize_mounts() # Pull the docker image if `force_pull` is set or image does not exist locally if self.force_pull or not self.cli.images(name=self.image): self.log.info("::group::Pulling docker image %s", self.image) diff --git a/providers/docker/src/airflow/providers/docker/operators/docker_swarm.py b/providers/docker/src/airflow/providers/docker/operators/docker_swarm.py index 539cb1c2f7ad1..0e844f9323afe 100644 --- a/providers/docker/src/airflow/providers/docker/operators/docker_swarm.py +++ b/providers/docker/src/airflow/providers/docker/operators/docker_swarm.py @@ -174,6 +174,7 @@ def __init__( self.log_driver_config = None def execute(self, context: Context) -> None: + self._normalize_mounts() self.environment["AIRFLOW_TMP_DIR"] = self.tmp_dir return self._run_service() diff --git a/providers/docker/tests/unit/docker/operators/test_docker_swarm.py b/providers/docker/tests/unit/docker/operators/test_docker_swarm.py index 65e83d45ddd27..9f8fdabd91048 100644 --- a/providers/docker/tests/unit/docker/operators/test_docker_swarm.py +++ b/providers/docker/tests/unit/docker/operators/test_docker_swarm.py @@ -209,6 +209,25 @@ def test_no_auto_remove(self, types_mock, docker_api_client_patcher): "Docker service being removed even when `auto_remove` set to `never`" ) + @mock.patch("airflow.providers.docker.operators.docker_swarm.types") + def test_dict_mounts_converted_at_execute(self, types_mock, docker_api_client_patcher): + client_mock = mock.Mock(spec=APIClient) + client_mock.create_service.return_value = {"ID": "some_id"} + client_mock.images.return_value = [] + client_mock.pull.return_value = [b'{"status":"pull log"}'] + client_mock.tasks.return_value = [{"ServiceID": "some_id", "Status": {"State": "complete"}}] + docker_api_client_patcher.return_value = client_mock + + operator = DockerSwarmOperator( + image="", + task_id="unittest", + enable_logging=False, + mounts=[{"source": "/host", "target": "/container", "type": "bind"}], + ) + operator.execute(None) + + assert all(isinstance(m, types.Mount) for m in operator.mounts) + @pytest.mark.parametrize("status", ["failed", "shutdown", "rejected", "orphaned", "remove"]) @mock.patch("airflow.providers.docker.operators.docker_swarm.types") def test_non_complete_service_raises_error(self, types_mock, docker_api_client_patcher, status): From 4ae16f02b2d32167ea34da04d9825df083901085 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Sat, 25 Jul 2026 00:33:31 -0700 Subject: [PATCH 5/5] Leave rendered Mount dicts alone in DockerOperator._normalize_mounts A user-supplied Mount is a dict subclass, so the templater flattens it to a plain API-cased dict ({Target, Source, Type, ReadOnly}) during rendering. _normalize_mounts then hit the else branch and called Mount(**m) with those keys, raising TypeError at execute() before the container started. A rendered Mount is already the shape docker-py wants, so leave it (and any real Mount) alone; only rebuild a Mount from a raw user dict. Add a render-then-execute regression test that fails on the pre-fix source. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../providers/docker/operators/docker.py | 5 ++++- .../unit/docker/operators/test_docker.py | 20 +++++++++++++++++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/providers/docker/src/airflow/providers/docker/operators/docker.py b/providers/docker/src/airflow/providers/docker/operators/docker.py index ec4ac81d14aac..e20f2cade37a2 100644 --- a/providers/docker/src/airflow/providers/docker/operators/docker.py +++ b/providers/docker/src/airflow/providers/docker/operators/docker.py @@ -487,7 +487,10 @@ def _copy_from_docker(self, container_id, src): return lib.load(file) def _normalize_mounts(self) -> None: - self.mounts = [m if isinstance(m, Mount) else Mount(**m) for m in self.mounts] + # A user dict is rebuilt into a Mount, but rendering flattens a Mount into a + # plain dict with API-cased keys (Target/Source/Type) that docker-py already + # accepts, so leave those (and any real Mount) alone. + self.mounts = [m if isinstance(m, Mount) or "Target" in m else Mount(**m) for m in self.mounts] def execute(self, context: Context) -> list[str] | str | None: self._normalize_mounts() diff --git a/providers/docker/tests/unit/docker/operators/test_docker.py b/providers/docker/tests/unit/docker/operators/test_docker.py index abf53d8277681..324c593782718 100644 --- a/providers/docker/tests/unit/docker/operators/test_docker.py +++ b/providers/docker/tests/unit/docker/operators/test_docker.py @@ -898,3 +898,23 @@ def test_dict_mounts_are_templated(self, create_task_instance_of_operator): ) rendered = ti.render_templates() assert rendered.mounts[0]["target"] == f"/{ti.run_id}" + + @pytest.mark.db_test + def test_mount_objects_survive_render_then_execute(self, create_task_instance_of_operator): + # A user-supplied Mount is a dict subclass, so the templater flattens it to a + # plain API-cased dict during rendering. Drive the real render -> execute path + # to prove execute() forwards it to docker-py instead of raising on it. + ti = create_task_instance_of_operator( + operator_class=DockerOperator, + dag_id="test", + task_id="test", + image="test", + mount_tmp_dir=False, + mounts=[Mount(source="workspace", target="/{{task_instance.run_id}}", type="volume")], + ) + task = ti.render_templates() + task.execute(None) + + passed_mounts = self.client_mock.create_host_config.call_args.kwargs["mounts"] + assert passed_mounts[0]["Target"] == f"/{ti.run_id}" + assert passed_mounts[0]["Source"] == "workspace"