mirror of
https://github.com/home-assistant/core.git
synced 2026-09-27 18:08:37 -04:00
Fix timeout of endpoint does not crash Portainer update (#180405)
Co-authored-by: Josef Zweck <josef@zweck.dev>
This commit is contained in:
co-authored by
Josef Zweck
parent
7db136df88
commit
4c6fba2712
@@ -230,7 +230,6 @@ class PortainerCoordinator(
|
||||
try:
|
||||
endpoints = await self.portainer.get_endpoints()
|
||||
except PortainerAuthenticationError as err:
|
||||
_LOGGER.error("Authentication error: %s", repr(err))
|
||||
raise ConfigEntryAuthFailed(
|
||||
translation_domain=DOMAIN,
|
||||
translation_key="invalid_auth",
|
||||
@@ -251,184 +250,195 @@ class PortainerCoordinator(
|
||||
)
|
||||
continue
|
||||
|
||||
(
|
||||
containers,
|
||||
docker_version,
|
||||
docker_info,
|
||||
docker_system_df,
|
||||
volumes,
|
||||
) = await asyncio.gather(
|
||||
self.portainer.get_containers(endpoint.id),
|
||||
self.portainer.docker_version(endpoint.id),
|
||||
self.portainer.docker_info(endpoint.id),
|
||||
self.portainer.docker_system_df(endpoint.id, verbose=True),
|
||||
self.portainer.get_volumes(endpoint.id),
|
||||
)
|
||||
|
||||
stack_requests = [self.portainer.get_stacks(endpoint_id=endpoint.id)]
|
||||
swarm_id = (
|
||||
docker_info.swarm.cluster.get("ID")
|
||||
if docker_info.swarm
|
||||
and docker_info.swarm.control_available
|
||||
and docker_info.swarm.cluster
|
||||
else None
|
||||
)
|
||||
if swarm_id:
|
||||
stack_requests.append(
|
||||
self.portainer.get_stacks(
|
||||
endpoint_id=endpoint.id, swarm_id=swarm_id
|
||||
)
|
||||
try:
|
||||
(
|
||||
containers,
|
||||
docker_version,
|
||||
docker_info,
|
||||
docker_system_df,
|
||||
volumes,
|
||||
) = await asyncio.gather(
|
||||
self.portainer.get_containers(endpoint.id),
|
||||
self.portainer.docker_version(endpoint.id),
|
||||
self.portainer.docker_info(endpoint.id),
|
||||
self.portainer.docker_system_df(endpoint.id, verbose=True),
|
||||
self.portainer.get_volumes(endpoint.id),
|
||||
)
|
||||
|
||||
stacks = [
|
||||
stack
|
||||
for result in await asyncio.gather(*stack_requests)
|
||||
for stack in result
|
||||
]
|
||||
|
||||
prev_endpoint = self.data.get(endpoint.id) if self.data else None
|
||||
container_map: dict[str, PortainerContainerData] = {}
|
||||
stack_map: dict[str, PortainerStackData] = {
|
||||
stack.name: PortainerStackData(stack=stack, container_count=0)
|
||||
for stack in stacks
|
||||
}
|
||||
|
||||
container_names = [
|
||||
sanitize_container_name(container.names[0]) for container in containers
|
||||
]
|
||||
container_inspects = dict(
|
||||
zip(
|
||||
container_names,
|
||||
await asyncio.gather(
|
||||
*(
|
||||
self.portainer.inspect_container(endpoint.id, container.id)
|
||||
for container in containers
|
||||
)
|
||||
),
|
||||
strict=False,
|
||||
)
|
||||
)
|
||||
local_images = dict(
|
||||
zip(
|
||||
container_inspects,
|
||||
await asyncio.gather(
|
||||
*(
|
||||
self._get_local_image(
|
||||
endpoint.id, str(container_inspect.image)
|
||||
)
|
||||
for container_inspect in container_inspects.values()
|
||||
)
|
||||
),
|
||||
strict=False,
|
||||
)
|
||||
)
|
||||
|
||||
# Map containers, started and stopped
|
||||
for container in containers:
|
||||
container_name = sanitize_container_name(container.names[0])
|
||||
prev_container = (
|
||||
prev_endpoint.containers.get(container_name)
|
||||
if prev_endpoint
|
||||
stack_requests = [self.portainer.get_stacks(endpoint_id=endpoint.id)]
|
||||
swarm_id = (
|
||||
docker_info.swarm.cluster.get("ID")
|
||||
if docker_info.swarm
|
||||
and docker_info.swarm.control_available
|
||||
and docker_info.swarm.cluster
|
||||
else None
|
||||
)
|
||||
|
||||
container_inspect = container_inspects[container_name]
|
||||
local_image = local_images[container_name]
|
||||
|
||||
image_status = (
|
||||
(
|
||||
result.status
|
||||
if (
|
||||
result := self.watcher.results.get(
|
||||
(endpoint.id, container.id)
|
||||
)
|
||||
if swarm_id:
|
||||
stack_requests.append(
|
||||
self.portainer.get_stacks(
|
||||
endpoint_id=endpoint.id, swarm_id=swarm_id
|
||||
)
|
||||
else None
|
||||
)
|
||||
if self.watcher
|
||||
else None
|
||||
)
|
||||
|
||||
# Check if container belongs to a stack via docker compose label
|
||||
stack_name: str | None = (
|
||||
container.labels.get("com.docker.compose.project")
|
||||
or container.labels.get("com.docker.stack.namespace")
|
||||
if container.labels
|
||||
else None
|
||||
)
|
||||
if stack_name and (stack_data := stack_map.get(stack_name)):
|
||||
stack_data.container_count += 1
|
||||
stacks = [
|
||||
stack
|
||||
for result in await asyncio.gather(*stack_requests)
|
||||
for stack in result
|
||||
]
|
||||
|
||||
container_map[container_name] = PortainerContainerData(
|
||||
container=container,
|
||||
container_inspect=container_inspect,
|
||||
local_image=local_image,
|
||||
stats=None,
|
||||
stats_pre=prev_container.stats if prev_container else None,
|
||||
image_status=image_status,
|
||||
stack=stack_map[stack_name].stack
|
||||
if stack_name and stack_name in stack_map
|
||||
else None,
|
||||
)
|
||||
prev_endpoint = self.data.get(endpoint.id) if self.data else None
|
||||
container_map: dict[str, PortainerContainerData] = {}
|
||||
stack_map: dict[str, PortainerStackData] = {
|
||||
stack.name: PortainerStackData(stack=stack, container_count=0)
|
||||
for stack in stacks
|
||||
}
|
||||
|
||||
volume_usage_map = {
|
||||
item["Name"]: item
|
||||
for item in (docker_system_df.volume_disk_usage.items or [])
|
||||
}
|
||||
volume_map: dict[str, PortainerVolumeData] = {}
|
||||
for volume in volumes:
|
||||
if item := volume_usage_map.get(volume.name):
|
||||
volume.usage_data = DockerVolumeUsageData(
|
||||
size=item["UsageData"]["Size"],
|
||||
ref_count=item["UsageData"]["RefCount"],
|
||||
)
|
||||
volume_map[volume.name] = PortainerVolumeData(volume=volume)
|
||||
|
||||
# Separately fetch stats for active containers
|
||||
active_containers = [
|
||||
container
|
||||
for container in containers
|
||||
if container.state
|
||||
in (DockerContainerState.RUNNING, DockerContainerState.PAUSED)
|
||||
]
|
||||
if active_containers:
|
||||
container_stats = dict(
|
||||
container_names = [
|
||||
sanitize_container_name(container.names[0])
|
||||
for container in containers
|
||||
]
|
||||
container_inspects = dict(
|
||||
zip(
|
||||
(
|
||||
sanitize_container_name(container.names[0])
|
||||
for container in active_containers
|
||||
),
|
||||
container_names,
|
||||
await asyncio.gather(
|
||||
*(
|
||||
self.portainer.container_stats(
|
||||
endpoint_id=endpoint.id,
|
||||
container_id=container.id,
|
||||
self.portainer.inspect_container(
|
||||
endpoint.id, container.id
|
||||
)
|
||||
for container in active_containers
|
||||
for container in containers
|
||||
)
|
||||
),
|
||||
strict=False,
|
||||
)
|
||||
)
|
||||
local_images = dict(
|
||||
zip(
|
||||
container_inspects,
|
||||
await asyncio.gather(
|
||||
*(
|
||||
self._get_local_image(
|
||||
endpoint.id, str(container_inspect.image)
|
||||
)
|
||||
for container_inspect in container_inspects.values()
|
||||
)
|
||||
),
|
||||
strict=False,
|
||||
)
|
||||
)
|
||||
|
||||
# Now assign stats to the containers
|
||||
for container_name, stats in container_stats.items():
|
||||
container_map[container_name].stats = stats
|
||||
# Map containers, started and stopped
|
||||
for container in containers:
|
||||
container_name = sanitize_container_name(container.names[0])
|
||||
prev_container = (
|
||||
prev_endpoint.containers.get(container_name)
|
||||
if prev_endpoint
|
||||
else None
|
||||
)
|
||||
|
||||
self._container_ids_by_endpoint[endpoint.id] = {
|
||||
data.container.id: name for name, data in container_map.items()
|
||||
}
|
||||
container_inspect = container_inspects[container_name]
|
||||
local_image = local_images[container_name]
|
||||
|
||||
mapped_endpoints[endpoint.id] = PortainerCoordinatorData(
|
||||
id=endpoint.id,
|
||||
name=endpoint.name,
|
||||
endpoint=endpoint,
|
||||
containers=container_map,
|
||||
docker_version=docker_version,
|
||||
docker_info=docker_info,
|
||||
volumes=volume_map,
|
||||
stacks=stack_map,
|
||||
)
|
||||
image_status = (
|
||||
(
|
||||
result.status
|
||||
if (
|
||||
result := self.watcher.results.get(
|
||||
(endpoint.id, container.id)
|
||||
)
|
||||
)
|
||||
else None
|
||||
)
|
||||
if self.watcher
|
||||
else None
|
||||
)
|
||||
|
||||
# Check if container belongs to a stack via docker compose label
|
||||
stack_name: str | None = (
|
||||
container.labels.get("com.docker.compose.project")
|
||||
or container.labels.get("com.docker.stack.namespace")
|
||||
if container.labels
|
||||
else None
|
||||
)
|
||||
if stack_name and (stack_data := stack_map.get(stack_name)):
|
||||
stack_data.container_count += 1
|
||||
|
||||
container_map[container_name] = PortainerContainerData(
|
||||
container=container,
|
||||
container_inspect=container_inspect,
|
||||
local_image=local_image,
|
||||
stats=None,
|
||||
stats_pre=prev_container.stats if prev_container else None,
|
||||
image_status=image_status,
|
||||
stack=stack_map[stack_name].stack
|
||||
if stack_name and stack_name in stack_map
|
||||
else None,
|
||||
)
|
||||
|
||||
volume_usage_map = {
|
||||
item["Name"]: item
|
||||
for item in (docker_system_df.volume_disk_usage.items or [])
|
||||
}
|
||||
volume_map: dict[str, PortainerVolumeData] = {}
|
||||
for volume in volumes:
|
||||
if item := volume_usage_map.get(volume.name):
|
||||
volume.usage_data = DockerVolumeUsageData(
|
||||
size=item["UsageData"]["Size"],
|
||||
ref_count=item["UsageData"]["RefCount"],
|
||||
)
|
||||
volume_map[volume.name] = PortainerVolumeData(volume=volume)
|
||||
|
||||
# Separately fetch stats for active containers
|
||||
active_containers = [
|
||||
container
|
||||
for container in containers
|
||||
if container.state
|
||||
in (DockerContainerState.RUNNING, DockerContainerState.PAUSED)
|
||||
]
|
||||
if active_containers:
|
||||
container_stats = dict(
|
||||
zip(
|
||||
(
|
||||
sanitize_container_name(container.names[0])
|
||||
for container in active_containers
|
||||
),
|
||||
await asyncio.gather(
|
||||
*(
|
||||
self.portainer.container_stats(
|
||||
endpoint_id=endpoint.id,
|
||||
container_id=container.id,
|
||||
)
|
||||
for container in active_containers
|
||||
)
|
||||
),
|
||||
strict=False,
|
||||
)
|
||||
)
|
||||
|
||||
# Now assign stats to the containers
|
||||
for container_name, stats in container_stats.items():
|
||||
container_map[container_name].stats = stats
|
||||
|
||||
self._container_ids_by_endpoint[endpoint.id] = {
|
||||
data.container.id: name for name, data in container_map.items()
|
||||
}
|
||||
|
||||
mapped_endpoints[endpoint.id] = PortainerCoordinatorData(
|
||||
id=endpoint.id,
|
||||
name=endpoint.name,
|
||||
endpoint=endpoint,
|
||||
containers=container_map,
|
||||
docker_version=docker_version,
|
||||
docker_info=docker_info,
|
||||
volumes=volume_map,
|
||||
stacks=stack_map,
|
||||
)
|
||||
except PortainerTimeoutError:
|
||||
_LOGGER.warning(
|
||||
"Timed out fetching data for endpoint: %s (ID: %d). Skipping data fetch",
|
||||
endpoint.name,
|
||||
endpoint.id,
|
||||
)
|
||||
continue
|
||||
|
||||
self._async_add_remove_endpoints(mapped_endpoints)
|
||||
self._async_sync_event_listeners(mapped_endpoints)
|
||||
@@ -707,5 +717,15 @@ class PortainerDockerDiskSpaceCoordinator(
|
||||
for endpoint in endpoints:
|
||||
if endpoint.status == EndpointStatus.DOWN:
|
||||
continue
|
||||
results[endpoint.id] = await self.portainer.docker_system_df(endpoint.id)
|
||||
try:
|
||||
results[endpoint.id] = await self.portainer.docker_system_df(
|
||||
endpoint.id
|
||||
)
|
||||
except PortainerTimeoutError:
|
||||
_LOGGER.warning(
|
||||
"Timed out fetching DF data for endpoint: %s (ID: %d). Skipping data fetch",
|
||||
endpoint.name,
|
||||
endpoint.id,
|
||||
)
|
||||
continue
|
||||
return results
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
"""Tests for the Portainer binary sensor platform."""
|
||||
|
||||
from typing import Any, cast
|
||||
from unittest.mock import AsyncMock, patch
|
||||
|
||||
from freezegun.api import FrozenDateTimeFactory
|
||||
@@ -8,9 +9,12 @@ from pyportainer.exceptions import (
|
||||
PortainerConnectionError,
|
||||
PortainerTimeoutError,
|
||||
)
|
||||
from pyportainer.models.docker import EndpointStatus
|
||||
from pyportainer.models.portainer import Endpoint
|
||||
import pytest
|
||||
from syrupy.assertion import SnapshotAssertion
|
||||
|
||||
from homeassistant.components.portainer.const import DOMAIN
|
||||
from homeassistant.components.portainer.coordinator import DEFAULT_SCAN_INTERVAL
|
||||
from homeassistant.config_entries import ConfigEntryState
|
||||
from homeassistant.const import STATE_UNAVAILABLE, Platform
|
||||
@@ -20,7 +24,12 @@ from homeassistant.util import dt as dt_util
|
||||
|
||||
from . import setup_integration
|
||||
|
||||
from tests.common import MockConfigEntry, async_fire_time_changed, snapshot_platform
|
||||
from tests.common import (
|
||||
MockConfigEntry,
|
||||
async_fire_time_changed,
|
||||
async_load_json_array_fixture,
|
||||
snapshot_platform,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
@@ -102,3 +111,45 @@ async def test_refresh_containers_exceptions(
|
||||
|
||||
assert (state := hass.states.get("binary_sensor.practical_morse_status"))
|
||||
assert state.state == STATE_UNAVAILABLE
|
||||
|
||||
|
||||
async def test_endpoint_timeout_only_marks_that_endpoint_unavailable(
|
||||
hass: HomeAssistant,
|
||||
mock_portainer_client: AsyncMock,
|
||||
mock_config_entry: MockConfigEntry,
|
||||
freezer: FrozenDateTimeFactory,
|
||||
) -> None:
|
||||
"""Test a timeout fetching one endpoint's data doesn't affect other endpoints."""
|
||||
endpoints = cast(
|
||||
list[dict[str, Any]],
|
||||
await async_load_json_array_fixture(hass, "endpoints.json", DOMAIN),
|
||||
)
|
||||
for endpoint in endpoints:
|
||||
endpoint["Status"] = EndpointStatus.UP
|
||||
mock_portainer_client.get_endpoints.return_value = [
|
||||
Endpoint.from_dict(endpoint) for endpoint in endpoints
|
||||
]
|
||||
|
||||
await setup_integration(hass, mock_config_entry)
|
||||
assert mock_config_entry.state is ConfigEntryState.LOADED
|
||||
assert (state := hass.states.get("binary_sensor.my_edge_offline_status"))
|
||||
assert state.state != STATE_UNAVAILABLE
|
||||
|
||||
docker_version = mock_portainer_client.docker_version.return_value
|
||||
|
||||
async def _docker_version(endpoint_id: int) -> Any:
|
||||
if endpoint_id == 42:
|
||||
raise PortainerTimeoutError("timeout")
|
||||
return docker_version
|
||||
|
||||
mock_portainer_client.docker_version.side_effect = _docker_version
|
||||
|
||||
freezer.tick(DEFAULT_SCAN_INTERVAL)
|
||||
async_fire_time_changed(hass, dt_util.utcnow())
|
||||
await hass.async_block_till_done(wait_background_tasks=True)
|
||||
|
||||
assert mock_config_entry.state is ConfigEntryState.LOADED
|
||||
assert (state := hass.states.get("binary_sensor.my_edge_offline_status"))
|
||||
assert state.state == STATE_UNAVAILABLE
|
||||
assert (state := hass.states.get("binary_sensor.my_environment_status"))
|
||||
assert state.state != STATE_UNAVAILABLE
|
||||
|
||||
@@ -1,17 +1,30 @@
|
||||
"""Tests for the Portainer sensor platform."""
|
||||
|
||||
from unittest.mock import patch
|
||||
from typing import Any, cast
|
||||
from unittest.mock import AsyncMock, patch
|
||||
|
||||
from freezegun.api import FrozenDateTimeFactory
|
||||
from pyportainer.exceptions import PortainerTimeoutError
|
||||
from pyportainer.models.docker import EndpointStatus
|
||||
from pyportainer.models.portainer import Endpoint
|
||||
import pytest
|
||||
from syrupy.assertion import SnapshotAssertion
|
||||
|
||||
from homeassistant.const import Platform
|
||||
from homeassistant.components.portainer.const import DOMAIN
|
||||
from homeassistant.components.portainer.coordinator import DEFAULT_DF_SCAN_INTERVAL
|
||||
from homeassistant.const import STATE_UNAVAILABLE, Platform
|
||||
from homeassistant.core import HomeAssistant
|
||||
from homeassistant.helpers import entity_registry as er
|
||||
from homeassistant.util import dt as dt_util
|
||||
|
||||
from . import setup_integration
|
||||
|
||||
from tests.common import MockConfigEntry, snapshot_platform
|
||||
from tests.common import (
|
||||
MockConfigEntry,
|
||||
async_fire_time_changed,
|
||||
async_load_json_array_fixture,
|
||||
snapshot_platform,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
@@ -38,3 +51,53 @@ async def test_all_entities(
|
||||
snapshot,
|
||||
mock_config_entry.entry_id,
|
||||
)
|
||||
|
||||
|
||||
async def test_df_endpoint_timeout_only_marks_that_endpoint_unavailable(
|
||||
hass: HomeAssistant,
|
||||
mock_portainer_client: AsyncMock,
|
||||
mock_config_entry: MockConfigEntry,
|
||||
freezer: FrozenDateTimeFactory,
|
||||
) -> None:
|
||||
"""Test a timeout fetching one endpoint's disk usage doesn't affect other endpoints."""
|
||||
endpoints = cast(
|
||||
list[dict[str, Any]],
|
||||
await async_load_json_array_fixture(hass, "endpoints.json", DOMAIN),
|
||||
)
|
||||
for endpoint in endpoints:
|
||||
endpoint["Status"] = EndpointStatus.UP
|
||||
mock_portainer_client.get_endpoints.return_value = [
|
||||
Endpoint.from_dict(endpoint) for endpoint in endpoints
|
||||
]
|
||||
|
||||
await setup_integration(hass, mock_config_entry)
|
||||
assert (
|
||||
state := hass.states.get("sensor.my_environment_image_disk_usage_total_size")
|
||||
)
|
||||
assert state.state != STATE_UNAVAILABLE
|
||||
assert (
|
||||
state := hass.states.get("sensor.my_edge_offline_image_disk_usage_total_size")
|
||||
)
|
||||
assert state.state != STATE_UNAVAILABLE
|
||||
|
||||
docker_system_df = mock_portainer_client.docker_system_df.return_value
|
||||
|
||||
async def _docker_system_df(endpoint_id: int) -> Any:
|
||||
if endpoint_id == 42:
|
||||
raise PortainerTimeoutError("timeout")
|
||||
return docker_system_df
|
||||
|
||||
mock_portainer_client.docker_system_df.side_effect = _docker_system_df
|
||||
|
||||
freezer.tick(DEFAULT_DF_SCAN_INTERVAL)
|
||||
async_fire_time_changed(hass, dt_util.utcnow())
|
||||
await hass.async_block_till_done(wait_background_tasks=True)
|
||||
|
||||
assert (
|
||||
state := hass.states.get("sensor.my_edge_offline_image_disk_usage_total_size")
|
||||
)
|
||||
assert state.state == STATE_UNAVAILABLE
|
||||
assert (
|
||||
state := hass.states.get("sensor.my_environment_image_disk_usage_total_size")
|
||||
)
|
||||
assert state.state != STATE_UNAVAILABLE
|
||||
|
||||
Reference in New Issue
Block a user