Files
NODEDC_MISSION_CORE/tests/test_perception_streaming_ingress.py

339 lines
13 KiB
Python

"""Bounded synthetic IPC fixtures, no model or local load tests."""
import json
import socket
import threading
from contextlib import contextmanager
from dataclasses import replace
import pytest
from k1link.compute.live_perception import LiveIngressEvent, LivePerceptionIngress
from k1link.perception import streaming_wire as wire
from k1link.perception.graph_contracts import GraphState
from k1link.perception.realtime_contract import StreamStart
from k1link.perception.streaming_ingress import StreamingIngress
from k1link.perception.streaming_lifecycle import StreamingLifecycle
from k1link.perception.streaming_queue import StreamMailbox
from k1link.perception.streaming_sender import StreamingSender
def start():
return StreamStart(
run_id="test-run",
source_id="recorded-or-live",
worker_id="worker-006",
epoch_id="epoch-1",
lease_generation=1,
profile_sha256="a" * 64,
image_sha256="b" * 64,
effective_config_sha256="c" * 64,
calibration_sha256="d" * 64,
clock_domain_id="original-host-clock",
input_mode="live",
)
def event(seq=1, *, modality="lidar", payload=b"points", source_sequence=None):
return LiveIngressEvent(
ingress_sequence=seq,
session_id="capture-1",
session_generation=1,
modality=modality,
source_id=f"raw/{modality}",
source_sequence=seq if source_sequence is None else source_sequence,
captured_at_epoch_ns=1_799_999_999_123_456_789,
received_monotonic_ns=9_007_199_254_740_993 + seq,
payload=payload,
)
def send_event(sock, observation, *, identity=None):
for header, payload in wire.event_packets(identity or start(), observation):
sock.sendall(header)
sock.sendall(payload)
@contextmanager
def harness(tmp_path, *, consume=None, byte_limit=16 * 1024 * 1024, clock=None, **kwargs):
mailbox = StreamMailbox(byte_limit=byte_limit)
run = StreamingLifecycle(
start(),
tmp_path,
mailbox,
threading.Event(),
**({"clock_ns": clock} if clock else {}),
)
left, right = socket.socketpair()
left.settimeout(1)
received, notices = [], []
run.ready()
receiver = StreamingIngress(
right, run, "capture-1", 1, consume or received.append, notices.append, **kwargs
)
receiver.start()
try:
yield left, receiver, run, received, notices
finally:
left.close()
run.request_stop("cancelled")
assert receiver.join()
assert run.close()
def hello(sock):
sock.sendall(wire.open_packet(start(), "capture-1", 1))
def test_unknown_duration_delivers_before_end_and_preserves_every_raw_byte(tmp_path):
delivered = threading.Event()
seen = []
def consume(value):
seen.append(value)
delivered.set()
with harness(tmp_path, consume=consume) as (sock, receiver, run, _, _):
hello(sock)
original = event(payload=b"x" * (wire.MAX_FRAGMENT + 19))
send_event(sock, original)
assert delivered.wait(1)
assert receiver.thread.is_alive() and not run.mailbox.done
assert seen == [original] # Includes int64 values above JS exact-number range.
sock.sendall(wire.terminal_packet(start()))
assert receiver.join()
assert receiver.terminal == "end" and receiver.error is None
assert receiver.counts["fragments"] == 2
assert receiver.snapshot()["incomplete_observations"] == 0
assert run.mailbox.bytes == 0
assert (
run.mailbox.peak_bytes
<= 2 * len(original.payload) + 3 * wire.MAX_HEADER + wire.PREFIX.size
)
@pytest.mark.parametrize("value", [0, 1, (1 << 53) + 1, (1 << 64) - 1])
def test_uint64_roundtrip(value):
assert wire.uint64(wire.decimal(value)) == value
@pytest.mark.parametrize("value", [1, True, "-1", "01", "1.0", str(1 << 64), "1e6"])
def test_uint64_rejects_lossy_or_ambiguous_projection(value):
with pytest.raises(wire.StreamWireError):
wire.uint64(value)
def test_duplicate_json_fields_and_unbounded_headers_fail_before_body(tmp_path):
with pytest.raises(wire.StreamWireError):
wire.parse_header(bytearray(b'{"binding":"a","binding":"b"}'))
with harness(tmp_path) as (sock, receiver, run, _, _):
sock.sendall(wire.PREFIX.pack(wire.MAGIC, wire.OPEN, wire.MAX_HEADER + 1))
assert receiver.join()
assert "header size" in receiver.error
assert run.state == GraphState.RUNNING # Unbound peer cannot stop owner.
assert run.mailbox.bytes == 0
@pytest.mark.parametrize(
"field,value", [("epoch_id", "old"), ("lease_generation", 2), ("image_sha256", "f" * 64)]
)
def test_wrong_handshake_does_not_cancel_current_owner(tmp_path, field, value):
with harness(tmp_path) as (sock, receiver, run, _, _):
sock.sendall(wire.open_packet(replace(start(), **{field: value}), "capture-1", 1))
assert receiver.join()
assert not receiver.opened and run.state == GraphState.RUNNING
run.check_current(start())
@pytest.mark.parametrize(
"mutation,expected",
[
(lambda h: h.update(binding="e" * 64), "another StreamStart"),
(lambda h: h.update(total_bytes=2 * wire.MAX_FRAGMENT + 1), "length"),
(lambda h: h.update(fragment_bytes=True), "length"),
(lambda h: h.update(offset=1), "length"),
(lambda h: h.update(payload_sha256="e" * 64), "observation integrity"),
(lambda h: h.update(fragment_sha256="e" * 64), "fragment integrity"),
(lambda h: h["event"].update(session_generation="2"), "acquisition"),
(lambda h: h["event"].update(received_monotonic_ns=9007199254740993), "decimal string"),
],
)
def test_invalid_fragments_fail_closed_without_delivery(tmp_path, mutation, expected):
with harness(tmp_path) as (sock, receiver, run, received, _):
hello(sock)
header, payload = next(wire.event_packets(start(), event()))
value = json.loads(header[wire.PREFIX.size :])
mutation(value)
sock.sendall(wire.packet(wire.FRAGMENT, value) + payload)
assert receiver.join()
assert expected in receiver.error and not received
assert run.state == GraphState.STOPPING and run.mailbox.bytes == 0
@pytest.mark.parametrize("case", ["truncated", "interleaved", "end", "timeout"])
def test_incomplete_reassembly_has_terminal_accounting_and_releases_bytes(tmp_path, case):
with harness(tmp_path, fragment_timeout=0.05) as (sock, receiver, run, received, _):
hello(sock)
pieces = list(wire.event_packets(start(), event(payload=b"x" * (wire.MAX_FRAGMENT + 1))))
sock.sendall(pieces[0][0])
sock.sendall(pieces[0][1])
if case == "truncated":
sock.shutdown(socket.SHUT_WR)
elif case == "interleaved":
send_event(sock, event(seq=2))
elif case == "end":
sock.sendall(wire.terminal_packet(start()))
assert receiver.join()
assert receiver.terminal == "failed" and not received
assert receiver.snapshot()["incomplete_observations"] == 1
assert run.mailbox.bytes == 0 and run.mailbox.quiescent
def test_raw_and_decoded_work_share_one_budget(tmp_path):
# Small bound proves rejection without creating large local pressure.
with harness(tmp_path, byte_limit=3 * wire.MAX_HEADER + 100) as (
sock,
receiver,
run,
received,
_,
):
active = {"sequence": 1, "payload_bytes": 85}
assert run.admit(start(), active)
assert run.mailbox.take() is active
hello(sock)
send_event(sock, event()) # Needs 12 extra bytes, only 6 remain.
assert receiver.join()
assert "byte budget" in receiver.error and not received
assert run.mailbox.bytes == 85 and not run.mailbox.quiescent
run.mailbox.release(active)
def test_stop_cannot_release_borrowed_payload_until_callback_exits(tmp_path):
entered, leave = threading.Event(), threading.Event()
def consume(_value):
entered.set()
assert leave.wait(2)
with harness(tmp_path, consume=consume) as (sock, receiver, run, _, _):
hello(sock)
send_event(sock, event())
assert entered.wait(1)
try:
assert not run.close("cancelled")
assert run.mailbox.bytes > 0 and not run.lease.released
finally:
leave.set()
assert receiver.join() and run.mailbox.bytes == 0
assert run.close() and run.lease.released
def test_lease_expiry_wakes_idle_receiver_without_new_bytes(tmp_path):
now = [100]
with harness(tmp_path, clock=lambda: now[0]) as (sock, receiver, run, _, _):
hello(sock)
now[0] += 3_000_000_000
assert receiver.join()
assert run.stop_event.is_set() and run.reason == "lease-lost"
assert run.mailbox.bytes == 0
def test_gaps_reset_camera_codec_and_cancel_is_not_successful_end(tmp_path):
with harness(tmp_path) as (sock, receiver, run, received, notices):
hello(sock)
send_event(sock, event(1, modality="camera-init"))
send_event(sock, event(2, modality="camera-frame"))
sock.sendall(
wire.gap_packet(start(), modality="camera-frame", reason="source-gap", count=1)
)
send_event(sock, event(4, modality="camera-init"))
send_event(sock, event(5, modality="camera-frame"))
sock.sendall(wire.terminal_packet(start(), cancel=True))
assert receiver.join()
assert receiver.terminal == "cancelled" and len(received) == 4
assert len(notices) == 1 and receiver.counts["declared_gap_observations"] == 1
assert receiver.counts["ingress_sequence_gaps"] == 1
assert run.reason == "cancelled"
@pytest.mark.parametrize("case", ["no-init", "duplicate", "camera-gap", "clock-regression"])
def test_source_order_and_codec_prerequisites_are_enforced(tmp_path, case):
with harness(tmp_path) as (sock, receiver, _, received, _):
hello(sock)
if case != "no-init":
send_event(sock, event(1, modality="camera-init"))
send_event(sock, event(2, modality="camera-frame"))
bad = event(
2 if case == "duplicate" else 3,
modality="camera-frame",
source_sequence=4 if case == "camera-gap" else None,
)
if case == "clock-regression":
bad = replace(bad, received_monotonic_ns=1)
send_event(sock, bad)
assert receiver.join() and receiver.terminal == "failed"
assert len(received) == (0 if case == "no-init" else 2)
def test_external_reservation_is_not_freed_by_cancel_and_cannot_double_release():
queue = StreamMailbox(byte_limit=100)
raw = queue.reserve_ingress(20)
queue.cancel()
assert not queue.quiescent and queue.bytes == 20
raw.release()
assert queue.quiescent and queue.bytes == 0
with pytest.raises(ValueError, match="ownership"):
raw.release()
def test_sender_streams_existing_events_and_explicit_gap_then_end(tmp_path):
with harness(tmp_path) as (sock, receiver, run, received, notices):
sender = StreamingSender(sock, start(), "capture-1", 1, lambda: run.check_current(start()))
sender.send(event())
sender.gap(modality="pose", reason="unavailable", count=0)
sender.send(event(2, modality="pose"))
sender.end()
assert receiver.join() and receiver.terminal == "end"
assert len(received) == 2 and notices[0]["reason"] == "unavailable"
with pytest.raises(wire.StreamWireError, match="closed"):
sender.send(event(3))
def test_sender_timeout_closes_transport_instead_of_building_backlog():
left, right = socket.socketpair()
left.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, 4096)
try:
sender = StreamingSender(left, start(), "capture-1", 1, lambda: None, timeout=0.01)
with pytest.raises(wire.StreamWireError, match="deadline"):
sender.send(event(payload=b"x" * 65536))
assert sender.closed and left.fileno() == -1
finally:
left.close()
right.close()
def test_existing_raw_first_ingress_events_feed_new_writer_without_wire_v1_changes(tmp_path):
source = LivePerceptionIngress()
source.begin_session("capture-1")
source.open_consumer("stream-adapter")
assert source.take_next("stream-adapter", timeout=0).modality == "control"
assert source.publish(
modality="lidar",
source_id="existing-raw-topic",
source_sequence=12,
captured_at_epoch_ns=1799999999123456789,
received_monotonic_ns=9007199254740993,
payload=b"unchanged raw payload",
)
raw = source.take_next("stream-adapter", timeout=0)
legacy = raw.wire_bytes()
with harness(tmp_path) as (sock, receiver, _, received, _):
sender = StreamingSender(sock, start(), "capture-1", 1, lambda: None)
sender.send(raw)
sender.end()
assert receiver.join() and receiver.terminal == "end"
assert received == [raw] and raw.wire_bytes() == legacy
source.close_consumer("stream-adapter")
source.close()