mirror of
https://github.com/esphome/esphome.git
synced 2026-09-30 00:10:22 +00:00
[api] Collapse the three protobuf decode virtuals into one (#19016)
This commit is contained in:
@@ -249,7 +249,7 @@ static APIBuffer build_infrared_rf_transmit_wire() {
|
||||
std::memcpy(bytes + len, packed, packed_len);
|
||||
len += packed_len;
|
||||
// field 6: modulation = 1 (non-zero so it's actually emitted and exercises
|
||||
// decode_varint for this field, matching the documented layout above).
|
||||
// decode_field for this field, matching the documented layout above).
|
||||
put_byte(0x30);
|
||||
put_varint(1);
|
||||
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
esphome:
|
||||
name: api-decode-wire-types-test
|
||||
host:
|
||||
api:
|
||||
logger:
|
||||
level: DEBUG
|
||||
|
||||
switch:
|
||||
- platform: template
|
||||
name: "Wire Switch"
|
||||
optimistic: true
|
||||
|
||||
output:
|
||||
- platform: template
|
||||
id: wire_dim
|
||||
type: float
|
||||
write_action:
|
||||
- lambda: ""
|
||||
|
||||
light:
|
||||
- platform: monochromatic
|
||||
name: "Wire Light"
|
||||
output: wire_dim
|
||||
default_transition_length: 0s
|
||||
effects:
|
||||
- pulse:
|
||||
name: Pulse
|
||||
|
||||
text:
|
||||
- platform: template
|
||||
name: "Wire Text"
|
||||
optimistic: true
|
||||
mode: text
|
||||
min_length: 0
|
||||
max_length: 255
|
||||
|
||||
number:
|
||||
- platform: template
|
||||
name: "Wire Number"
|
||||
optimistic: true
|
||||
min_value: -1000
|
||||
max_value: 1000
|
||||
step: 0.5
|
||||
@@ -125,11 +125,12 @@ class RawApiClient:
|
||||
await self.read_until_frame(MESSAGE_TYPE_OF[api_pb2.HelloResponse])
|
||||
|
||||
async def send_message(self, msg: message.Message) -> None:
|
||||
await self.send_raw(MESSAGE_TYPE_OF[type(msg)], msg.SerializeToString())
|
||||
|
||||
async def send_raw(self, msg_type: int, payload: bytes) -> None:
|
||||
"""Send a frame with a hand built payload, for shapes protobuf will not serialize."""
|
||||
loop = asyncio.get_running_loop()
|
||||
await loop.sock_sendall(
|
||||
self._sock,
|
||||
encode_frame(MESSAGE_TYPE_OF[type(msg)], msg.SerializeToString()),
|
||||
)
|
||||
await loop.sock_sendall(self._sock, encode_frame(msg_type, payload))
|
||||
|
||||
async def read_until_frame(self, msg_type: int, timeout: float = 10.0) -> None:
|
||||
"""Read until at least one frame of msg_type has been received."""
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
"""decode_field() must take fields that match their declared wire type, drop the ones that do
|
||||
not, skip unknown fields, and handle two byte tags, varints and length prefixes."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Callable
|
||||
import struct
|
||||
|
||||
from aioesphomeapi import (
|
||||
EntityState,
|
||||
LightState,
|
||||
NumberState,
|
||||
SwitchState,
|
||||
TextState,
|
||||
api_pb2,
|
||||
)
|
||||
import pytest
|
||||
|
||||
from .raw_api_client import MESSAGE_TYPE_OF, RawApiClient, encode_varint
|
||||
from .state_utils import InitialStateHelper, StateWaiter, require_entity
|
||||
from .types import APIClientConnectedFactory, RunCompiledFunction
|
||||
|
||||
SWITCH_COMMAND = MESSAGE_TYPE_OF[api_pb2.SwitchCommandRequest]
|
||||
WIRE_VARINT, WIRE_LENGTH, WIRE_FIXED32 = 0, 2, 5
|
||||
|
||||
|
||||
def tag(field: int, wire_type: int) -> bytes:
|
||||
return encode_varint((field << 3) | wire_type)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_api_decode_wire_types(
|
||||
yaml_config: str,
|
||||
run_compiled: RunCompiledFunction,
|
||||
api_client_connected: APIClientConnectedFactory,
|
||||
unused_tcp_port: int,
|
||||
) -> None:
|
||||
async with (
|
||||
run_compiled(yaml_config),
|
||||
api_client_connected() as client,
|
||||
RawApiClient(unused_tcp_port) as raw,
|
||||
):
|
||||
entities, _ = await client.list_entities_services()
|
||||
switch = require_entity(entities, "wire_switch")
|
||||
light = require_entity(entities, "wire_light")
|
||||
text = require_entity(entities, "wire_text")
|
||||
number = require_entity(entities, "wire_number")
|
||||
key = tag(1, WIRE_FIXED32) + struct.pack("<I", switch.key)
|
||||
on, off = tag(2, WIRE_VARINT) + b"\x01", tag(2, WIRE_VARINT) + b"\x00"
|
||||
|
||||
switch_states: list[bool] = []
|
||||
waiter = StateWaiter()
|
||||
|
||||
def on_state(state: EntityState) -> None:
|
||||
if isinstance(state, SwitchState) and state.key == switch.key:
|
||||
switch_states.append(state.state)
|
||||
waiter.on_state(state)
|
||||
|
||||
def switch_is(value: bool) -> Callable[[EntityState], bool]:
|
||||
return lambda s: (
|
||||
isinstance(s, SwitchState) and s.key == switch.key and s.state is value
|
||||
)
|
||||
|
||||
def number_is(value: float) -> Callable[[EntityState], bool]:
|
||||
return lambda s: (
|
||||
isinstance(s, NumberState) and s.key == number.key and s.state == value
|
||||
)
|
||||
|
||||
initial = InitialStateHelper(entities)
|
||||
client.subscribe_states(initial.on_state_wrapper(on_state))
|
||||
await initial.wait_for_initial_states()
|
||||
await raw.connect()
|
||||
|
||||
# A well formed command: fixed32 key, varint state
|
||||
await raw.send_raw(SWITCH_COMMAND, key + on)
|
||||
await waiter.expect(switch_is(True))
|
||||
await raw.send_raw(SWITCH_COMMAND, key + off)
|
||||
await waiter.expect(switch_is(False))
|
||||
|
||||
# The same field with the wrong wire type is dropped, and a varint key never matches an
|
||||
# entity; each of these would turn the switch on if the payload were read as a varint
|
||||
seen = len(switch_states)
|
||||
await raw.send_raw(SWITCH_COMMAND, key + tag(2, WIRE_LENGTH) + b"\x01\x01")
|
||||
await raw.send_raw(
|
||||
SWITCH_COMMAND, key + tag(2, WIRE_FIXED32) + b"\x01\x00\x00\x00"
|
||||
)
|
||||
await raw.send_raw(
|
||||
SWITCH_COMMAND, tag(1, WIRE_VARINT) + encode_varint(switch.key) + on
|
||||
)
|
||||
# Ordered on the raw socket itself: this frame cannot be parsed before the bad ones, so
|
||||
# the only switch state since the marker must be the one it produces
|
||||
await raw.send_raw(SWITCH_COMMAND, key + on)
|
||||
await waiter.expect(switch_is(True), label="switch on after wrong wire types")
|
||||
assert switch_states[seen:] == [True]
|
||||
await raw.send_raw(SWITCH_COMMAND, key + off)
|
||||
await waiter.expect(switch_is(False))
|
||||
|
||||
# Truncated bodies stop the decode loop without taking the connection down: a tag with its
|
||||
# continuation bit set and nothing after it, a length prefix past the end of the payload,
|
||||
# and a fixed32 with two of its four bytes
|
||||
seen = len(switch_states)
|
||||
await raw.send_raw(SWITCH_COMMAND, key + b"\x80")
|
||||
await raw.send_raw(SWITCH_COMMAND, key + tag(2, WIRE_LENGTH) + b"\x7f" + b"ab")
|
||||
await raw.send_raw(SWITCH_COMMAND, tag(1, WIRE_FIXED32) + b"\x01\x02")
|
||||
await raw.send_raw(SWITCH_COMMAND, key + on)
|
||||
await waiter.expect(switch_is(True), label="switch on after truncated frames")
|
||||
assert switch_states[seen:] == [True]
|
||||
await raw.send_raw(SWITCH_COMMAND, key + off)
|
||||
await waiter.expect(switch_is(False))
|
||||
|
||||
# A negative number goes through the fixed32 float path of a normal client
|
||||
client.number_command(number.key, -77.5)
|
||||
await waiter.expect(number_is(-77.5))
|
||||
|
||||
# An unknown field ahead of the known ones is skipped; field 200 needs a two byte tag
|
||||
await raw.send_raw(
|
||||
SWITCH_COMMAND, tag(200, WIRE_VARINT) + encode_varint(300) + key + on
|
||||
)
|
||||
await waiter.expect(switch_is(True))
|
||||
|
||||
# Two byte tags (effect fields 18 and 19) and a two byte varint (300 ms transition)
|
||||
client.light_command(
|
||||
light.key, state=True, brightness=0.5, transition_length=0.3, effect="Pulse"
|
||||
)
|
||||
await waiter.expect(
|
||||
lambda s: (
|
||||
isinstance(s, LightState) and s.key == light.key and s.effect == "Pulse"
|
||||
)
|
||||
)
|
||||
client.light_command(light.key, effect="None", state=False)
|
||||
await waiter.expect(
|
||||
lambda s: isinstance(s, LightState) and s.key == light.key and not s.state
|
||||
)
|
||||
|
||||
# A string whose length prefix needs two varint bytes
|
||||
long_text = "w" * 200
|
||||
client.text_command(text.key, long_text)
|
||||
await waiter.expect(
|
||||
lambda s: (
|
||||
isinstance(s, TextState) and s.key == text.key and s.state == long_text
|
||||
)
|
||||
)
|
||||
@@ -18,7 +18,9 @@ sys.path.insert(0, str(Path(__file__).parents[4] / "script" / "api_protobuf"))
|
||||
import aioesphomeapi.api_options_pb2 as pb # noqa: E402
|
||||
from api_protobuf import ( # noqa: E402
|
||||
MAX_MESSAGE_ID,
|
||||
SOURCE_CLIENT,
|
||||
_make_ifdef_line,
|
||||
build_message_type,
|
||||
create_field_type_info,
|
||||
get_varint64_ifdef,
|
||||
validate_message_id,
|
||||
@@ -36,12 +38,15 @@ def _file_with_messages(
|
||||
file_desc = descriptor_pb2.FileDescriptorProto(name="test.proto")
|
||||
for name, field_type, deprecated in messages:
|
||||
msg = file_desc.message_type.add(name=name)
|
||||
field = msg.field.add(name="value", number=1, type=field_type)
|
||||
field = msg.field.add()
|
||||
field.CopyFrom(_field(field_type))
|
||||
field.options.deprecated = deprecated
|
||||
return file_desc
|
||||
|
||||
|
||||
UINT64 = descriptor_pb2.FieldDescriptorProto.TYPE_UINT64
|
||||
MESSAGE = descriptor_pb2.FieldDescriptorProto.TYPE_MESSAGE
|
||||
DOUBLE = descriptor_pb2.FieldDescriptorProto.TYPE_DOUBLE
|
||||
INT64 = descriptor_pb2.FieldDescriptorProto.TYPE_INT64
|
||||
SINT64 = descriptor_pb2.FieldDescriptorProto.TYPE_SINT64
|
||||
UINT32 = descriptor_pb2.FieldDescriptorProto.TYPE_UINT32
|
||||
@@ -182,3 +187,106 @@ def test_multi_byte_tag_fixed32_falls_back_to_the_generic_helper(
|
||||
content = _encode_field(field_type, number=16)
|
||||
assert "write_tag_and_fixed32" not in content, content
|
||||
assert content.startswith("pos = ProtoEncode::encode_"), content
|
||||
|
||||
|
||||
def _decode_case(field_type: int, number: int, *, repeated: bool = False) -> str:
|
||||
"""Return the decode_field() case the generator emits for one decoded field."""
|
||||
field = _field(field_type, number, repeated=repeated)
|
||||
if field_type == MESSAGE:
|
||||
field.type_name = ".Sub"
|
||||
return create_field_type_info(
|
||||
field, needs_decode=True, needs_encode=False
|
||||
).decode_content
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("field_type", "number", "wire_type", "accessor"),
|
||||
[
|
||||
(UINT32, 2, "WIRE_TYPE_VARINT", "value.as_varint()"),
|
||||
(BOOL, 3, "WIRE_TYPE_VARINT", "value.as_bool()"),
|
||||
(STRING, 1, "WIRE_TYPE_LENGTH_DELIMITED", "value.data()"),
|
||||
(FLOAT, 4, "WIRE_TYPE_FIXED32", "value.as_float()"),
|
||||
(FIXED32, 5, "WIRE_TYPE_FIXED32", "value.as_fixed32()"),
|
||||
],
|
||||
)
|
||||
def test_decode_cases_carry_field_number_and_wire_type(
|
||||
field_type: int, number: int, wire_type: str, accessor: str
|
||||
) -> None:
|
||||
"""Each decoded field yields one case keyed on its number and declared wire type."""
|
||||
case = _decode_case(field_type, number)
|
||||
lines = case.splitlines()
|
||||
assert lines[0] == f"case proto_tag({number}, {wire_type}):", case
|
||||
assert accessor in lines[1], case
|
||||
assert lines[-1].strip() == "break;", case
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("field_type", "repeated", "wire_type", "store"),
|
||||
[
|
||||
(UINT32, True, "WIRE_TYPE_VARINT", "this->value.push_back(value.as_varint());"),
|
||||
(
|
||||
STRING,
|
||||
True,
|
||||
"WIRE_TYPE_LENGTH_DELIMITED",
|
||||
"this->value.push_back(value.as_string());",
|
||||
),
|
||||
(
|
||||
MESSAGE,
|
||||
False,
|
||||
"WIRE_TYPE_LENGTH_DELIMITED",
|
||||
"value.decode_to_message(this->value);",
|
||||
),
|
||||
(
|
||||
MESSAGE,
|
||||
True,
|
||||
"WIRE_TYPE_LENGTH_DELIMITED",
|
||||
"value.decode_to_message(this->value.back());",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_repeated_and_message_fields_decode_through_the_same_case_shape(
|
||||
field_type: int, repeated: bool, wire_type: str, store: str
|
||||
) -> None:
|
||||
"""Repeated and sub message fields land in the one switch with their own store."""
|
||||
case = _decode_case(field_type, 7, repeated=repeated)
|
||||
lines = case.splitlines()
|
||||
assert lines[0] == f"case proto_tag(7, {wire_type}):", case
|
||||
assert store in case, case
|
||||
if field_type == MESSAGE and repeated:
|
||||
assert "this->value.emplace_back();" in case, case
|
||||
assert lines[-1].strip() == "break;", case
|
||||
|
||||
|
||||
def test_a_fixed64_field_fails_at_generation_time() -> None:
|
||||
"""The decode loop has no 64 bit wire type path, so such a field must never reach it silently."""
|
||||
desc = descriptor_pb2.DescriptorProto(name="Wide")
|
||||
desc.field.add(name="ratio", number=1, type=DOUBLE)
|
||||
with pytest.raises(
|
||||
ValueError, match="64-bit type 'double' .*ratio.* not supported"
|
||||
):
|
||||
build_message_type(desc, {}, {"Wide": SOURCE_CLIENT})
|
||||
|
||||
|
||||
def test_message_gets_a_single_decode_field_override() -> None:
|
||||
"""All wire types of a decoded message land in one decode_field() switch."""
|
||||
desc = descriptor_pb2.DescriptorProto(name="Mixed")
|
||||
desc.field.add(name="name", number=1, type=STRING)
|
||||
desc.field.add(name="count", number=2, type=UINT32)
|
||||
desc.field.add(name="level", number=3, type=FLOAT)
|
||||
header, cpp, _ = build_message_type(desc, {}, {"Mixed": SOURCE_CLIENT})
|
||||
decl = "void decode_field(uint32_t tag, const uint8_t *data, proto_varint_value_t scalar) override;"
|
||||
assert header.count(decl) == 1
|
||||
assert (
|
||||
cpp.count(
|
||||
"void Mixed::decode_field(uint32_t tag, const uint8_t *data, proto_varint_value_t scalar) {"
|
||||
)
|
||||
== 1
|
||||
)
|
||||
assert "switch (tag) {" in cpp
|
||||
assert "const ProtoFieldValue value(data, scalar);" in cpp
|
||||
for number, wire_type in (
|
||||
(1, "WIRE_TYPE_LENGTH_DELIMITED"),
|
||||
(2, "WIRE_TYPE_VARINT"),
|
||||
(3, "WIRE_TYPE_FIXED32"),
|
||||
):
|
||||
assert f"case proto_tag({number}, {wire_type}):" in cpp, cpp
|
||||
|
||||
Reference in New Issue
Block a user