Skip to content

Commit 6e8119e

Browse files
committed
tests: integration: harden Forward worker assertions
Signed-off-by: Eduardo Silva <eduardo@chronosphere.io>
1 parent 97789a5 commit 6e8119e

1 file changed

Lines changed: 21 additions & 6 deletions

File tree

tests/integration/scenarios/in_forward/tests/test_in_forward_001.py

Lines changed: 21 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -750,11 +750,23 @@ def _send_forward_payload(service, payload, use_tls):
750750

751751
def _recv_msgpack_value(sock):
752752
sock.settimeout(5)
753-
data = sock.recv(4096)
754-
assert data
755-
value, offset = _unpack_obj(data, 0)
756-
assert offset == len(data)
757-
return value
753+
data = bytearray()
754+
755+
while True:
756+
chunk = sock.recv(4096)
757+
assert chunk
758+
data.extend(chunk)
759+
760+
try:
761+
value, offset = _unpack_obj(data, 0)
762+
except IndexError:
763+
continue
764+
765+
if offset > len(data):
766+
continue
767+
768+
assert offset == len(data)
769+
return value
758770

759771

760772
def _decode_str_like(raw):
@@ -1048,10 +1060,13 @@ def test_in_forward_workers_forward_mode_preserves_valid_prefix():
10481060
],
10491061
])
10501062
_send_tcp_payload(service.flb_listener_port, payload)
1051-
records = service.wait_for_record_count(1, timeout=20)
1063+
service.wait_for_log_message("event decoder or encoder failure", timeout=20)
1064+
service.wait_for_record_count(1, timeout=20)
10521065
finally:
10531066
service.stop()
10541067

1068+
records = service.flattened_records()
1069+
assert len(records) == 1
10551070
assert records[0]["message"] == "valid-prefix"
10561071

10571072

0 commit comments

Comments
 (0)