feat(lab): publish E31 through E33 results

This commit is contained in:
DCCONSTRUCTIONS
2026-07-27 14:42:22 +03:00
parent fa2d66ab0e
commit 2c016b117d
11 changed files with 1544 additions and 2 deletions
@@ -142,6 +142,25 @@ class _E30Chain:
human_decisions: tuple[dict[str, Any], ...]
def read_e31_source_qualification(root: Path) -> E31SourceQualification:
"""Open and fully validate one immutable E31 result."""
resolved = _safe_result_root(root, _RESULT_ID)
manifest = _read_json(resolved / E31_MANIFEST_NAME)
identity = _object(manifest.get("identity"), "E31 identity")
result = _read_existing(resolved, identity)
if (
result.report.get("schema_version") != E31_SOURCE_QUALIFICATION_REPORT_SCHEMA
or result.report.get("result_id") != resolved.name
or result.report.get("authority") != _diagnostic_authority()
or result.manifest.get("status") != result.report.get("status")
or result.manifest.get("eligible_for_e32")
!= result.report.get("eligible_for_e32")
):
raise E31SourceQualificationError("E31 result report is inconsistent")
return result
def build_e31_source_qualification(
*,
source_pack_root: Path,
+418
View File
@@ -0,0 +1,418 @@
from __future__ import annotations
import copy
import re
from collections.abc import Callable
from functools import lru_cache
from pathlib import Path
from typing import Any, Final
from fastapi import APIRouter, Query
from k1link.compute.e31_source_qualification import (
E31SourceQualification,
E31SourceQualificationError,
read_e31_source_qualification,
)
from k1link.compute.e32_track_geometry_replay import (
E32TrackGeometryReplay,
E32TrackGeometryReplayError,
read_e32_track_geometry_replay,
)
from k1link.compute.e33_worker_shadow import (
E33WorkerShadowError,
E33WorkerShadowResult,
read_e33_worker_shadow_result,
)
LABORATORY_ADVANCED_CATALOG_SCHEMA: Final = (
"missioncore.laboratory-advanced-catalog/v1"
)
_E31_RESULT_ID = re.compile(r"^e31-source-qualification-[a-f0-9]{64}$")
_E32_RESULT_ID = re.compile(r"^e32-track-geometry-[a-f0-9]{64}$")
_E33_RESULT_ID = re.compile(r"^e33-worker-shadow-[a-f0-9]{64}$")
RootProvider = Callable[[], Path | None]
def _result_signature(root: Path) -> tuple[int, ...]:
signature: list[int] = []
for path in sorted(root.iterdir(), key=lambda item: item.name):
if not path.is_file() or path.is_symlink():
continue
stat = path.stat()
signature.extend((stat.st_size, stat.st_mtime_ns))
return tuple(signature)
@lru_cache(maxsize=16)
def _read_e31_cached(
root_text: str,
signature: tuple[int, ...],
) -> E31SourceQualification:
del signature
return read_e31_source_qualification(Path(root_text))
@lru_cache(maxsize=16)
def _read_e32_cached(
root_text: str,
signature: tuple[int, ...],
) -> E32TrackGeometryReplay:
del signature
return read_e32_track_geometry_replay(Path(root_text))
@lru_cache(maxsize=16)
def _read_e33_cached(
root_text: str,
signature: tuple[int, ...],
) -> E33WorkerShadowResult:
del signature
return read_e33_worker_shadow_result(Path(root_text))
def _configured_root(provider: RootProvider) -> Path | None:
value = provider()
if value is None:
return None
candidate = value.expanduser().absolute()
if candidate.is_symlink():
return None
try:
root = candidate.resolve(strict=True)
except OSError:
return None
if not root.is_dir():
return None
return root
def _candidates(root: Path, pattern: re.Pattern[str]) -> list[Path]:
return sorted(
(
candidate
for candidate in root.iterdir()
if candidate.is_dir()
and not candidate.is_symlink()
and pattern.fullmatch(candidate.name) is not None
),
key=lambda candidate: candidate.stat().st_mtime_ns,
reverse=True,
)
def _object(value: object, label: str) -> dict[str, Any]:
if not isinstance(value, dict):
raise ValueError(f"{label} is invalid")
return value
def _project_e31(result: E31SourceQualification) -> dict[str, object]:
identity = _object(result.manifest.get("identity"), "E31 identity")
source = _object(identity.get("source"), "E31 source")
method = _object(identity.get("method"), "E31 method")
source_time = _object(result.report.get("source_time"), "E31 source time")
lidar_camera = _object(
source_time.get("lidar_camera_abs_delta_ms"),
"E31 LiDAR/camera timing",
)
pose_point = _object(
source_time.get("pose_point_abs_delta_ms"),
"E31 pose/point timing",
)
offset = _object(result.report.get("offset_sensitivity"), "E31 offset")
self_mask = _object(result.report.get("self_mask"), "E31 self-mask")
semantic_mask = _object(self_mask.get("semantic_mask"), "E31 semantic mask")
geometry_mask = _object(
self_mask.get("geometry_point_mask"),
"E31 geometry mask",
)
return {
"result_id": result.result_id,
"created_at_utc": result.manifest.get("created_at_utc"),
"source_session_id": source.get("session_id"),
"status": result.report.get("status"),
"eligible_for_e32": result.report.get("eligible_for_e32"),
"profile_id": method.get("profile_id"),
"producer_sha256": identity.get("producer_sha256"),
"metrics": {
"frame_count": source_time.get("frame_count"),
"available_binding_count": source_time.get("available_binding_count"),
"available_fraction": source_time.get("available_fraction"),
"lidar_camera_p95_ms": lidar_camera.get("p95"),
"pose_point_p95_ms": pose_point.get("p95"),
"selected_offset_ms": offset.get("selected_offset_ms"),
"correspondence_count": offset.get("correspondence_count"),
"supported_fraction": offset.get("baseline_supported_fraction"),
"semantic_self_sample_count": semantic_mask.get("source_sample_count"),
"semantic_self_collateral_count": semantic_mask.get(
"collateral_item_count"
),
"exact_geometry_correction_count": len(
self_mask.get("exact_correction_item_ids", [])
),
"geometry_point_mask_status": geometry_mask.get("status"),
},
"limitations": copy.deepcopy(result.report.get("limitations")),
"authority": copy.deepcopy(result.report.get("authority")),
"access": "read-only",
}
def _project_e32(result: E32TrackGeometryReplay) -> dict[str, object]:
identity = _object(result.manifest.get("identity"), "E32 identity")
source = _object(identity.get("source"), "E32 source")
method = _object(identity.get("method"), "E32 method")
binding = _object(
identity.get("track_geometry_binding"),
"E32 track geometry binding",
)
metrics = _object(result.report.get("metrics"), "E32 metrics")
frames = _object(metrics.get("frames"), "E32 frames")
objects = _object(metrics.get("objects"), "E32 objects")
ownership = _object(metrics.get("point_ownership"), "E32 point ownership")
point_rows = _object(metrics.get("qualified_point_rows"), "E32 point rows")
runtime = _object(metrics.get("runtime"), "E32 runtime")
frame_processing = _object(
runtime.get("frame_processing_ms"),
"E32 frame processing",
)
return {
"result_id": result.result_id,
"created_at_utc": result.manifest.get("created_at_utc"),
"source_session_id": binding.get("source_session_id"),
"status": result.report.get("status"),
"profile_id": method.get("profile_id"),
"producer_sha256": identity.get("producer_sha256"),
"e31_result_id": source.get("e31_result_id"),
"metrics": {
"frames_total": frames.get("total"),
"frames_source_available": frames.get("source_available"),
"semantic_published": objects.get("semantic_published"),
"semantic_masked": objects.get("semantic_masked"),
"geometry_published": objects.get("geometry_published"),
"observations_arbitrated": ownership.get("observations_arbitrated"),
"overlapping_claims_removed": ownership.get(
"overlapping_claims_removed"
),
"qualified_points_published": point_rows.get("e32_published"),
"qualified_points_withheld": point_rows.get(
"unqualified_ranges_withheld"
),
"frame_processing_p95_ms": frame_processing.get("p95"),
"build_elapsed_ms": runtime.get("build_elapsed_ms"),
},
"decision": copy.deepcopy(result.report.get("decision")),
"authority": copy.deepcopy(result.report.get("authority")),
"access": "read-only",
}
def _project_e33(
result: E33WorkerShadowResult,
linked_e32: E32TrackGeometryReplay,
) -> dict[str, object]:
identity = _object(result.result.get("identity"), "E33 identity")
linked_identity = _object(
linked_e32.manifest.get("identity"),
"linked E32 identity",
)
binding = _object(
linked_identity.get("track_geometry_binding"),
"linked E32 binding",
)
profile = _object(identity.get("profile"), "E33 profile")
worker = _object(identity.get("worker"), "E33 worker")
metrics = _object(result.report.get("metrics"), "E33 metrics")
accounting = _object(metrics.get("accounting"), "E33 accounting")
processing = _object(metrics.get("processing_ms"), "E33 processing")
release_lag = _object(metrics.get("release_lag_ms"), "E33 release lag")
result_age = _object(metrics.get("result_age_ms"), "E33 result age")
resources = _object(metrics.get("resources"), "E33 resources")
process_rss = _object(resources.get("process_rss_mib"), "E33 RSS")
gpu_utilization = _object(
resources.get("gpu_utilization_percent"),
"E33 GPU utilization",
)
work_queue = _object(metrics.get("work_queue"), "E33 work queue")
result_queue = _object(metrics.get("result_queue"), "E33 result queue")
return {
"result_id": result.result_id,
"created_at_utc": result.result.get("created_at_utc"),
"source_session_id": binding.get("source_session_id"),
"status": result.result.get("acceptance_state"),
"e32_result_id": identity.get("e32_result_id"),
"pipeline_id": identity.get("pipeline"),
"mode": profile.get("mode"),
"worker": {
"node": worker.get("worker_node"),
"container_image": worker.get("container_image"),
"python": worker.get("python"),
"numpy": worker.get("numpy"),
},
"metrics": {
"source_frames": accounting.get("source_frames"),
"delivered_frames": accounting.get("delivered"),
"input_superseded": accounting.get("input_superseded"),
"result_superseded": accounting.get("result_superseded"),
"effective_delivery_fps": metrics.get("effective_delivery_fps"),
"deadline_miss_fraction": metrics.get("deadline_miss_fraction"),
"processing_p95_ms": processing.get("p95"),
"release_lag_p95_ms": release_lag.get("p95"),
"result_age_p95_ms": result_age.get("p95"),
"process_rss_p95_mib": process_rss.get("p95"),
"gpu_utilization_p95_percent": gpu_utilization.get("p95"),
"gpu_visible": resources.get("gpu_visible"),
"work_queue_capacity": work_queue.get("capacity"),
"result_queue_capacity": result_queue.get("capacity"),
"wall_to_ideal_ratio": metrics.get("wall_to_ideal_ratio"),
},
"acceptance": copy.deepcopy(result.report.get("acceptance")),
"authority": copy.deepcopy(result.report.get("authority")),
"access": "read-only",
}
def _empty_catalog(configured: bool) -> dict[str, object]:
return {
"schema_version": LABORATORY_ADVANCED_CATALOG_SCHEMA,
"configured": configured,
"items": [],
"candidate_total": 0,
"invalid_total": 0,
"access": "read-only",
}
def build_advanced_laboratory_router(
*,
e31_root_provider: RootProvider = lambda: None,
e32_root_provider: RootProvider = lambda: None,
e33_root_provider: RootProvider = lambda: None,
) -> APIRouter:
router = APIRouter(prefix="/api/v1/laboratory", tags=["laboratory"])
@router.get("/e31/results")
def list_e31_results(
limit: int = Query(default=1, ge=1, le=10),
) -> dict[str, object]:
root = _configured_root(e31_root_provider)
if root is None:
return _empty_catalog(False)
candidates = _candidates(root, _E31_RESULT_ID)
items: list[dict[str, object]] = []
invalid_total = 0
for candidate in candidates:
try:
result = _read_e31_cached(
str(candidate.resolve()),
_result_signature(candidate),
)
if result.report.get("eligible_for_e32") is not True:
raise ValueError("E31 result is not accepted")
if len(items) < limit:
items.append(_project_e31(result))
except (
E31SourceQualificationError,
KeyError,
OSError,
TypeError,
ValueError,
):
invalid_total += 1
return {
**_empty_catalog(True),
"items": items,
"candidate_total": len(candidates),
"invalid_total": invalid_total,
}
@router.get("/e32/results")
def list_e32_results(
limit: int = Query(default=1, ge=1, le=10),
) -> dict[str, object]:
root = _configured_root(e32_root_provider)
if root is None:
return _empty_catalog(False)
candidates = _candidates(root, _E32_RESULT_ID)
items: list[dict[str, object]] = []
invalid_total = 0
for candidate in candidates:
try:
result = _read_e32_cached(
str(candidate.resolve()),
_result_signature(candidate),
)
if (
result.report.get("status")
!= "accepted-diagnostic-track-geometry-replay"
):
raise ValueError("E32 result is not accepted")
if len(items) < limit:
items.append(_project_e32(result))
except (
E32TrackGeometryReplayError,
KeyError,
OSError,
TypeError,
ValueError,
):
invalid_total += 1
return {
**_empty_catalog(True),
"items": items,
"candidate_total": len(candidates),
"invalid_total": invalid_total,
}
@router.get("/e33/results")
def list_e33_results(
limit: int = Query(default=1, ge=1, le=10),
) -> dict[str, object]:
root = _configured_root(e33_root_provider)
e32_root = _configured_root(e32_root_provider)
if root is None or e32_root is None:
return _empty_catalog(False)
candidates = _candidates(root, _E33_RESULT_ID)
items: list[dict[str, object]] = []
invalid_total = 0
for candidate in candidates:
try:
result = _read_e33_cached(
str(candidate.resolve()),
_result_signature(candidate),
)
identity = _object(result.result.get("identity"), "E33 identity")
e32_result_id = identity.get("e32_result_id")
if (
not result.accepted
or not isinstance(e32_result_id, str)
or _E32_RESULT_ID.fullmatch(e32_result_id) is None
):
raise ValueError("E33 result is not accepted")
linked_root = e32_root / e32_result_id
linked_e32 = _read_e32_cached(
str(linked_root.resolve(strict=True)),
_result_signature(linked_root),
)
if len(items) < limit:
items.append(_project_e33(result, linked_e32))
except (
E32TrackGeometryReplayError,
E33WorkerShadowError,
KeyError,
OSError,
TypeError,
ValueError,
):
invalid_total += 1
return {
**_empty_catalog(True),
"items": items,
"candidate_total": len(candidates),
"invalid_total": invalid_total,
}
return router
+26
View File
@@ -32,6 +32,7 @@ from k1link.sessions import (
SessionRecordingPreparationManager,
SessionStore,
)
from k1link.web.advanced_laboratory_api import build_advanced_laboratory_router
from k1link.web.device_plugin_composition import load_installed_device_plugins
from k1link.web.e30_engineering_api import build_e30_engineering_router
from k1link.web.e30_human_review_api import build_e30_human_review_router
@@ -496,6 +497,31 @@ app.include_router(
),
)
)
app.include_router(
build_advanced_laboratory_router(
e31_root_provider=lambda: (
REPOSITORY_ROOT
/ ".runtime"
/ "compute-experiments"
/ "e31"
/ "source-qualifications"
),
e32_root_provider=lambda: (
REPOSITORY_ROOT
/ ".runtime"
/ "compute-experiments"
/ "e32"
/ "results"
),
e33_root_provider=lambda: (
REPOSITORY_ROOT
/ ".runtime"
/ "compute-experiments"
/ "e33"
/ "results"
),
)
)
app.include_router(
build_e30_review_router(
materialization_root_provider=lambda: (