Skip to content

Fix Beam async hook launching pipelines through a shell - #72510

Merged
potiuk merged 3 commits into
apache:mainfrom
r3wretrhy:fix-beam-async-no-shell
Oct 4, 2026
Merged

potiuk merged 3 commits into
apache:mainfrom
r3wretrhy:fix-beam-async-no-shell

Conversation

@r3wretrhy

Copy link
Copy Markdown
Contributor

BeamAsyncHook started pipeline processes by joining the argv list with POSIX shlex quoting and running the result through a shell. The sync hook already uses subprocess.Popen with shell=False and the original list. The async path could therefore split interpreter or pipeline paths that contain spaces, especially on Windows where POSIX quoting is not what the shell expects.

This runs _beam_version and run_beam_command_async with asyncio.create_subprocess_exec so each argument stays one argv entry. Missing interpreters still raise AirflowException. Tests assert exec is used (not shell) and that an argument containing spaces survives a live process.


Was generative AI tooling used to co-author this PR?
  • Yes (Grok)

Generated-by: Grok following the guidelines


  • Tests run locally: python verification of run_beam_command_async with a spaced argument, exec-vs-shell construction, and _beam_version error wrapping. Full provider suite via CI: providers/apache/beam/tests/unit/apache/beam/hooks/test_beam.py.

@boring-cyborg

boring-cyborg Bot commented Sep 4, 2026

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contributors' Guide
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our prek-hooks will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example Dag that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@potiuk

potiuk commented Sep 8, 2026

Copy link
Copy Markdown
Member

this must wait for #66952

@r3wretrhy
r3wretrhy force-pushed the fix-beam-async-no-shell branch from dff9c46 to 145d86c Compare September 9, 2026 00:59
@r3wretrhy

Copy link
Copy Markdown
Contributor Author

Hi @potiuk — #66952 has merged. Is there anything else needed before this PR can proceed, or can review continue from here? Thanks.

@r3wretrhy
r3wretrhy force-pushed the fix-beam-async-no-shell branch from 145d86c to db09881 Compare September 21, 2026 14:49
r3wretrhy and others added 2 commits October 4, 2026 01:23
The test passes on main too, because shlex.quote plus sh already keeps an argument with spaces together, and it spawns a real interpreter. test_run_beam_command_async_uses_exec_with_argv already covers the switch to create_subprocess_exec.

Generated-by: Claude Opus 5
@potiuk
potiuk force-pushed the fix-beam-async-no-shell branch from db09881 to bfba8a0 Compare October 3, 2026 23:24
asyncio's Process.returncode is typed int | None, so assigning it to the int inferred from the OSError branch failed mypy-providers.

Generated-by: Claude Opus 5

@potiuk potiuk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for waiting on #66952, and for the fix. I've rebased the branch onto current main (no conflicts) and pushed small follow-up commits on top of yours:

  • Dropped test_run_beam_command_async_preserves_arguments_with_spaces. It also passes on main on Linux/macOS, because shlex.quote plus sh already keeps "hello world" together, so it doesn't exercise the change, and it starts a real interpreter. test_run_beam_command_async_uses_exec_with_argv already covers the behaviour.
  • Fixed a mypy error in _beam_version (dcda56fba2): asyncio's Process.returncode is typed int | None, so returncode needed that annotation; CI's mypy-providers caught it once the branch was on current main.

The change itself looks right to me. create_subprocess_exec(*cmd, cwd=working_directory) passes argv through unchanged, which matches what the sync run_beam_command already does with Popen(cmd, shell=False). The environment is inherited in both cases, and the logged command line is unchanged. A missing interpreter in run_beam_command_async now raises FileNotFoundError, which matches the sync hook, and the triggers already turn any exception into an error event.

Merging once CI is green.


Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting

@potiuk
potiuk merged commit b2c215c into apache:main Oct 4, 2026
80 checks passed
@boring-cyborg

boring-cyborg Bot commented Oct 4, 2026

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants