feat(observatory): add portable calculation profiles

This commit is contained in:
DCCONSTRUCTIONS
2026-08-31 15:42:56 +03:00
parent d1b75efcea
commit 9beb534108
75 changed files with 24419 additions and 348 deletions
+204 -31
View File
@@ -47,20 +47,38 @@ from k1link.observatory.m49_queue_binding import (
M49QueueBindingError,
M49RecordedQueueBindingService,
)
from k1link.observatory.portable_queue_binding import (
PortableQueueBindingError,
PortableRecordedQueueBindingService,
)
from k1link.observatory.portable_result_contract import (
PortableCalculationProfileRegistry,
PortableResultContractValidatorRegistry,
)
from k1link.observatory.portable_result_publisher import (
resolve_published_portable_calculation_profile,
)
from k1link.observatory.portable_run_definitions import (
PortableRunDefinitionRegistry,
PortableRunDefinitionRegistryError,
)
from k1link.observatory.portable_setup_projection import (
PORTABLE_LAB_V1_SETUP_ID,
PortableLabV1SetupProjector,
PortableSetupProjectionError,
PortableSetupProjector,
portable_calculation_profile_registry,
)
from k1link.observatory.portable_worker_integration import (
PortableObservatoryWorkerIntegration,
PortableWorkerIntegrationError,
PortableWorkerStorageRoots,
build_portable_observatory_worker_integration,
portable_result_validator_registry,
)
from k1link.observatory.recorded_jobs import (
ObservatoryRecordedJobQueue,
ObservatoryRecordedQueueError,
RecordedRunDefinitionRegistry,
)
from k1link.observatory.source_admission import RecordedK1SourceAdmissionService
from k1link.sessions import (
MaterializedRecording,
RecordedCameraFrameService,
@@ -73,6 +91,7 @@ from k1link.sessions import (
SessionRecordingPreparationManager,
SessionStore,
)
from k1link.sessions.models import SessionSummary
from k1link.simulation.projects import SimulationProjectService, SimulationProjectStore
from k1link.web.advanced_laboratory_api import build_advanced_laboratory_router
from k1link.web.artifact_health_api import build_artifact_health_router
@@ -261,6 +280,55 @@ plugin_catalog: DevicePluginCatalog = plugin_environment.catalog
plugin_dispatcher: DevicePluginDispatcher = plugin_environment.dispatcher
session_store = SessionStore(REPOSITORY_ROOT)
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY: PortableRunDefinitionRegistry | None
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY_ERROR: str | None
OBSERVATORY_PORTABLE_CALCULATION_PROFILES: PortableCalculationProfileRegistry | None
OBSERVATORY_PORTABLE_RESULT_VALIDATORS: PortableResultContractValidatorRegistry | None
try:
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY = PortableRunDefinitionRegistry.from_file(
REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json"
)
OBSERVATORY_PORTABLE_CALCULATION_PROFILES = portable_calculation_profile_registry(
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY
)
OBSERVATORY_PORTABLE_RESULT_VALIDATORS = portable_result_validator_registry(
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY
)
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY_ERROR = None
except (
PortableRunDefinitionRegistryError,
PortableWorkerIntegrationError,
OSError,
ValueError,
) as exc:
# Portable definitions are an optional observation-only slice. Registry
# drift cannot affect K1, Simulation, legacy LAB, or the exact M49 binding.
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY = None
OBSERVATORY_PORTABLE_CALCULATION_PROFILES = None
OBSERVATORY_PORTABLE_RESULT_VALIDATORS = None
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY_ERROR = str(exc)
def _resolve_observatory_calculation_profile(
summary: SessionSummary,
) -> dict[str, object] | None:
if OBSERVATORY_LABORATORY_SETUP_REGISTRY is not None:
legacy = OBSERVATORY_LABORATORY_SETUP_REGISTRY.observatory_calculation_profile(
summary
)
if legacy is not None:
return legacy
if (
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY is None
or OBSERVATORY_PORTABLE_CALCULATION_PROFILES is None
):
return None
return resolve_published_portable_calculation_profile(
summary,
definitions=OBSERVATORY_PORTABLE_DEFINITION_REGISTRY,
calculation_profiles=OBSERVATORY_PORTABLE_CALCULATION_PROFILES,
)
def _load_optional_observatory_worker_authentication(
recorded_job_queue: ObservatoryRecordedJobQueue | None,
@@ -299,9 +367,14 @@ try:
REPOSITORY_ROOT / "config" / "observatory-m49-recorded-queue-binding.json"
),
)
recorded_definitions = list(OBSERVATORY_RECORDED_BINDING_SERVICE.definitions.definitions)
if OBSERVATORY_PORTABLE_DEFINITION_REGISTRY is not None:
recorded_definitions.extend(
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY.ready_recorded_definitions()
)
OBSERVATORY_RECORDED_JOB_QUEUE = ObservatoryRecordedJobQueue(
session_store.data_dir,
definitions=OBSERVATORY_RECORDED_BINDING_SERVICE.definitions,
definitions=RecordedRunDefinitionRegistry(tuple(recorded_definitions)),
)
OBSERVATORY_RECORDED_JOB_QUEUE_ERROR = None
except (M49QueueBindingError, ObservatoryRecordedQueueError, OSError, ValueError) as exc:
@@ -316,11 +389,14 @@ OBSERVATORY_WORKER_CLAIM_LEASE_READY = False
OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY = False
OBSERVATORY_WORKER_PRODUCTION_API_ENABLED = False
OBSERVATORY_WORKER_AUTHENTICATION: ObservatoryWorkerAuthentication | None
OBSERVATORY_WORKER_AUTHENTICATION_ERROR: str | None
OBSERVATORY_WORKER_API_ERROR: str | None
OBSERVATORY_WORKER_AUTHENTICATION = None
OBSERVATORY_WORKER_API_ERROR = (
"Worker pull API is hard-disabled until claim leases and a verified "
"Observatory result publisher are implemented and accepted"
(
OBSERVATORY_WORKER_AUTHENTICATION,
OBSERVATORY_WORKER_AUTHENTICATION_ERROR,
) = _load_optional_observatory_worker_authentication(
OBSERVATORY_RECORDED_JOB_QUEUE,
token_path=OBSERVATORY_WORKER_TOKEN_PATH,
)
simulation_project_store = SimulationProjectStore(session_store.data_dir)
simulation_project_service = SimulationProjectService(simulation_project_store)
@@ -336,37 +412,122 @@ session_recording_materializer = SessionRecordingMaterializer(
session_recorded_media_inspector = RecordedMediaInspector(
session_store.data_dir / "recorded-media-preparations"
)
OBSERVATORY_PORTABLE_SETUP_PROJECTOR: PortableLabV1SetupProjector | None
OBSERVATORY_PORTABLE_WORKER_INTEGRATION: PortableObservatoryWorkerIntegration | None
OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR: str | None
OBSERVATORY_PORTABLE_WORKER_STORAGE_ROOTS: PortableWorkerStorageRoots | None = None
try:
if OBSERVATORY_PORTABLE_DEFINITION_REGISTRY is None:
raise PortableWorkerIntegrationError(
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY_ERROR
or "portable definition registry is unavailable"
)
if OBSERVATORY_PORTABLE_CALCULATION_PROFILES is None:
raise PortableWorkerIntegrationError(
"portable calculation profile registry is unavailable"
)
if OBSERVATORY_PORTABLE_RESULT_VALIDATORS is None:
raise PortableWorkerIntegrationError(
"portable result validator registry is unavailable"
)
if OBSERVATORY_RECORDED_JOB_QUEUE is None:
raise PortableWorkerIntegrationError(
OBSERVATORY_RECORDED_JOB_QUEUE_ERROR
or "Observatory recorded-job queue is unavailable"
)
if session_artifact_gateway is None:
raise PortableWorkerIntegrationError(
"central artifact store is not configured"
)
if session_artifact_gateway.status().central_status != "ready":
raise PortableWorkerIntegrationError(
"central artifact store is unavailable"
)
OBSERVATORY_PORTABLE_WORKER_STORAGE_ROOTS = (
PortableWorkerStorageRoots.from_environment(
artifact_store_root=session_artifact_gateway.store.root,
)
)
OBSERVATORY_PORTABLE_WORKER_INTEGRATION = (
build_portable_observatory_worker_integration(
queue=OBSERVATORY_RECORDED_JOB_QUEUE,
session_store=session_store,
media_inspector=session_recorded_media_inspector,
definitions=OBSERVATORY_PORTABLE_DEFINITION_REGISTRY,
artifact_store=session_artifact_gateway.store,
calculation_profiles=OBSERVATORY_PORTABLE_CALCULATION_PROFILES,
validators=OBSERVATORY_PORTABLE_RESULT_VALIDATORS,
source_cas_root=(
OBSERVATORY_PORTABLE_WORKER_STORAGE_ROOTS.source_cas_root
),
result_staging_root=(
OBSERVATORY_PORTABLE_WORKER_STORAGE_ROOTS.result_staging_root
),
)
)
OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR = None
except (PortableWorkerIntegrationError, OSError, ValueError) as exc:
# Constructing this dormant foundation does not enable the Worker router.
# Failure remains isolated from K1, Simulation and legacy LAB.
OBSERVATORY_PORTABLE_WORKER_INTEGRATION = None
OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR = str(exc)
OBSERVATORY_WORKER_API_ERROR = (
"Worker pull API is hard-disabled pending sealed installed executors, "
"a configured Worker credential, explicit integration acceptance and the "
"production gate"
+ (
""
if OBSERVATORY_WORKER_AUTHENTICATION_ERROR is None
else (
"; authentication unavailable: "
f"{OBSERVATORY_WORKER_AUTHENTICATION_ERROR}"
)
)
+ (
""
if OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR is None
else f"; integration unavailable: {OBSERVATORY_PORTABLE_WORKER_INTEGRATION_ERROR}"
)
)
OBSERVATORY_WORKER_DISPATCH_READY = (
OBSERVATORY_WORKER_PRODUCTION_API_ENABLED
and OBSERVATORY_WORKER_CLAIM_LEASE_READY
and OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY
and OBSERVATORY_RECORDED_JOB_QUEUE is not None
and OBSERVATORY_WORKER_AUTHENTICATION is not None
and OBSERVATORY_PORTABLE_WORKER_INTEGRATION is not None
)
OBSERVATORY_PORTABLE_BINDING_SERVICE: PortableRecordedQueueBindingService | None
OBSERVATORY_PORTABLE_SETUP_PROJECTOR: PortableSetupProjector | None
OBSERVATORY_PORTABLE_SETUP_PROJECTOR_ERROR: str | None
try:
portable_definition_registry = PortableRunDefinitionRegistry.from_file(
REPOSITORY_ROOT / "config" / "observatory-portable-run-definitions.json"
)
portable_lab_v1_definition = next(
definition
for definition in portable_definition_registry.definitions
if definition.setup_id == PORTABLE_LAB_V1_SETUP_ID
)
portable_source_capability_service = RecordedK1SourceAdmissionService(
if OBSERVATORY_PORTABLE_DEFINITION_REGISTRY is None:
raise PortableSetupProjectionError(
OBSERVATORY_PORTABLE_DEFINITION_REGISTRY_ERROR
or "portable definition registry is unavailable"
)
OBSERVATORY_PORTABLE_BINDING_SERVICE = PortableRecordedQueueBindingService(
data_dir=session_store.data_dir,
session_store=session_store,
media_inspector=session_recorded_media_inspector,
requirements=portable_lab_v1_definition.to_source_admission_requirements(),
definitions=OBSERVATORY_PORTABLE_DEFINITION_REGISTRY,
queue=OBSERVATORY_RECORDED_JOB_QUEUE,
)
OBSERVATORY_PORTABLE_SETUP_PROJECTOR = PortableLabV1SetupProjector(
registry=portable_definition_registry,
capability_probe=portable_source_capability_service,
OBSERVATORY_PORTABLE_SETUP_PROJECTOR = PortableSetupProjector(
registry=OBSERVATORY_PORTABLE_DEFINITION_REGISTRY,
capability_probe=OBSERVATORY_PORTABLE_BINDING_SERVICE,
dispatch_available=OBSERVATORY_WORKER_DISPATCH_READY,
)
OBSERVATORY_PORTABLE_SETUP_PROJECTOR_ERROR = None
except (
PortableQueueBindingError,
PortableRunDefinitionRegistryError,
PortableSetupProjectionError,
OSError,
StopIteration,
ValueError,
) as exc:
# Portable LAB V1 is an optional observation-only slice. A drifted
# Portable setup execution is an optional observation-only slice. A drifted
# registry cannot affect K1, Simulation, legacy LAB, or the exact M49 queue.
OBSERVATORY_PORTABLE_BINDING_SERVICE = None
OBSERVATORY_PORTABLE_SETUP_PROJECTOR = None
OBSERVATORY_PORTABLE_SETUP_PROJECTOR_ERROR = str(exc)
_ffmpeg = _resolve_media_tool("ffmpeg")
@@ -829,6 +990,14 @@ app.include_router(
perception_overlay_provider=session_perception_overlay_store,
perception_media_provider=session_perception_epoch_store,
point_color_renderers=plugin_environment.point_color_renderers,
lab_calculation_profile_resolver=(
None
if (
OBSERVATORY_LABORATORY_SETUP_REGISTRY is None
and OBSERVATORY_PORTABLE_DEFINITION_REGISTRY is None
)
else _resolve_observatory_calculation_profile
),
)
)
app.include_router(
@@ -843,19 +1012,23 @@ app.include_router(
recorded_job_queue_error=OBSERVATORY_RECORDED_JOB_QUEUE_ERROR,
portable_setup_projector=OBSERVATORY_PORTABLE_SETUP_PROJECTOR,
portable_setup_projector_error=OBSERVATORY_PORTABLE_SETUP_PROJECTOR_ERROR,
portable_binding_service=OBSERVATORY_PORTABLE_BINDING_SERVICE,
)
)
if (
OBSERVATORY_WORKER_PRODUCTION_API_ENABLED
and OBSERVATORY_WORKER_CLAIM_LEASE_READY
and OBSERVATORY_WORKER_VERIFIED_RESULT_PUBLISHER_READY
and OBSERVATORY_RECORDED_JOB_QUEUE is not None
and OBSERVATORY_WORKER_AUTHENTICATION is not None
):
if OBSERVATORY_WORKER_DISPATCH_READY:
assert OBSERVATORY_RECORDED_JOB_QUEUE is not None
assert OBSERVATORY_WORKER_AUTHENTICATION is not None
assert OBSERVATORY_PORTABLE_WORKER_INTEGRATION is not None
app.include_router(
build_observatory_worker_router(
OBSERVATORY_RECORDED_JOB_QUEUE,
authentication=OBSERVATORY_WORKER_AUTHENTICATION,
artifact_transport=(
OBSERVATORY_PORTABLE_WORKER_INTEGRATION.artifact_transport
),
result_publisher=(
OBSERVATORY_PORTABLE_WORKER_INTEGRATION.result_publisher
),
)
)
app.include_router(
+265 -8
View File
@@ -22,9 +22,19 @@ from k1link.observatory.m49_queue_binding import (
M49QueueBindingIntegrityError,
M49RecordedQueueBindingService,
)
from k1link.observatory.portable_queue_binding import (
PortableQueueBindingError,
PortableQueueBindingIntegrityError,
PortableQueueBindingStaleCheckError,
PortableRecordedQueueBindingService,
)
from k1link.observatory.portable_run_definitions import (
PortableRunDefinitionUnavailableError,
)
from k1link.observatory.portable_setup_projection import (
PortableLabV1SetupProjector,
PortableSetupProjectionError,
PortableSetupProjector,
)
from k1link.observatory.recorded_jobs import (
ObservatoryRecordedJobQueue,
@@ -33,6 +43,7 @@ from k1link.observatory.recorded_jobs import (
ObservatoryRecordedQueueError,
ObservatoryRecordedQueueNotFoundError,
)
from k1link.observatory.source_admission import PortableSourceAdmissionError
from k1link.sessions import SessionIntegrityError, SessionNotFoundError, SessionStore
from k1link.sessions.models import SessionSummary
@@ -130,6 +141,8 @@ class ObservatoryRecordedRunSubmitRequest(_StrictApiModel):
max_length=96,
pattern=r"^[a-z][a-z0-9-]{2,95}$",
)
definition_sha256: str | None = Field(default=None, pattern=r"^[a-f0-9]{64}$")
check_sha256: str | None = Field(default=None, pattern=r"^[a-f0-9]{64}$")
def build_observatory_router(
@@ -142,8 +155,9 @@ def build_observatory_router(
recorded_binding_service: M49RecordedQueueBindingService | None = None,
recorded_job_queue: ObservatoryRecordedJobQueue | None = None,
recorded_job_queue_error: str | None = None,
portable_setup_projector: PortableLabV1SetupProjector | None = None,
portable_setup_projector: PortableSetupProjector | PortableLabV1SetupProjector | None = None,
portable_setup_projector_error: str | None = None,
portable_binding_service: PortableRecordedQueueBindingService | None = None,
) -> APIRouter:
"""Build bounded catalog-only mutations for typed Observatory projections."""
@@ -214,6 +228,126 @@ def build_observatory_router(
available.add(result_id)
return frozenset(available)
def portable_run_preflight(
source: SessionSummary,
request: ObservatoryRunPreflightRequest,
) -> dict[str, object] | None:
projector = portable_setup_projector
if projector is None or not projector.has_setup(request.setup_id):
return None
try:
projected = projector.project(source, setup_id=request.setup_id)
except PortableSetupProjectionError as exc:
raise HTTPException(
status_code=503,
detail="Portable-каталог сетапов нарушил контракт целостности.",
) from exc
definition = projected.get("run_definition")
compatibility = projected.get("source_compatibility")
executor = projected.get("executor")
if (
not isinstance(definition, dict)
or not isinstance(compatibility, dict)
or not isinstance(executor, dict)
):
raise HTTPException(
status_code=503,
detail="Portable-каталог сетапов нарушил контракт целостности.",
)
expected_digest = definition.get("definition_sha256")
if request.definition_sha256 != expected_digest:
raise HTTPException(
status_code=409,
detail="Идентичность RunDefinition изменилась; обновите каталог.",
)
compatible = compatibility.get("compatible") is True
executor_ready = executor.get("state") == "ready" and executor.get("ready") is True
checked = None
check_reason: str | None = None
if (
compatible
and executor_ready
and portable_binding_service is not None
and recorded_job_queue is not None
and isinstance(expected_digest, str)
):
try:
checked = portable_binding_service.check(
source_session_id=request.source_session_id,
setup_id=request.setup_id,
definition_sha256=expected_digest,
)
except PortableRunDefinitionUnavailableError as exc:
check_reason = str(exc)
except (PortableQueueBindingError, PortableSourceAdmissionError, ValueError):
check_reason = (
"Источник или исполняемый portable-релиз не прошёл проверку целостности."
)
elif compatible and executor_ready:
check_reason = "Portable dispatch-контур или durable-очередь недоступны."
check_sha256 = None if checked is None else checked.check_sha256
queueable = check_sha256 is not None
checks: list[dict[str, Any]] = [
{
"check_id": "source-compatibility",
"outcome": "pass" if compatible else "fail",
"reason_code": "source-compatible" if compatible else "source-incompatible",
"message": str(compatibility.get("reason")),
},
{
"check_id": "executor",
"outcome": "pass" if executor_ready else "fail",
"reason_code": (
"executor-release-sealed"
if executor_ready
else str(executor.get("reason_code"))
),
"message": (
"Исполняемый portable-релиз и image запечатаны."
if executor_ready
else str(executor.get("reason"))
),
},
{
"check_id": "definition-check",
"outcome": "pass" if checked is not None else "fail",
"reason_code": (
"portable-check-sealed" if checked is not None else "portable-check-unavailable"
),
"message": (
"Источник и RunDefinition связаны одноразовым check SHA."
if checked is not None
else check_reason
or "Portable-проверка недоступна до установки executor-релиза."
),
},
{
"check_id": "durable-queue",
"outcome": "pass" if queueable else "fail",
"reason_code": (
"durable-queue-ready" if queueable else "durable-queue-unavailable"
),
"message": (
"Durable-очередь готова принять расчёт по check SHA."
if queueable
else "Расчёт нельзя поставить в очередь."
),
},
]
return {
"schema_version": OBSERVATORY_RUN_PREFLIGHT_SCHEMA,
"source_session_id": request.source_session_id,
"setup_id": request.setup_id,
"definition_sha256": expected_digest,
"check_sha256": check_sha256,
"outcome": "queueable" if queueable else "blocked",
"submission_allowed": queueable,
"checks": checks,
"existing_result_ids": [],
"executor": executor,
"authority": projected.get("authority", dict(_OBSERVATION_ONLY_AUTHORITY)),
}
if portable_setup_projector is not None:
@router.get("/api/v1/observatory/portable-laboratory-setups")
@@ -230,7 +364,7 @@ def build_observatory_router(
except PortableSetupProjectionError as exc:
raise HTTPException(
status_code=503,
detail="Portable-каталог LAB V1 нарушил контракт целостности.",
detail="Portable-каталог сетапов нарушил контракт целостности.",
) from exc
elif portable_setup_projector_error is not None:
@@ -246,10 +380,10 @@ def build_observatory_router(
del source_session_id
raise HTTPException(
status_code=503,
detail="Portable-каталог LAB V1 недоступен.",
detail="Portable-каталог сетапов недоступен.",
)
if setup_registry is not None:
if setup_registry is not None or portable_setup_projector is not None:
@router.get("/api/v1/observatory/laboratory-setups")
def list_observatory_laboratory_setups(
@@ -259,6 +393,11 @@ def build_observatory_router(
pattern=r"^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$",
),
) -> dict[str, object]:
if setup_registry is None:
raise HTTPException(
status_code=503,
detail="Каталог legacy-сетапов Обсерватории недоступен.",
)
source = source_summary(source_session_id)
return setup_registry.catalog(
source,
@@ -270,6 +409,11 @@ def build_observatory_router(
request: ObservatoryRunPreflightRequest,
) -> dict[str, object]:
source = source_summary(request.source_session_id)
portable = portable_run_preflight(source, request)
if portable is not None:
return portable
if setup_registry is None:
raise HTTPException(status_code=404, detail="Сетап лаборатории не найден.")
try:
setup_registry.setup(request.setup_id)
except KeyError as exc:
@@ -639,16 +783,18 @@ def build_observatory_router(
detail="Подготовка расчётов Обсерватории недоступна.",
)
if (
setup_registry is not None
and recorded_binding_service is not None
and recorded_job_queue is not None
if recorded_job_queue is not None and (
recorded_binding_service is not None or portable_binding_service is not None
):
@router.post("/api/v1/observatory/runs", status_code=202)
def submit_observatory_recorded_run(
request: ObservatoryRecordedRunSubmitRequest,
) -> dict[str, object]:
portable_request = (
portable_setup_projector is not None
and portable_setup_projector.has_setup(request.setup_id)
)
try:
existing_job = recorded_job_queue.get_by_idempotency_key(request.idempotency_key)
except ObservatoryRecordedQueueNotFoundError:
@@ -662,6 +808,10 @@ def build_observatory_router(
if (
existing_job.source_session_id != request.source_session_id
or existing_job.setup_id != request.setup_id
or (
portable_request
and existing_job.definition_sha256 != request.definition_sha256
)
):
raise HTTPException(
status_code=409,
@@ -670,6 +820,113 @@ def build_observatory_router(
return existing_job.as_dict()
source = source_summary(request.source_session_id)
if portable_request:
if portable_binding_service is None:
raise HTTPException(
status_code=503,
detail="Portable dispatch-контур недоступен.",
)
if request.definition_sha256 is None or request.check_sha256 is None:
raise HTTPException(
status_code=409,
detail=(
"Для portable-расчёта требуются актуальные definition SHA "
"и check SHA из preflight."
),
)
assert portable_setup_projector is not None
try:
portable_projection = portable_setup_projector.project(
source,
setup_id=request.setup_id,
)
except PortableSetupProjectionError as exc:
raise HTTPException(
status_code=503,
detail="Portable-каталог сетапов нарушил контракт целостности.",
) from exc
projected_definition = portable_projection.get("run_definition")
projected_executor = portable_projection.get("executor")
if not isinstance(projected_definition, dict) or not isinstance(
projected_executor, dict
):
raise HTTPException(
status_code=503,
detail="Portable-каталог сетапов нарушил контракт целостности.",
)
if projected_definition.get("definition_sha256") != request.definition_sha256:
raise HTTPException(
status_code=409,
detail="Идентичность RunDefinition изменилась; повторите preflight.",
)
if (
projected_executor.get("state") != "ready"
or projected_executor.get("ready") is not True
):
raise HTTPException(
status_code=409,
detail="Portable executor-релиз не установлен.",
)
try:
job, _created = portable_binding_service.submit(
source_session_id=request.source_session_id,
setup_id=request.setup_id,
definition_sha256=request.definition_sha256,
expected_check_sha256=request.check_sha256,
idempotency_key=request.idempotency_key,
)
return job.as_dict()
except PortableQueueBindingStaleCheckError as exc:
raise HTTPException(
status_code=409,
detail="Источник или RunDefinition изменились; повторите preflight.",
) from exc
except PortableRunDefinitionUnavailableError as exc:
raise HTTPException(
status_code=409,
detail="Portable executor-релиз не установлен.",
) from exc
except (
PortableQueueBindingIntegrityError,
PortableSourceAdmissionError,
) as exc:
raise HTTPException(
status_code=409,
detail="Portable-привязка источника не прошла проверку целостности.",
) from exc
except ObservatoryRecordedQueueConflictError as exc:
try:
raced = recorded_job_queue.get_by_idempotency_key(request.idempotency_key)
except ObservatoryRecordedQueueNotFoundError:
raced = None
if (
raced is not None
and raced.source_session_id == request.source_session_id
and raced.setup_id == request.setup_id
and raced.definition_sha256 == request.definition_sha256
):
return raced.as_dict()
raise HTTPException(
status_code=409,
detail="Ключ идемпотентности уже связан с другим расчётом.",
) from exc
except ObservatoryRecordedQueueCapacityError as exc:
raise HTTPException(
status_code=503,
detail="Квота durable-очереди расчётов исчерпана.",
) from exc
except (
PortableQueueBindingError,
ObservatoryRecordedQueueError,
ValueError,
) as exc:
raise HTTPException(
status_code=503,
detail="Portable dispatch-контур недоступен.",
) from exc
if setup_registry is None or recorded_binding_service is None:
raise HTTPException(status_code=404, detail="Сетап лаборатории не найден.")
try:
setup_registry.setup(request.setup_id)
except KeyError as exc:
+292 -3
View File
@@ -18,11 +18,23 @@ from dataclasses import dataclass
from pathlib import Path
from typing import Annotated, Final, Literal
from fastapi import APIRouter, Depends, Header, HTTPException, Response
from fastapi import APIRouter, Depends, Header, HTTPException, Request, Response
from fastapi import Path as ApiPath
from fastapi.responses import FileResponse
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
from pydantic import BaseModel, ConfigDict, Field
from k1link.observatory.portable_artifact_transport import (
MAX_RESULT_MANIFEST_BYTES,
PortableArtifactTransportError,
PortableArtifactTransportIntegrityError,
PortableArtifactTransportUnavailableError,
PortableObservatoryArtifactTransport,
)
from k1link.observatory.portable_result_contract import PortableResultPublisherError
from k1link.observatory.portable_result_publisher import (
PortableObservatoryResultPublisher,
)
from k1link.observatory.recorded_jobs import (
ObservatoryRecordedCheckpointError,
ObservatoryRecordedJobQueue,
@@ -38,6 +50,7 @@ from k1link.observatory.recorded_jobs import (
OBSERVATORY_WORKER_CLAIM_REQUEST_SCHEMA: Final = "missioncore.observatory-worker-claim-request/v1"
OBSERVATORY_WORKER_START_REQUEST_SCHEMA: Final = "missioncore.observatory-worker-start-request/v1"
OBSERVATORY_WORKER_RENEW_REQUEST_SCHEMA: Final = "missioncore.observatory-worker-renew-request/v1"
OBSERVATORY_WORKER_CHECKPOINT_REQUEST_SCHEMA: Final = (
"missioncore.observatory-worker-checkpoint-request/v1"
)
@@ -46,6 +59,10 @@ OBSERVATORY_WORKER_SUCCEED_REQUEST_SCHEMA: Final = (
)
OBSERVATORY_WORKER_FAIL_REQUEST_SCHEMA: Final = "missioncore.observatory-worker-fail-request/v1"
OBSERVATORY_WORKER_CONTOUR_HEADER: Final = "X-Mission-Core-Contour-Id"
OBSERVATORY_WORKER_CLAIM_TOKEN_HEADER: Final = "X-Mission-Core-Claim-Token"
OBSERVATORY_WORKER_CLAIM_GENERATION_HEADER: Final = (
"X-Mission-Core-Claim-Generation"
)
_IDENTIFIER = re.compile(r"^[a-z][a-z0-9-]{2,95}$")
_SHA256 = re.compile(r"^[a-f0-9]{64}$")
@@ -138,6 +155,13 @@ class ObservatoryWorkerStartRequest(_StrictWorkerRequest):
claim_token: str = Field(pattern=_CLAIM_TOKEN_PATTERN)
class ObservatoryWorkerRenewRequest(_StrictWorkerRequest):
schema_version: Literal["missioncore.observatory-worker-renew-request/v1"]
claim_token: str = Field(pattern=_CLAIM_TOKEN_PATTERN)
claim_generation: int = Field(ge=1)
heartbeat_sequence: int = Field(ge=1)
class ObservatoryWorkerCheckpointRequest(_StrictWorkerRequest):
schema_version: Literal["missioncore.observatory-worker-checkpoint-request/v1"]
claim_token: str = Field(pattern=_CLAIM_TOKEN_PATTERN)
@@ -174,6 +198,8 @@ def build_observatory_worker_router(
queue: ObservatoryRecordedJobQueue,
*,
authentication: ObservatoryWorkerAuthentication,
artifact_transport: PortableObservatoryArtifactTransport | None = None,
result_publisher: PortableObservatoryResultPublisher | None = None,
) -> APIRouter:
"""Build the bounded Worker pull/state-transition router.
@@ -182,6 +208,9 @@ def build_observatory_worker_router(
hashed, and is compared to the configured digest in constant time.
"""
if result_publisher is not None and artifact_transport is None:
raise ValueError("portable result publisher requires artifact transport")
def require_configured_worker(
credentials: Annotated[
HTTPAuthorizationCredentials | None,
@@ -245,6 +274,20 @@ def build_observatory_worker_router(
) -> dict[str, object]:
return _queue_call(lambda: queue.start(job_id, claim_token=request.claim_token)).as_dict()
@router.post("/recorded-jobs/{job_id}/lease/renew")
def renew_job_claim(
request: ObservatoryWorkerRenewRequest,
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
) -> dict[str, object]:
return _queue_call(
lambda: queue.renew_claim(
job_id,
claim_token=request.claim_token,
claim_generation=request.claim_generation,
heartbeat_sequence=request.heartbeat_sequence,
)
).as_dict()
@router.post("/recorded-jobs/{job_id}/checkpoint")
def checkpoint_job(
request: ObservatoryWorkerCheckpointRequest,
@@ -263,14 +306,39 @@ def build_observatory_worker_router(
request: ObservatoryWorkerSucceedRequest,
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
) -> dict[str, object]:
return _queue_call(
if artifact_transport is not None:
_artifact_call(
lambda: artifact_transport.require_completed_for_success(
job_id=job_id,
result_id=request.result_id,
result_sha256=request.result_sha256,
claim_token=request.claim_token,
claimant_id=authentication.contour_id,
)
)
succeeded = _queue_call(
lambda: queue.succeed(
job_id,
claim_token=request.claim_token,
result_id=request.result_id,
result_sha256=request.result_sha256,
)
).as_dict()
)
if artifact_transport is not None and result_publisher is not None:
package_root = _artifact_call(
lambda: artifact_transport.package_root_for_terminal(succeeded)
)
try:
result_publisher.publish(job=succeeded, package_root=package_root)
except PortableResultPublisherError as exc:
raise HTTPException(
status_code=503,
detail=(
"Recorded result is sealed but its verified publication "
"requires reconciliation."
),
) from exc
return succeeded.as_dict()
@router.post("/recorded-jobs/{job_id}/fail")
def fail_job(
@@ -286,6 +354,164 @@ def build_observatory_worker_router(
)
).as_dict()
if artifact_transport is not None:
@router.get("/recorded-jobs/{job_id}/source-materialization")
def source_materialization(
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
claim_token: Annotated[
str,
Header(
alias=OBSERVATORY_WORKER_CLAIM_TOKEN_HEADER,
pattern=_CLAIM_TOKEN_PATTERN,
),
],
claim_generation: Annotated[
int,
Header(alias=OBSERVATORY_WORKER_CLAIM_GENERATION_HEADER, ge=1),
],
) -> dict[str, object]:
return _artifact_call(
lambda: artifact_transport.source_manifest(
job_id=job_id,
claim_token=claim_token,
claim_generation=claim_generation,
claimant_id=authentication.contour_id,
)
.as_dict()
)
@router.get("/recorded-jobs/{job_id}/source-members/{member_id}")
def source_member(
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
member_id: Annotated[str, ApiPath(pattern=r"^[a-f0-9]{64}$")],
claim_token: Annotated[
str,
Header(
alias=OBSERVATORY_WORKER_CLAIM_TOKEN_HEADER,
pattern=_CLAIM_TOKEN_PATTERN,
),
],
claim_generation: Annotated[
int,
Header(alias=OBSERVATORY_WORKER_CLAIM_GENERATION_HEADER, ge=1),
],
) -> FileResponse:
member, path = _artifact_call(
lambda: artifact_transport.materialize_source_member(
job_id=job_id,
member_id=member_id,
claim_token=claim_token,
claim_generation=claim_generation,
claimant_id=authentication.contour_id,
)
)
return FileResponse(
path,
media_type=member.media_type,
headers={
"ETag": f'"{member.sha256}"',
"Cache-Control": "private, no-store",
"X-Content-Type-Options": "nosniff",
"X-Mission-Core-Content-Sha256": member.sha256,
},
)
@router.put(
"/recorded-jobs/{job_id}/result-packages/{result_sha256}/manifest"
)
async def stage_result_manifest(
request: Request,
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
result_sha256: Annotated[str, ApiPath(pattern=r"^[a-f0-9]{64}$")],
claim_token: Annotated[
str,
Header(
alias=OBSERVATORY_WORKER_CLAIM_TOKEN_HEADER,
pattern=_CLAIM_TOKEN_PATTERN,
),
],
claim_generation: Annotated[
int,
Header(alias=OBSERVATORY_WORKER_CLAIM_GENERATION_HEADER, ge=1),
],
) -> dict[str, object]:
payload = await _read_bounded_body(request, MAX_RESULT_MANIFEST_BYTES)
return _artifact_call(
lambda: artifact_transport.stage_result_manifest(
job_id=job_id,
result_sha256=result_sha256,
manifest_payload=payload,
claim_token=claim_token,
claim_generation=claim_generation,
claimant_id=authentication.contour_id,
)
.as_dict()
)
@router.put(
"/recorded-jobs/{job_id}/result-packages/{result_sha256}/members/{member_id}"
)
async def upload_result_member(
request: Request,
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
result_sha256: Annotated[str, ApiPath(pattern=r"^[a-f0-9]{64}$")],
member_id: Annotated[str, ApiPath(pattern=r"^[a-f0-9]{64}$")],
claim_token: Annotated[
str,
Header(
alias=OBSERVATORY_WORKER_CLAIM_TOKEN_HEADER,
pattern=_CLAIM_TOKEN_PATTERN,
),
],
claim_generation: Annotated[
int,
Header(alias=OBSERVATORY_WORKER_CLAIM_GENERATION_HEADER, ge=1),
],
) -> dict[str, object]:
try:
plan = await artifact_transport.upload_result_member(
job_id=job_id,
result_sha256=result_sha256,
member_id=member_id,
chunks=request.stream(),
claim_token=claim_token,
claim_generation=claim_generation,
claimant_id=authentication.contour_id,
)
except Exception as exc:
_raise_artifact_or_queue_error(exc)
return plan.as_dict()
@router.post(
"/recorded-jobs/{job_id}/result-packages/{result_sha256}/complete"
)
def complete_result_package(
job_id: Annotated[str, ApiPath(pattern=_JOB_ID_PATTERN)],
result_sha256: Annotated[str, ApiPath(pattern=r"^[a-f0-9]{64}$")],
claim_token: Annotated[
str,
Header(
alias=OBSERVATORY_WORKER_CLAIM_TOKEN_HEADER,
pattern=_CLAIM_TOKEN_PATTERN,
),
],
claim_generation: Annotated[
int,
Header(alias=OBSERVATORY_WORKER_CLAIM_GENERATION_HEADER, ge=1),
],
) -> dict[str, object]:
return _artifact_call(
lambda: artifact_transport.complete_result_upload(
job_id=job_id,
result_sha256=result_sha256,
claim_token=claim_token,
claim_generation=claim_generation,
claimant_id=authentication.contour_id,
)
.as_dict()
)
return router
@@ -341,3 +567,66 @@ def _queue_call[T](operation: Callable[[], T]) -> T:
status_code=503,
detail="Recorded-job queue is unavailable.",
) from exc
def _artifact_call[T](operation: Callable[[], T]) -> T:
try:
return _queue_call(operation)
except PortableArtifactTransportIntegrityError as exc:
raise HTTPException(
status_code=409,
detail="Worker artifact identity was rejected.",
) from exc
except PortableArtifactTransportUnavailableError as exc:
raise HTTPException(
status_code=409,
detail="Worker artifact member is unavailable for this claim.",
) from exc
except PortableArtifactTransportError as exc:
raise HTTPException(
status_code=503,
detail="Worker artifact transport is unavailable.",
) from exc
def _raise_artifact_or_queue_error(exc: Exception) -> None:
if isinstance(exc, PortableArtifactTransportIntegrityError):
raise HTTPException(
status_code=409,
detail="Worker artifact identity was rejected.",
) from exc
if isinstance(exc, PortableArtifactTransportUnavailableError):
raise HTTPException(
status_code=409,
detail="Worker artifact member is unavailable for this claim.",
) from exc
if isinstance(exc, PortableArtifactTransportError):
raise HTTPException(
status_code=503,
detail="Worker artifact transport is unavailable.",
) from exc
_queue_call(lambda: _raise(exc))
raise AssertionError("unreachable")
def _raise(exc: Exception) -> None:
raise exc
async def _read_bounded_body(request: Request, maximum_bytes: int) -> bytes:
content_length = request.headers.get("content-length")
if content_length is not None:
try:
declared = int(content_length)
except ValueError as exc:
raise HTTPException(status_code=400, detail="Content-Length is invalid.") from exc
if declared < 1 or declared > maximum_bytes:
raise HTTPException(status_code=413, detail="Request body is outside bounds.")
payload = bytearray()
async for chunk in request.stream():
payload.extend(chunk)
if len(payload) > maximum_bytes:
raise HTTPException(status_code=413, detail="Request body is outside bounds.")
if not payload:
raise HTTPException(status_code=400, detail="Request body is empty.")
return bytes(payload)
+24 -5
View File
@@ -41,6 +41,7 @@ from k1link.sessions.canonical_lab_spatial import (
CANONICAL_LAB_SPATIAL_PROFILE,
canonical_lab_spatial_frame,
)
from k1link.sessions.models import SessionSummary
from k1link.sessions.plugin_contract import RecordedPointColorRenderer
from k1link.viewer.recorded import (
APPLICATION_ID as RECORDED_APPLICATION_ID,
@@ -328,6 +329,9 @@ def build_session_router(
perception_overlay_provider: RecordedPerceptionOverlayProvider | None = None,
perception_media_provider: RecordedPerceptionMediaProvider | None = None,
point_color_renderers: Mapping[str, RecordedPointColorRenderer] | None = None,
lab_calculation_profile_resolver: (
Callable[[SessionSummary], Mapping[str, object] | None] | None
) = None,
allow_synchronous_recording_fallback: bool = False,
replay_action_id: str = DEFAULT_REPLAY_ACTION_ID,
) -> APIRouter:
@@ -336,12 +340,29 @@ def build_session_router(
router = APIRouter(tags=["observation-sessions"])
recorded_media_inspector = media_inspector or RecordedMediaInspector()
def lab_catalog_document(
summary: SessionSummary,
contract: Literal["v1", "v2", "v3"],
) -> dict[str, Any]:
lab = summary.lab
if lab is None:
raise ValueError("LAB catalog document requires a LAB summary")
document = lab.as_dict(include_replay_capability=contract in ("v2", "v3"))
if contract == "v3":
profile = (
None
if lab_calculation_profile_resolver is None
else lab_calculation_profile_resolver(summary)
)
document["calculation_profile"] = None if profile is None else dict(profile)
return document
@router.get("/api/v1/observation-sessions")
def list_observation_sessions(
limit: int = Query(default=20, ge=1, le=100),
cursor: str | None = Query(default=None, max_length=128),
scope: Literal["all", "source", "laboratory"] = "all",
lab_contract: Literal["v1", "v2"] = "v1",
lab_contract: Literal["v1", "v2", "v3"] = "v1",
) -> dict[str, Any]:
try:
_refresh_catalog(catalog_refresher)
@@ -349,7 +370,7 @@ def build_session_router(
limit=limit,
cursor=cursor,
scope=scope,
include_capability_projections=lab_contract == "v2",
include_capability_projections=lab_contract in ("v2", "v3"),
)
return {
"items": [
@@ -364,9 +385,7 @@ def build_session_router(
"replayable": item.replayable,
**(
{
"lab": item.lab.as_dict(
include_replay_capability=lab_contract == "v2"
)
"lab": lab_catalog_document(item, lab_contract)
}
if item.lab is not None
else {}