"""Small synthetic memory/clock/pipe checks; no real decode or models on the Mac.""" import json import socket import struct from contextlib import contextmanager from pathlib import Path from types import SimpleNamespace import numpy as np import pytest from k1link.perception.streaming_pipe_rpc import BoundedPipeRpc, PipeRpcError from k1link.perception.streaming_queue import StreamMailbox from k1link.perception.streaming_sensors import CausalSensorWindow, normalized_sensor def bundle(seq=0, size=10): return {"sequence": seq, "payload_bytes": size} def test_reservation_handoff_is_atomic_and_does_not_double_charge(): mailbox = StreamMailbox(byte_limit=10) ticket = mailbox.reserve_ingress(10) value = bundle() assert mailbox.put_reserved(value, ticket) assert mailbox.bytes == mailbox.peak_bytes == 10 with pytest.raises(ValueError, match="ownership"): ticket.release() assert mailbox.take() is value mailbox.cancel() assert mailbox.bytes == 10 and not mailbox.quiescent mailbox.release(value) assert mailbox.quiescent @pytest.mark.parametrize("case", ["closed", "size", "foreign", "sequence"]) def test_rejected_handoff_keeps_caller_ownership(case): mailbox = StreamMailbox(byte_limit=40) ticket = mailbox.reserve_ingress(10) value = bundle() if case == "closed": mailbox.cancel() assert not mailbox.put_reserved(value, ticket) else: if case == "size": value["payload_bytes"] = 11 if case == "sequence": mailbox.put(bundle(1)) with pytest.raises(ValueError): (StreamMailbox() if case == "foreign" else mailbox).put_reserved(value, ticket) assert mailbox.bytes >= 10 ticket.size = 99999 # Only the stored reservation size is authoritative. ticket.release() mailbox.cancel() assert mailbox.bytes == 0 and mailbox.quiescent def test_handoff_overflow_drops_pending_only(): mailbox = StreamMailbox(byte_limit=40) mailbox.put(bundle(0)) active = mailbox.take() mailbox.put(bundle(1, 5)) mailbox.put(bundle(2, 5)) ticket = mailbox.reserve_ingress(20) assert mailbox.put_reserved(bundle(3, 20), ticket) assert mailbox.dropped == [{"sequence": 1, "reason": "pending-overflow"}] assert mailbox.bytes == 35 mailbox.cancel() assert mailbox.bytes == 10 mailbox.release(active) assert mailbox.quiescent @contextmanager def pipe(response=None, check=lambda: None): client, server = socket.socketpair() try: if response: server.sendall(response) yield BoundedPipeRpc(client.fileno(), client.fileno(), check), server finally: client.close() server.close() def message(header, payload=b""): raw = json.dumps({**header, "payload_bytes": len(payload)}).encode() return struct.pack("