From 162089a8740ec541cbacfb4930cb25fc5fcd5136 Mon Sep 17 00:00:00 2001 From: Martin Hjelmare Date: Tue, 15 Sep 2026 13:44:53 +0200 Subject: [PATCH] Refactor command_line notify to async subprocess (#182282) --- .../components/command_line/notify.py | 87 ++++---- tests/components/command_line/test_notify.py | 196 +++++++++++++++--- 2 files changed, 213 insertions(+), 70 deletions(-) diff --git a/homeassistant/components/command_line/notify.py b/homeassistant/components/command_line/notify.py index fc40398da6d6..21871e015194 100644 --- a/homeassistant/components/command_line/notify.py +++ b/homeassistant/components/command_line/notify.py @@ -1,7 +1,7 @@ """Support for command line notification services.""" -import logging -import subprocess +import asyncio +from contextlib import suppress from typing import Any, override from homeassistant.components.notify import ( @@ -12,13 +12,10 @@ from homeassistant.const import CONF_COMMAND from homeassistant.core import HomeAssistant from homeassistant.exceptions import HomeAssistantError from homeassistant.helpers.typing import ConfigType, DiscoveryInfoType -from homeassistant.util.process import kill_subprocess from .const import CONF_COMMAND_TIMEOUT, DOMAIN, LOGGER from .utils import create_platform_yaml_not_supported_issue, render_template_args -_LOGGER = logging.getLogger(__name__) - async def async_get_service( hass: HomeAssistant, @@ -46,40 +43,58 @@ class CommandLineNotificationService(BaseNotificationService): self._timeout = timeout @override - def send_message(self, message: str = "", **kwargs: Any) -> None: + async def async_send_message(self, message: str = "", **kwargs: Any) -> None: """Send a message to a command line.""" if not (command := render_template_args(self.hass, self.command)): return LOGGER.debug("Running with message: %s", message) - with subprocess.Popen( # noqa: S602 # shell by design - command, - universal_newlines=True, - stdin=subprocess.PIPE, - close_fds=False, # required for posix_spawn - shell=True, - ) as proc: - try: - proc.communicate(input=message, timeout=self._timeout) - if proc.returncode != 0: - _LOGGER.error( - "Command failed (with return code %s): %s", - proc.returncode, - command, - ) - except subprocess.TimeoutExpired as err: - _LOGGER.debug("Timeout for command: %s", command) - kill_subprocess(proc) - raise HomeAssistantError( - translation_domain=DOMAIN, - translation_key="timeout_error", - translation_placeholders={"command": command}, - ) from err - except subprocess.SubprocessError as err: - _LOGGER.debug("Error trying to exec command: %s", command) - raise HomeAssistantError( - translation_domain=DOMAIN, - translation_key="command_error", - translation_placeholders={"command": command, "error": str(err)}, - ) from err + try: + proc = await asyncio.create_subprocess_shell( # shell by design + command, + stdin=asyncio.subprocess.PIPE, + close_fds=False, # required for posix_spawn + ) + except OSError as err: + LOGGER.debug("Error trying to exec command: %s", command) + raise HomeAssistantError( + translation_domain=DOMAIN, + translation_key="command_error", + translation_placeholders={"command": command, "error": str(err)}, + ) from err + + try: + async with asyncio.timeout(self._timeout): + await proc.communicate(input=message.encode()) + except TimeoutError as err: + LOGGER.debug("Timeout for command: %s", command) + with suppress(ProcessLookupError): + # The command may have exited between the timeout and the kill. + proc.kill() + if (stdin := proc.stdin) is not None and ( + not stdin.is_closing() or stdin.transport.get_write_buffer_size() + ): + # A still connected stdin pipe keeps proc.wait() pending forever, + # see https://bugs.python.org/issue43884. + stdin.transport.abort() + await proc.wait() + raise HomeAssistantError( + translation_domain=DOMAIN, + translation_key="timeout_error", + translation_placeholders={"command": command}, + ) from err + except asyncio.CancelledError: + # Kill synchronously so the child isn't orphaned; the event loop + # reaps it without awaiting wait(), which cancellation would + # interrupt anyway. + with suppress(ProcessLookupError): + proc.kill() + raise + + if proc.returncode != 0: + LOGGER.error( + "Command failed (with return code %s): %s", + proc.returncode, + command, + ) diff --git a/tests/components/command_line/test_notify.py b/tests/components/command_line/test_notify.py index 1fe070b216b8..571e7bfda85b 100644 --- a/tests/components/command_line/test_notify.py +++ b/tests/components/command_line/test_notify.py @@ -1,15 +1,16 @@ """The tests for the command line notification platform.""" +import asyncio import os from pathlib import Path -import subprocess import tempfile -from unittest.mock import patch +from unittest.mock import AsyncMock, MagicMock, patch import pytest from homeassistant import setup from homeassistant.components.command_line import DOMAIN +from homeassistant.components.command_line.notify import CommandLineNotificationService from homeassistant.components.notify import DOMAIN as NOTIFY_DOMAIN from homeassistant.core import HomeAssistant from homeassistant.exceptions import HomeAssistantError @@ -210,8 +211,8 @@ async def test_error_for_none_zero_exit_code( "command_line": [ { "notify": { - "command": "sleep 10000", - "command_timeout": 0.0000001, + "command": "sleep 5", + "command_timeout": 0, "name": "Test5", } } @@ -223,15 +224,55 @@ async def test_timeout( caplog: pytest.LogCaptureFixture, hass: HomeAssistant, load_yaml_integration: None ) -> None: """Test blocking is not forever.""" - with pytest.raises( - HomeAssistantError, match="Timeout trying to execute command: sleep 10000" - ): + with pytest.raises(HomeAssistantError) as exc_info: await hass.services.async_call( NOTIFY_DOMAIN, "test5", {"message": "error"}, blocking=True ) + assert exc_info.value.translation_key == "timeout_error" + assert exc_info.value.translation_placeholders == {"command": "sleep 5"} assert "Timeout" in caplog.text +@pytest.mark.parametrize( + "message", + [ + pytest.param("x" * 100000, id="stdin_buffer_below_high_water_mark"), + pytest.param("x" * 200000, id="stdin_buffer_above_high_water_mark"), + ], +) +@pytest.mark.parametrize( + "get_config", + [ + { + "command_line": [ + { + "notify": { + "command": "sleep 5", + "command_timeout": 1, + "name": "Test7", + } + } + ] + } + ], +) +@pytest.mark.usefixtures("load_yaml_integration") +async def test_timeout_with_unflushed_stdin(hass: HomeAssistant, message: str) -> None: + """Test a timeout is raised when the command never drains stdin. + + A message larger than the pipe buffer leaves data queued in the stdin + transport when the timeout cancels communicate(). The outer timeout keeps a + regression from hanging the test run instead of failing it. + """ + with pytest.raises(HomeAssistantError) as exc_info: + async with asyncio.timeout(3): + await hass.services.async_call( + NOTIFY_DOMAIN, "test7", {"message": message}, blocking=True + ) + assert exc_info.value.translation_key == "timeout_error" + assert exc_info.value.translation_placeholders == {"command": "sleep 5"} + + @pytest.mark.parametrize( "get_config", [ @@ -247,34 +288,121 @@ async def test_timeout( } ], ) -async def test_subprocess_exceptions( +async def test_spawn_error( caplog: pytest.LogCaptureFixture, hass: HomeAssistant, load_yaml_integration: None ) -> None: - """Test that notify subprocess exceptions are handled correctly.""" + """Test that a failure to spawn the command is handled correctly.""" - with patch( - "homeassistant.components.command_line.notify.subprocess.Popen" - ) as check_output: - check_output.return_value.__enter__ = check_output - check_output.return_value.communicate.side_effect = [ - subprocess.TimeoutExpired("cmd", 10), - None, - subprocess.SubprocessError(), - ] - with pytest.raises( - HomeAssistantError, match="Timeout trying to execute command: exit 0" - ): - await hass.services.async_call( - NOTIFY_DOMAIN, "test6", {"message": "error"}, blocking=True - ) - assert check_output.call_count == 2 - assert "Timeout for command" in caplog.text + with ( + patch( + "homeassistant.components.command_line.notify.asyncio.create_subprocess_shell", + side_effect=OSError("exec failed"), + ), + pytest.raises(HomeAssistantError) as exc_info, + ): + await hass.services.async_call( + NOTIFY_DOMAIN, "test6", {"message": "error"}, blocking=True + ) + assert exc_info.value.translation_key == "command_error" + assert exc_info.value.translation_placeholders == { + "command": "exit 0", + "error": "exec failed", + } + assert "Error trying to exec command" in caplog.text - with pytest.raises( - HomeAssistantError, match="Error trying to execute command: exit 0. Error: " - ): - await hass.services.async_call( - NOTIFY_DOMAIN, "test6", {"message": "error"}, blocking=True - ) - assert check_output.call_count == 4 - assert "Error trying to exec command" in caplog.text + +@pytest.mark.parametrize( + ("is_closing", "write_buffer_size", "expected_abort"), + [ + pytest.param(True, 0, False, id="stdin_already_closed"), + pytest.param(True, 4096, True, id="stdin_closing_with_queued_data"), + pytest.param(False, 0, True, id="stdin_still_connected"), + ], +) +@pytest.mark.parametrize( + "get_config", + [ + { + "command_line": [ + { + "notify": { + "command": "exit 0", + "name": "Test6", + } + } + ] + } + ], +) +@pytest.mark.usefixtures("load_yaml_integration") +async def test_timeout_cleanup( + caplog: pytest.LogCaptureFixture, + hass: HomeAssistant, + is_closing: bool, + write_buffer_size: int, + expected_abort: bool, +) -> None: + """Test the stdin pipe is only aborted while it can still block wait(). + + The command is assumed to have exited between the timeout and the kill, so + this also covers that kill() raises ProcessLookupError. + """ + mock_proc = AsyncMock() + mock_proc.communicate.side_effect = TimeoutError + mock_proc.kill = MagicMock(side_effect=ProcessLookupError) + mock_proc.stdin = MagicMock() + mock_proc.stdin.is_closing.return_value = is_closing + mock_proc.stdin.transport.get_write_buffer_size.return_value = write_buffer_size + + with ( + patch( + "homeassistant.components.command_line.notify.asyncio.create_subprocess_shell", + return_value=mock_proc, + ), + pytest.raises(HomeAssistantError) as exc_info, + ): + await hass.services.async_call( + NOTIFY_DOMAIN, "test6", {"message": "error"}, blocking=True + ) + assert exc_info.value.translation_key == "timeout_error" + assert exc_info.value.translation_placeholders == {"command": "exit 0"} + mock_proc.kill.assert_called_once() + mock_proc.wait.assert_awaited_once() + assert mock_proc.stdin.transport.abort.called is expected_abort + assert "Timeout for command" in caplog.text + + +@pytest.mark.parametrize( + "kill_side_effect", + [ + pytest.param(None, id="process_running"), + # The command may have exited before the kill. + pytest.param(ProcessLookupError, id="process_already_gone"), + ], +) +async def test_cancelled_kills_process( + hass: HomeAssistant, kill_side_effect: type[Exception] | None +) -> None: + """Test the subprocess is killed and the cancellation is re-raised. + + The event loop reaps the killed child on its own, so the cancellation path + does not await wait(). + """ + mock_proc = AsyncMock() + mock_proc.communicate.side_effect = asyncio.CancelledError + mock_proc.kill = MagicMock(side_effect=kill_side_effect) + + service = CommandLineNotificationService("exit 0", 15) + service.hass = hass + + with ( + patch( + "homeassistant.components.command_line.notify.asyncio.create_subprocess_shell", + return_value=mock_proc, + ), + pytest.raises(asyncio.CancelledError), + ): + await service.async_send_message("error") + + mock_proc.kill.assert_called_once() + mock_proc.wait.assert_not_awaited()