mirror of
https://github.com/home-assistant/core.git
synced 2026-10-06 14:29:21 -04:00
Fix InfluxDB dropping a batch when a string field contains a line break (#182640)
This commit is contained in:
@@ -209,10 +209,19 @@ CONFIG_SCHEMA = probatio.Schema(
|
||||
)
|
||||
|
||||
|
||||
def _single_line(value: str) -> str:
|
||||
"""Return the value with its line breaks replaced by spaces.
|
||||
|
||||
Line protocol has no escape for a line break, in a string field it ends
|
||||
the point, and the escape the clients use for tags is not one either.
|
||||
"""
|
||||
return " ".join(value.splitlines())
|
||||
|
||||
|
||||
def _generate_event_to_json(conf: dict) -> Callable[[Event], dict[str, Any] | None]:
|
||||
"""Build event to json converter and add to config."""
|
||||
entity_filter = convert_include_exclude_filter(conf)
|
||||
tags = conf.get(CONF_TAGS)
|
||||
tags = {key: _single_line(value) for key, value in conf[CONF_TAGS].items()}
|
||||
tags_attributes: list[str] = conf[CONF_TAGS_ATTRIBUTES]
|
||||
default_measurement = conf.get(CONF_DEFAULT_MEASUREMENT)
|
||||
measurement_attr: str = conf[CONF_MEASUREMENT_ATTR]
|
||||
@@ -276,6 +285,9 @@ def _generate_event_to_json(conf: dict) -> Callable[[Event], dict[str, Any] | No
|
||||
else:
|
||||
include_uom = measurement_attr != "unit_of_measurement"
|
||||
|
||||
if isinstance(measurement, str):
|
||||
measurement = _single_line(measurement)
|
||||
|
||||
json: dict[str, Any] = {
|
||||
INFLUX_CONF_MEASUREMENT: measurement,
|
||||
INFLUX_CONF_TAGS: {
|
||||
@@ -286,7 +298,7 @@ def _generate_event_to_json(conf: dict) -> Callable[[Event], dict[str, Any] | No
|
||||
INFLUX_CONF_FIELDS: {},
|
||||
}
|
||||
if _include_state:
|
||||
json[INFLUX_CONF_FIELDS][INFLUX_CONF_STATE] = state.state
|
||||
json[INFLUX_CONF_FIELDS][INFLUX_CONF_STATE] = _single_line(state.state)
|
||||
if _include_value:
|
||||
json[INFLUX_CONF_FIELDS][INFLUX_CONF_VALUE] = _state_as_value
|
||||
|
||||
@@ -294,7 +306,9 @@ def _generate_event_to_json(conf: dict) -> Callable[[Event], dict[str, Any] | No
|
||||
ignore_attributes.update(global_ignore_attributes)
|
||||
for key, value in state.attributes.items():
|
||||
if key in tags_attributes:
|
||||
json[INFLUX_CONF_TAGS][key] = value
|
||||
json[INFLUX_CONF_TAGS][key] = (
|
||||
_single_line(value) if isinstance(value, str) else value
|
||||
)
|
||||
elif (
|
||||
(key != CONF_UNIT_OF_MEASUREMENT or include_uom)
|
||||
and (key != "device_class" or include_dc)
|
||||
@@ -311,7 +325,7 @@ def _generate_event_to_json(conf: dict) -> Callable[[Event], dict[str, Any] | No
|
||||
json[INFLUX_CONF_FIELDS][key] = float(value)
|
||||
except ValueError, TypeError:
|
||||
new_key = f"{key}_str"
|
||||
new_value = str(value)
|
||||
new_value = _single_line(str(value))
|
||||
json[INFLUX_CONF_FIELDS][new_key] = new_value
|
||||
|
||||
if RE_DIGIT_TAIL.match(new_value):
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
"""The tests for the InfluxDB component."""
|
||||
|
||||
from collections.abc import Generator
|
||||
from collections.abc import Callable, Generator
|
||||
from dataclasses import dataclass
|
||||
import datetime
|
||||
from http import HTTPStatus
|
||||
@@ -621,6 +621,77 @@ async def test_event_listener(
|
||||
write_api.reset_mock()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("hass_config", "mock_client", "config_ext", "get_write_api", "get_mock_call"),
|
||||
[
|
||||
(
|
||||
{
|
||||
"influxdb": {
|
||||
"override_measurement": "state\nlog",
|
||||
"tags_attributes": ["room"],
|
||||
"tags": {"site": "first\nfloor"},
|
||||
}
|
||||
},
|
||||
influxdb.DEFAULT_API_VERSION,
|
||||
BASE_V1_CONFIG,
|
||||
_get_write_api_mock_v1,
|
||||
influxdb.DEFAULT_API_VERSION,
|
||||
),
|
||||
(
|
||||
{
|
||||
"influxdb": {
|
||||
"override_measurement": "state\nlog",
|
||||
"tags_attributes": ["room"],
|
||||
"tags": {"site": "first\nfloor"},
|
||||
}
|
||||
},
|
||||
influxdb.API_VERSION_2,
|
||||
BASE_V2_CONFIG,
|
||||
_get_write_api_mock_v2,
|
||||
influxdb.API_VERSION_2,
|
||||
),
|
||||
],
|
||||
indirect=["mock_client", "get_mock_call"],
|
||||
)
|
||||
async def test_event_listener_multiline_strings(
|
||||
hass: HomeAssistant,
|
||||
mock_client: MagicMock,
|
||||
config_ext: dict[str, Any],
|
||||
get_write_api: Callable[[MagicMock], MagicMock],
|
||||
get_mock_call: Callable[..., Any],
|
||||
) -> None:
|
||||
"""Test line breaks in strings are replaced, line protocol has no escape for them."""
|
||||
await _setup(hass, mock_client, config_ext, get_write_api)
|
||||
|
||||
hass.states.async_set(
|
||||
"fake.entity_id",
|
||||
"Avenida de Logroño, 50\n28002 Madrid\r\nEspaña",
|
||||
{"address": "line one\nline two", "room": "living\nroom"},
|
||||
)
|
||||
await hass.async_block_till_done()
|
||||
await async_wait_for_queue_to_process(hass)
|
||||
|
||||
body = [
|
||||
{
|
||||
"measurement": "state log",
|
||||
"tags": {
|
||||
"domain": "fake",
|
||||
"entity_id": "entity_id",
|
||||
"room": "living room",
|
||||
"site": "first floor",
|
||||
},
|
||||
"time": ANY,
|
||||
"fields": {
|
||||
"state": "Avenida de Logroño, 50 28002 Madrid España",
|
||||
"address_str": "line one line two",
|
||||
},
|
||||
}
|
||||
]
|
||||
write_api = get_write_api(mock_client)
|
||||
assert write_api.call_count == 1
|
||||
assert write_api.call_args == get_mock_call(body)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("hass_config", "mock_client", "config_ext", "get_write_api", "get_mock_call"),
|
||||
[
|
||||
|
||||
Reference in New Issue
Block a user