2 Commits
7 changed files with 372 additions and 10 deletions
@@ -261,6 +261,18 @@ export function MseFmp4WebSocketPlayer({
}); });
activeDiagnosticRef.current = diagnosticLifecycle; activeDiagnosticRef.current = diagnosticLifecycle;
diagnosticLifecycle.verifyBuild(); diagnosticLifecycle.verifyBuild();
diagnosticLifecycle.post({
eventCode: "live_camera_player_effect_started",
streamId: delivery.id,
transportEpoch,
});
if (!recoveryAuthorityIdentity) {
diagnosticLifecycle.post({
eventCode: "live_camera_playback_authority_missing",
streamId: delivery.id,
transportEpoch,
});
}
const transportIsCurrent = () => cameraTransportCallbackIsCurrent( const transportIsCurrent = () => cameraTransportCallbackIsCurrent(
transportEpochRef.current, transportEpochRef.current,
transportEpoch, transportEpoch,
@@ -473,6 +485,11 @@ export function MseFmp4WebSocketPlayer({
const onSourceOpen = () => { const onSourceOpen = () => {
if (!transportIsCurrent() || failed) return; if (!transportIsCurrent() || failed) return;
diagnosticLifecycle.post({
eventCode: "live_camera_media_source_open",
streamId: delivery.id,
transportEpoch,
});
try { try {
sourceBuffer = mediaSource.addSourceBuffer(mediaType); sourceBuffer = mediaSource.addSourceBuffer(mediaType);
} catch { } catch {
@@ -525,6 +542,11 @@ export function MseFmp4WebSocketPlayer({
socket.binaryType = "arraybuffer"; socket.binaryType = "arraybuffer";
socket.addEventListener("open", () => { socket.addEventListener("open", () => {
if (!transportIsCurrent() || failed) return; if (!transportIsCurrent() || failed) return;
diagnosticLifecycle.post({
eventCode: "live_camera_websocket_open",
streamId: delivery.id,
transportEpoch,
});
startupWatchdog?.armFirstMedia(); startupWatchdog?.armFirstMedia();
setStatus("buffering"); setStatus("buffering");
setMessage(recovering setMessage(recovering
@@ -7,6 +7,10 @@ export type LiveViewerDiagnosticEventCode =
| "live_receiver_error" | "live_receiver_error"
| "live_camera_source_projected" | "live_camera_source_projected"
| "live_camera_window_admitted" | "live_camera_window_admitted"
| "live_camera_player_effect_started"
| "live_camera_playback_authority_missing"
| "live_camera_media_source_open"
| "live_camera_websocket_open"
| "live_camera_transport_restart_requested" | "live_camera_transport_restart_requested"
| "live_camera_transport_playing"; | "live_camera_transport_playing";
@@ -203,6 +203,24 @@ test("live camera registers sourceopen before assigning the MSE object URL", asy
); );
}); });
test("live camera reports each browser startup boundary without device actions", async () => {
const source = await readFile(
new URL("../src/components/MseFmp4WebSocketPlayer.tsx", import.meta.url),
"utf8",
);
for (const eventCode of [
"live_camera_player_effect_started",
"live_camera_playback_authority_missing",
"live_camera_media_source_open",
"live_camera_websocket_open",
]) {
assert.ok(
source.includes(`eventCode: "${eventCode}"`),
`missing bounded player diagnostic ${eventCode}`,
);
}
});
test("recorded observation timeline clamps relative seek time without epoch precision in the UI", () => { test("recorded observation timeline clamps relative seek time without epoch precision in the UI", () => {
const range = normalizeTimelineRange({ const range = normalizeTimelineRange({
min: 0, min: 0,
+103 -2
View File
@@ -28455,6 +28455,91 @@ class XgridsK1CompatibilityService:
) )
) )
@staticmethod
def _matching_accepted_start_scan_over_before_active(
physical_command_proof: Mapping[str, Any] | None,
*,
acquisition_id: str | None,
start_operation_id: str | None,
verified_control: Mapping[str, Any] | None,
) -> bool:
"""Recognize an accepted START followed by bound pre-active SCAN_OVER.
K1 can accept the canonical START and then leave calibration before it
ever reports the target SCANNING state (for example, after losing
power). The physical ledger correctly remains unresolved because no
active target was observed. Local capture may nevertheless terminate
once the same DeviceInfo-bound generation reports fresh, non-retained
SCAN_OVER. This proof never resolves the physical ledger or authorizes
another START/STOP; a later explicit read-only reconciliation remains
mandatory.
"""
if (
not isinstance(physical_command_proof, Mapping)
or not isinstance(acquisition_id, str)
or not acquisition_id
or not isinstance(start_operation_id, str)
or not start_operation_id
or not isinstance(verified_control, Mapping)
or physical_command_proof.get("status") != "unresolved"
or physical_command_proof.get("requires_reconciliation") is not True
or physical_command_proof.get("runtime_bound") is not True
or physical_command_proof.get("reconciliation_ready") is not True
or physical_command_proof.get("observed_session_state") != "scan_over"
or physical_command_proof.get("active_operation_id")
!= start_operation_id
):
return False
record = physical_command_proof.get("record")
if not isinstance(record, Mapping):
return False
expected_binding = (
XgridsK1CompatibilityService._exact_control_binding_document(
verified_control
)
)
baseline = record.get("baseline_status")
response = record.get("application_response")
last_status = record.get("last_status")
last_status_valid = bool(
last_status is None
or (
isinstance(last_status, Mapping)
and last_status.get("session_state") == "scanning"
and last_status.get("project_bound") is True
and last_status.get("init_ready") is True
and last_status.get("mqtt_retained") is False
and last_status.get("system_error_code") is None
)
)
return bool(
expected_binding is not None
and isinstance(record.get("connection"), Mapping)
and dict(record["connection"]) == expected_binding
and record.get("operation_id") == start_operation_id
and record.get("acquisition_id") == acquisition_id
and record.get("action") == "start"
and record.get("stage") == "observing"
and record.get("resolution") is None
and record.get("publish_call_returned") is True
and record.get("qos2_completed") is True
and isinstance(record.get("packet_id"), int)
and not isinstance(record.get("packet_id"), bool)
and int(record["packet_id"]) > 0
and isinstance(baseline, Mapping)
and baseline.get("session_state") == "ready"
and baseline.get("project_bound") is False
and baseline.get("init_ready") is False
and baseline.get("mqtt_retained") is False
and baseline.get("system_error_code") is None
and isinstance(response, Mapping)
and response.get("operation_id") == start_operation_id
and response.get("action") == "start"
and response.get("success") is True
and last_status_valid
)
@staticmethod @staticmethod
def _matching_classified_prepared_stop_active( def _matching_classified_prepared_stop_active(
physical_command_proof: Mapping[str, Any] | None, physical_command_proof: Mapping[str, Any] | None,
@@ -30254,6 +30339,18 @@ class XgridsK1CompatibilityService:
if isinstance(terminal_control_proof, Mapping) if isinstance(terminal_control_proof, Mapping)
else None else None
) )
accepted_start_scan_over_before_active = (
self._matching_accepted_start_scan_over_before_active(
physical_command_proof,
acquisition_id=current_acquisition_id,
start_operation_id=canonical_start_operation_id,
verified_control=(
terminal_verified_control
if isinstance(terminal_verified_control, Mapping)
else None
),
)
)
device_reported_scan_over_without_stop = bool( device_reported_scan_over_without_stop = bool(
current is not None current is not None
and current.control_mode == "plugin-commanded" and current.control_mode == "plugin-commanded"
@@ -30262,7 +30359,8 @@ class XgridsK1CompatibilityService:
and isinstance(terminal_control_proof, Mapping) and isinstance(terminal_control_proof, Mapping)
and terminal_control_proof.get("state") == "failed" and terminal_control_proof.get("state") == "failed"
and isinstance(terminal_control_failure, Mapping) and isinstance(terminal_control_failure, Mapping)
and terminal_control_failure.get("failed_phase") == "scanning" and terminal_control_failure.get("failed_phase")
in {"initializing", "scanning"}
and terminal_control_failure.get("modeling_command_attempted") is True and terminal_control_failure.get("modeling_command_attempted") is True
and terminal_control_failure.get("stop_command_attempted") is False and terminal_control_failure.get("stop_command_attempted") is False
and terminal_control_failure.get("diagnostic_snapshot_unavailable") == [] and terminal_control_failure.get("diagnostic_snapshot_unavailable") == []
@@ -30281,7 +30379,8 @@ class XgridsK1CompatibilityService:
and physical_command_proof.get("runtime_bound") is True and physical_command_proof.get("runtime_bound") is True
and physical_command_proof.get("reconciliation_ready") is True and physical_command_proof.get("reconciliation_ready") is True
and physical_command_proof.get("observed_session_state") == "scan_over" and physical_command_proof.get("observed_session_state") == "scan_over"
and self._matching_start_active_confirmed( and (
self._matching_start_active_confirmed(
physical_command_proof, physical_command_proof,
acquisition_id=current_acquisition_id, acquisition_id=current_acquisition_id,
start_operation_id=canonical_start_operation_id, start_operation_id=canonical_start_operation_id,
@@ -30291,6 +30390,8 @@ class XgridsK1CompatibilityService:
else None else None
), ),
) )
or accepted_start_scan_over_before_active
)
) )
# Camera-only transport loss is supervised by the backend producer # Camera-only transport loss is supervised by the backend producer
# watchdog and exact camera CAS. Snapshot polling is observational: it # watchdog and exact camera CAS. Snapshot polling is observational: it
+4
View File
@@ -20,6 +20,10 @@ LiveViewerEventCode = Literal[
"live_receiver_error", "live_receiver_error",
"live_camera_source_projected", "live_camera_source_projected",
"live_camera_window_admitted", "live_camera_window_admitted",
"live_camera_player_effect_started",
"live_camera_playback_authority_missing",
"live_camera_media_source_open",
"live_camera_websocket_open",
"live_camera_transport_restart_requested", "live_camera_transport_restart_requested",
"live_camera_transport_playing", "live_camera_transport_playing",
] ]
+10
View File
@@ -247,6 +247,10 @@ def test_live_viewer_diagnostic_endpoint_accepts_only_bounded_events(
for event_code in ( for event_code in (
"live_camera_source_projected", "live_camera_source_projected",
"live_camera_window_admitted", "live_camera_window_admitted",
"live_camera_player_effect_started",
"live_camera_playback_authority_missing",
"live_camera_media_source_open",
"live_camera_websocket_open",
): ):
boundary_event = LiveViewerDiagnosticEvent( boundary_event = LiveViewerDiagnosticEvent(
schema_version="missioncore.live-viewer-diagnostic/v2", schema_version="missioncore.live-viewer-diagnostic/v2",
@@ -256,6 +260,12 @@ def test_live_viewer_diagnostic_endpoint_accepts_only_bounded_events(
viewer_instance_id="00000000-0000-4000-8000-000000000004", viewer_instance_id="00000000-0000-4000-8000-000000000004",
lifecycle_generation=1, lifecycle_generation=1,
stream_id="camera-preview-2", stream_id="camera-preview-2",
transport_epoch=(
1
if event_code
not in {"live_camera_source_projected", "live_camera_window_admitted"}
else None
),
) )
with caplog.at_level( with caplog.at_level(
logging.INFO, logging.INFO,
+203
View File
@@ -17194,6 +17194,209 @@ def test_bound_scan_over_without_stop_dominates_late_empty_camera_epoch_failure(
assert control.stop_calls == stop_calls_before == 0 assert control.stop_calls == stop_calls_before == 0
def test_accepted_start_then_scan_over_during_calibration_stops_local_wait_only(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
service, runtime = service_with_fake_runtime(tmp_path)
control = FakeInteractiveControlSession()
service._application_control_session = control # type: ignore[assignment] # noqa: SLF001
binding = _seed_supervised_connection(service)
prepared = service.prepare_acquisition(
_prepare_request(
project_name="POWERLOSS001",
host=binding.target_ipv4,
compatibility_attestation=ATTESTATION,
)
)
acquisition_id = str(prepared["acquisition"]["acquisition_id"])
service.start_acquisition(
_start_request(
acquisition_id=acquisition_id,
physical_acceptance=PHYSICAL_ACCEPTANCE,
)
)
start_operation_id = service._acquisition_start_operation_id # noqa: SLF001
assert isinstance(start_operation_id, str)
physical = _exact_start_physical_proof(
operation_id=start_operation_id,
acquisition_id=acquisition_id,
binding=binding,
resolved=False,
)
physical.update(
{
"runtime_bound": True,
"reconciliation_ready": True,
"observed_session_state": "scan_over",
"active_operation_id": start_operation_id,
}
)
record = physical["record"]
assert isinstance(record, dict)
record["last_status"] = None
monkeypatch.setattr(
service._physical_command_coordinator, # noqa: SLF001
"snapshot",
lambda: physical,
)
runtime.mark_ready()
awaiting = service.state()
assert awaiting["acquisition"]["state"] == "awaiting_external_start"
assert runtime.pcl_frames == 0
original_control_snapshot = control.snapshot
def failed_during_calibration_snapshot() -> dict[str, object]:
snapshot = original_control_snapshot()
snapshot["outcome_unknown"] = True
snapshot["failure"] = {
"reason_code": "mqtt_network_loop_failed",
"failed_phase": "initializing",
"modeling_command_attempted": True,
"stop_command_attempted": False,
"diagnostic_snapshot_unavailable": [],
"diagnostic_evidence_unavailable": [],
"safe_to_retry": False,
}
snapshot["transport"] = {
"state": "failed",
"publish_attempts": 12,
"device_status_reports": 4,
"latest_device_session_state": "scan_over",
"latest_device_project_bound": True,
"latest_device_init_ready": False,
"latest_system_error_code": None,
"automatic_retry": False,
"automatic_reconnect": False,
}
return snapshot
monkeypatch.setattr(control, "snapshot", failed_during_calibration_snapshot)
control.state = "failed"
control.state_revision += 1
monkeypatch.setattr(service, "_seal_acquisition_capture_clock", lambda: None)
camera_stop_calls: list[str] = []
monkeypatch.setattr(
service.camera_preview,
"stop_current",
lambda: camera_stop_calls.append("stop") or {},
)
start_projects_before = list(control.start_projects)
physical_record_before = json.dumps(record, sort_keys=True)
terminal = service.state()
start_operation = next(
operation
for operation in terminal["operations"]
if operation["operation_id"] == start_operation_id
)
assert terminal["acquisition"]["state"] == "interrupted"
assert terminal["acquisition"]["message_code"] == (
"acquisition.recovery.device_standby_observed"
)
assert terminal["acquisition"]["result"] == {
"receiver_stopped": True,
"device_state": "scan_over",
"device_stop": "not-sent",
"automatic_command_retry": False,
"read_only_recovery": True,
"physical_reconciliation_required": True,
}
assert terminal["acquisition"]["cleanup_pending"] is False
assert start_operation["status"] == "interrupted"
assert start_operation["stage_code"] == (
"physical-start-active-but-no-point-before-standby"
)
assert start_operation["error"]["side_effect_status"] == "succeeded"
assert start_operation["error"]["automatic_replay_allowed"] is False
assert camera_stop_calls == ["stop"]
assert runtime.stop_calls == 1
assert control.start_projects == start_projects_before == ["POWERLOSS001"]
assert control.stop_calls == 0
assert physical["status"] == "unresolved"
assert physical["requires_reconciliation"] is True
assert json.dumps(record, sort_keys=True) == physical_record_before
repeated = service.state()
assert repeated["acquisition"] == terminal["acquisition"]
assert repeated["physical_command"]["status"] == "unresolved"
assert repeated["physical_command"]["requires_reconciliation"] is True
assert camera_stop_calls == ["stop"]
assert runtime.stop_calls == 1
assert control.start_projects == ["POWERLOSS001"]
assert control.stop_calls == 0
@pytest.mark.parametrize(
"invalid_proof",
[
"active-owner",
"application-response",
"binding-generation",
"qos2",
"retained-target",
],
)
def test_pre_active_scan_over_classifier_rejects_inexact_start_proof(
tmp_path: Path,
invalid_proof: str,
) -> None:
service, _runtime = service_with_fake_runtime(tmp_path)
control = FakeInteractiveControlSession()
service._application_control_session = control # type: ignore[assignment] # noqa: SLF001
binding = _seed_supervised_connection(service)
operation_id = "op-00000000-0000-4000-8000-000000009901"
acquisition_id = "acq-00000000-0000-4000-8000-000000009901"
physical = _exact_start_physical_proof(
operation_id=operation_id,
acquisition_id=acquisition_id,
binding=binding,
resolved=False,
)
physical.update(
{
"runtime_bound": True,
"reconciliation_ready": True,
"observed_session_state": "scan_over",
"active_operation_id": operation_id,
}
)
record = physical["record"]
assert isinstance(record, dict)
record["last_status"] = None
verified_control = _verified_control_for_binding(binding)
if invalid_proof == "active-owner":
physical["active_operation_id"] = "op-other-start"
elif invalid_proof == "application-response":
response = record["application_response"]
assert isinstance(response, dict)
response["operation_id"] = "op-other-start"
elif invalid_proof == "binding-generation":
verified_control["producer_generation"] = 2
elif invalid_proof == "qos2":
record["qos2_completed"] = False
else:
record["last_status"] = {
"session_state": "scanning",
"project_bound": True,
"init_ready": True,
"mqtt_retained": True,
"system_error_code": None,
}
assert (
service._matching_accepted_start_scan_over_before_active( # noqa: SLF001
physical,
acquisition_id=acquisition_id,
start_operation_id=operation_id,
verified_control=verified_control,
)
is False
)
def test_active_stream_recovery_reopens_exact_camera_epoch_once_and_fences_stale_lineage( def test_active_stream_recovery_reopens_exact_camera_epoch_once_and_fences_stale_lineage(
tmp_path: Path, tmp_path: Path,
monkeypatch: pytest.MonkeyPatch, monkeypatch: pytest.MonkeyPatch,