NODEDC_MISSION_CORE/tests/test_xgrids_application_acc...

567 lines
20 KiB
Python

from __future__ import annotations
import hashlib
from collections.abc import Collection, Sequence
import pytest
from pydantic import JsonValue, TypeAdapter
from k1link.device_plugins.xgrids_k1.protocol.application_acceptance import (
ApplicationAcceptanceError,
OperatorDialogueCheckpoint,
PhysicalAcceptanceChecklist,
PhysicalAcceptanceDialogueExecutor,
PhysicalAcceptancePermit,
ScanInitializationTimeout,
)
from k1link.device_plugins.xgrids_k1.protocol.application_bootstrap import (
DEVICE_CONFIG_RESPONSE_TOPIC,
ApplicationControlAuthority,
ShadowApplicationBootstrapOrchestrator,
build_canonical_post_start_observation,
)
from k1link.device_plugins.xgrids_k1.protocol.application_publish import (
OneShotPublishEnvelope,
)
from k1link.device_plugins.xgrids_k1.protocol.modeling_control import (
OPENAPI_SUCCESS,
CommandHeaderIdentity,
ModelingAction,
MountType,
RecordMode,
ScanMode,
encode_modeling_start,
encode_modeling_stop,
)
from k1link.device_plugins.xgrids_k1.protocol.modeling_safety import (
ShadowModelingCommand,
)
APPLICATION_KEY = "11111111-2222-3333-4444-555555555555"
VENDOR_DEVICE_ID = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee"
def _varint(value: int) -> bytes:
encoded = bytearray()
while value > 0x7F:
encoded.append((value & 0x7F) | 0x80)
value >>= 7
encoded.append(value)
return bytes(encoded)
def _uint(number: int, value: int) -> bytes:
return _varint(number << 3) + _varint(value)
def _bytes(number: int, value: bytes) -> bytes:
return _varint((number << 3) | 2) + _varint(len(value)) + value
def _text(number: int, value: str) -> bytes:
return _bytes(number, value.encode())
def _header(session_id: str) -> bytes:
return _text(4, VENDOR_DEVICE_ID) + _text(5, session_id) + _text(6, APPLICATION_KEY)
def _device_info_response(
session_id: str,
*,
device_model: str = "LixelKity K1",
device_type: str = "A4",
) -> bytes:
base_info = b"".join(
(
_text(2, "V3.0.2_20250624.122658"),
_text(3, "V3.0.2"),
_text(6, device_model),
_text(7, "K1SERIAL01"),
_text(8, device_type),
)
)
working_status = _uint(1, 1)
device_info = _bytes(2, base_info) + _bytes(7, working_status)
error = _uint(1, OPENAPI_SUCCESS)
return _bytes(1, _header(session_id)) + _bytes(2, device_info) + _bytes(15, error)
def _generic_response(session_id: str) -> bytes:
return _bytes(1, _header(session_id)) + _bytes(15, _uint(1, OPENAPI_SUCCESS))
def _modeling_response(action: ModelingAction) -> bytes:
session = f"{VENDOR_DEVICE_ID}:ModelingRequest"
return _bytes(1, _header(session)) + _uint(2, action) + _bytes(15, _uint(1, OPENAPI_SUCCESS))
class SyntheticAcceptanceTransport:
def __init__(
self,
clock: FakeClock | None = None,
*,
device_model: str = "LixelKity K1",
device_type: str = "A4",
) -> None:
self.batches: list[tuple[str, ...]] = []
self.maintain_calls: list[tuple[float, tuple[str, ...]]] = []
self.clock = clock
self.pre_start_maintain_calls = 0
self.active_maintain_calls = 0
self.post_stop_maintain_calls = 0
self.start_emitted = False
self.stop_emitted = False
self.device_model = device_model
self.device_type = device_type
def exchange_batch_once(
self,
envelopes: Sequence[OneShotPublishEnvelope],
*,
required_response_operation_keys: Collection[str],
) -> dict[str, bytes]:
self.batches.append(tuple(envelope.operation_key for envelope in envelopes))
responses: dict[str, bytes] = {}
modeling_operation = next(
(
operation_key
for operation_key in required_response_operation_keys
if operation_key in {"modeling:start", "modeling:stop"}
),
None,
)
if modeling_operation is not None:
action = (
ModelingAction.START
if modeling_operation == "modeling:start"
else ModelingAction.STOP
)
responses[modeling_operation] = _modeling_response(action)
if action is ModelingAction.START:
self.start_emitted = True
else:
self.stop_emitted = True
return responses
for envelope in envelopes:
prefix, ordinal_text, message_type = envelope.operation_key.split(":", 2)
assert prefix in {"bootstrap", "dialogue"}
ordinal = int(ordinal_text)
if message_type == "DeviceInfoRequest":
session = (
":DeviceInfoRequest"
if ordinal == 1
else f"{VENDOR_DEVICE_ID}:DeviceInfoRequest"
)
response = _device_info_response(
session,
device_model=self.device_model,
device_type=self.device_type,
)
else:
session = (
f":{message_type}" if ordinal <= 3 else f"{VENDOR_DEVICE_ID}:{message_type}"
)
if message_type == "DeviceConfigRequest":
session += ":Publish_Proto_DeviceConfig_SetTime"
response = _generic_response(session)
if envelope.operation_key in required_response_operation_keys:
responses[envelope.operation_key] = response
return responses
def maintain_open_for(
self,
duration_seconds: float,
*,
allowed_response_topics: Collection[str] = (),
) -> None:
self.maintain_calls.append(
(duration_seconds, tuple(sorted(allowed_response_topics)))
)
if self.clock is not None:
self.clock.now += duration_seconds
if self.stop_emitted:
self.post_stop_maintain_calls += 1
elif self.start_emitted:
self.active_maintain_calls += 1
else:
self.pre_start_maintain_calls += 1
def scan_initialization_complete(self, _binding: object) -> bool:
return self.active_maintain_calls >= 2
def pre_start_ready(self, _binding: object) -> bool:
return not self.start_emitted
def standby_complete(self, _binding: object) -> bool:
return self.stop_emitted and self.post_stop_maintain_calls >= 2
def validate_bound_status(self, _binding: object) -> None:
return None
def _checklist(action: ModelingAction) -> PhysicalAcceptanceChecklist:
return PhysicalAcceptanceChecklist(
action=action,
operator_present=True,
owner_controlled_device=True,
lixelgo_closed=True,
battery_storage_confirmed=True,
expected_physical_state_confirmed=True,
)
class FakeClock:
def __init__(self) -> None:
self.now = 100.0
def __call__(self) -> float:
return self.now
def test_canonical_session_owns_start_active_scan_stop_and_save_boundary() -> None:
clock = FakeClock()
transport = SyntheticAcceptanceTransport(clock)
executor = PhysicalAcceptanceDialogueExecutor(transport)
orchestrator = ShadowApplicationBootstrapOrchestrator(
ApplicationControlAuthority(openapi_key=APPLICATION_KEY),
epoch_seconds=1_752_680_000,
timezone_name="Europe/Moscow",
)
with pytest.raises(ApplicationAcceptanceError, match="collapsed bootstrap is disabled"):
executor.run_bootstrap(orchestrator)
binding = executor.run_connection_stage(orchestrator)
with pytest.raises(ApplicationAcceptanceError, match="not issued by this canonical session"):
executor.run_workspace_entry_stage(
orchestrator,
OperatorDialogueCheckpoint(
event="workspace-entered",
operator_initiated=True,
owner_token=object(),
),
)
assert [len(batch) for batch in transport.batches] == [1, 5]
assert transport.batches[1] == (
"bootstrap:2:ModelingStatusRequest",
"bootstrap:3:GetRtkAdvanceRequest",
"bootstrap:4:DeviceConfigRequest",
"bootstrap:5:DeviceInfoRequest",
"bootstrap:6:GetRtkAdvanceRequest",
)
workspace_checks = iter((False, False, True))
workspace_checkpoint = executor.wait_for_operator_checkpoint(
"workspace-entered",
lambda: next(workspace_checks),
)
executor.run_workspace_entry_stage(
orchestrator,
workspace_checkpoint,
)
project_checks = iter((False, False, True))
project_checkpoint = executor.wait_for_operator_checkpoint(
"project-prompt-opened",
lambda: next(project_checks),
)
executor.run_project_prompt_stage(
orchestrator,
project_checkpoint,
)
command = ShadowModelingCommand.from_command(
encode_modeling_start(
CommandHeaderIdentity(
device_id=binding.vendor_device_id,
openapi_key=APPLICATION_KEY,
),
project_name="SAFE_PROJECT",
record_mode=RecordMode.RECORD_AND_CALCULATE,
scan_mode=ScanMode.LCC,
mount_type=MountType.HANDHELD,
)
)
with pytest.raises(ApplicationAcceptanceError, match="standalone modeling commands"):
executor.execute_modeling(command)
start_wait_polls = 0
def start_confirmed() -> bool:
nonlocal start_wait_polls
start_wait_polls += 1
return start_wait_polls > 202
start_checkpoint = executor.wait_for_operator_checkpoint(
"start-confirmed",
start_confirmed,
)
# Human/UI time is serviced on the socket and precedes the command permit.
start_permit = PhysicalAcceptancePermit(
_checklist(ModelingAction.START),
ttl_seconds=120.0,
monotonic=clock,
)
post_start = build_canonical_post_start_observation(
ApplicationControlAuthority(openapi_key=APPLICATION_KEY),
binding,
)
response = executor.execute_canonical_start(
command,
post_start,
authority=ApplicationControlAuthority(openapi_key=APPLICATION_KEY),
binding=binding,
permit=start_permit,
checkpoint=start_checkpoint,
)
assert response.action is ModelingAction.START
stop_checks = iter((False, False, True))
executor.maintain_active_until_stop_requested(lambda: next(stop_checks))
stop_permit = PhysicalAcceptancePermit(
_checklist(ModelingAction.STOP),
ttl_seconds=120.0,
monotonic=clock,
)
stop_command = ShadowModelingCommand.from_command(
encode_modeling_stop(
CommandHeaderIdentity(
device_id=binding.vendor_device_id,
openapi_key=APPLICATION_KEY,
)
)
)
stop_response = executor.execute_canonical_stop(stop_command, stop_permit)
executor.maintain_post_stop_until_standby()
assert stop_response.action is ModelingAction.STOP
assert [len(batch) for batch in transport.batches] == [1, 5, 1, 3, 1, 1, 2, 1]
assert len(transport.maintain_calls) == 212
assert set(transport.maintain_calls) == {
(1.0, ("lixel/application/response/modeling_status",))
}
assert start_permit.snapshot()["consumed"] is True
assert stop_permit.snapshot()["consumed"] is True
assert executor.snapshot()["bootstrap_complete"] is True
assert executor.snapshot()["command_complete"] is True
assert executor.snapshot()["start_complete"] is True
assert executor.snapshot()["stop_complete"] is True
assert executor.snapshot()["dialogue_stage"] == "standby-confirmed"
response_evidence = executor.snapshot()["response_evidence"]
assert isinstance(response_evidence, list)
assert len(response_evidence) == 13
assert response_evidence[0]["operation_key"] == "bootstrap:1:DeviceInfoRequest"
assert response_evidence[-1]["operation_key"] == "modeling:stop"
assert all("payload" not in item for item in response_evidence)
assert APPLICATION_KEY not in str(executor.snapshot())
assert VENDOR_DEVICE_ID not in str(executor.snapshot())
TypeAdapter(JsonValue).validate_python(executor.snapshot())
def test_start_initialization_wait_has_fail_closed_watchdog() -> None:
class NeverInitializedTransport(SyntheticAcceptanceTransport):
def scan_initialization_complete(self, _binding: object) -> bool:
return False
clock = FakeClock()
transport = NeverInitializedTransport(clock)
executor = PhysicalAcceptanceDialogueExecutor(
transport,
monotonic=clock,
scan_initialization_timeout_seconds=2.0,
)
authority = ApplicationControlAuthority(openapi_key=APPLICATION_KEY)
orchestrator = ShadowApplicationBootstrapOrchestrator(
authority,
epoch_seconds=1_752_680_000,
timezone_name="Europe/Moscow",
)
binding = executor.run_connection_stage(orchestrator)
executor.run_workspace_entry_stage(
orchestrator,
executor.wait_for_operator_checkpoint("workspace-entered", lambda: True),
)
executor.run_project_prompt_stage(
orchestrator,
executor.wait_for_operator_checkpoint("project-prompt-opened", lambda: True),
)
start_checkpoint = executor.wait_for_operator_checkpoint(
"start-confirmed",
lambda: True,
)
command = ShadowModelingCommand.from_command(
encode_modeling_start(
CommandHeaderIdentity(
device_id=binding.vendor_device_id,
openapi_key=APPLICATION_KEY,
),
project_name="SAFE_PROJECT",
record_mode=RecordMode.RECORD_AND_CALCULATE,
scan_mode=ScanMode.LCC,
mount_type=MountType.HANDHELD,
)
)
permit = PhysicalAcceptancePermit(
_checklist(ModelingAction.START),
monotonic=clock,
)
with pytest.raises(ScanInitializationTimeout):
executor.execute_canonical_start(
command,
build_canonical_post_start_observation(authority, binding),
authority=authority,
binding=binding,
permit=permit,
checkpoint=start_checkpoint,
)
assert transport.start_emitted
assert not transport.stop_emitted
assert executor.snapshot()["dialogue_stage"] == "initializing"
assert executor.snapshot()["start_attempted"] is True
assert executor.snapshot()["start_complete"] is False
@pytest.mark.parametrize(
("device_model", "device_type"),
(
("LixelKity K2", "A4"),
("LixelKity K1", "K1"),
),
)
def test_connection_stage_rejects_other_device_model_or_type_before_preparation(
device_model: str,
device_type: str,
) -> None:
transport = SyntheticAcceptanceTransport(
device_model=device_model,
device_type=device_type,
)
executor = PhysicalAcceptanceDialogueExecutor(transport)
orchestrator = ShadowApplicationBootstrapOrchestrator(
ApplicationControlAuthority(openapi_key=APPLICATION_KEY),
epoch_seconds=1_752_680_000,
timezone_name="Europe/Moscow",
)
with pytest.raises(ApplicationAcceptanceError, match="incompatible with the selected"):
executor.run_connection_stage(orchestrator)
assert transport.batches == [("bootstrap:1:DeviceInfoRequest",)]
assert not transport.start_emitted
assert orchestrator.binding is None
snapshot = executor.snapshot()
assert snapshot["correlation_failure"] is None
compatibility_failure = snapshot["compatibility_failure"]
assert isinstance(compatibility_failure, dict)
assert compatibility_failure["reason_code"] == "compatibility_profile_mismatch"
assert compatibility_failure["expected"] == {
"device_model": "LixelKity K1",
"platform_type": "A4",
"firmware": "3.0.2",
"is_activated": True,
}
def test_bootstrap_correlation_failure_records_only_redacted_response_evidence() -> None:
class CorruptDeviceConfigTransport(SyntheticAcceptanceTransport):
def exchange_batch_once(
self,
envelopes: Sequence[OneShotPublishEnvelope],
*,
required_response_operation_keys: Collection[str],
) -> dict[str, bytes]:
responses = super().exchange_batch_once(
envelopes,
required_response_operation_keys=required_response_operation_keys,
)
operation_key = "bootstrap:4:DeviceConfigRequest"
if operation_key in responses:
responses[operation_key] = _generic_response(
f"{VENDOR_DEVICE_ID}:wrong-session"
)
return responses
executor = PhysicalAcceptanceDialogueExecutor(
CorruptDeviceConfigTransport(),
)
orchestrator = ShadowApplicationBootstrapOrchestrator(
ApplicationControlAuthority(openapi_key=APPLICATION_KEY),
epoch_seconds=1_752_680_000,
timezone_name="Europe/Moscow",
)
with pytest.raises(
ApplicationAcceptanceError,
match=r"ordinal 4 \(DeviceConfigRequest\).*session correlation",
):
executor.run_connection_stage(orchestrator)
snapshot = executor.snapshot()
failure = snapshot["correlation_failure"]
assert failure == {
"phase": "bootstrap",
"operation_key": "bootstrap:4:DeviceConfigRequest",
"response_topic": DEVICE_CONFIG_RESPONSE_TOPIC,
"reason_code": "response_session_mismatch",
"reason": "application response session correlation failed",
}
evidence = snapshot["response_evidence"]
assert isinstance(evidence, list)
assert evidence[-1]["payload_sha256"] == hashlib.sha256(
_generic_response(f"{VENDOR_DEVICE_ID}:wrong-session")
).hexdigest()
assert evidence[-1]["payload_bytes"] > 0
assert all("payload" not in item for item in evidence)
assert APPLICATION_KEY not in str(snapshot)
assert VENDOR_DEVICE_ID not in str(snapshot)
TypeAdapter(JsonValue).validate_python(snapshot)
def test_standalone_start_and_stop_are_both_disabled() -> None:
transport = SyntheticAcceptanceTransport()
start = PhysicalAcceptanceDialogueExecutor(transport)
start_command = ShadowModelingCommand.from_command(
encode_modeling_start(
CommandHeaderIdentity(device_id=VENDOR_DEVICE_ID, openapi_key=APPLICATION_KEY),
project_name="SAFE_PROJECT",
record_mode=RecordMode.RECORD_AND_CALCULATE,
scan_mode=ScanMode.LCC,
mount_type=MountType.HANDHELD,
)
)
with pytest.raises(ApplicationAcceptanceError, match="standalone modeling commands"):
start.execute_modeling(start_command)
stop = PhysicalAcceptanceDialogueExecutor(transport)
stop_command = ShadowModelingCommand.from_command(
encode_modeling_stop(
CommandHeaderIdentity(device_id=VENDOR_DEVICE_ID, openapi_key=APPLICATION_KEY)
)
)
with pytest.raises(ApplicationAcceptanceError, match="standalone modeling commands"):
stop.execute_modeling(stop_command)
assert transport.batches == []
def test_acceptance_permit_rejects_implicit_or_expired_authority() -> None:
with pytest.raises(ApplicationAcceptanceError, match="explicitly true"):
PhysicalAcceptanceChecklist(
action=ModelingAction.START,
operator_present=False, # type: ignore[arg-type]
owner_controlled_device=True,
lixelgo_closed=True,
battery_storage_confirmed=True,
expected_physical_state_confirmed=True,
)
now = [100.0]
permit = PhysicalAcceptancePermit(
_checklist(ModelingAction.STOP),
ttl_seconds=15.0,
monotonic=lambda: now[0],
)
now[0] = 115.0
with pytest.raises(ApplicationAcceptanceError, match="expired"):
permit.consume(ModelingAction.STOP)