Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
116 changes: 111 additions & 5 deletions backend/cortex_backend/api/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,12 @@
from cortex_backend.execution.coordinator import DurableFakeCoordinator
from cortex_backend.execution.fake import FakeExecutionPlan
from cortex_backend.execution.models import ExecutionJob, ExecutionEvent, TerminalExecutionStatus
from cortex_backend.execution.recipe_coordinator import (
RECIPE_IMAGE_PROFILE,
RecipeExecutionCoordinator,
RecipeExecutionError,
RecipeImageRequest,
)
from cortex_backend.execution.repository import (
ApprovalPolicyError,
ApprovalTransitionError,
Expand All @@ -49,6 +55,8 @@
ExecutionAccepted,
ExecutionApprovalDecisionRequest,
ExecutionPreviewRequest,
RecipeImageTransformAccepted,
RecipeImageTransformRequest,
ExecutionSSEEvent,
ExecutionStatusResponse,
ExecutionTaskListResponse,
Expand Down Expand Up @@ -408,7 +416,7 @@ def start_fake_execution(
payload: ExecutionPreviewRequest,
principal: SessionPrincipal = Depends(require_session),
) -> ExecutionAccepted:
coordinator = _execution_coordinator(request)
coordinator = _fake_execution_coordinator(request)
try:
job = coordinator.start(
owner=_execution_owner(principal),
Expand All @@ -429,6 +437,50 @@ def start_fake_execution(
sequence=job.sequence,
)

@router.post(
"/execution/recipe/image",
response_model=RecipeImageTransformAccepted,
status_code=status.HTTP_202_ACCEPTED,
)
def start_recipe_image_transform(
request: Request,
payload: RecipeImageTransformRequest,
principal: SessionPrincipal = Depends(require_session),
) -> RecipeImageTransformAccepted:
"""Start one explicitly qualified, owner-scoped image recipe.

The route is intentionally unavailable unless the app was built with
the explicit qualification lifecycle and that lifecycle completed its
health-gated startup. Attachment staging is a separate trusted boundary;
callers provide only its opaque artifact identifier here.
"""

coordinator = _recipe_coordinator(request)
try:
job = coordinator.start_image_transform(
RecipeImageRequest(
owner=_execution_owner(principal),
request_id=payload.request_id,
source_artifact_id=payload.source_artifact_id,
plan=payload.plan,
retention_seconds=payload.retention_seconds,
)
)
except RecipeExecutionError as exc:
_raise_recipe_request_error(exc)
except (TypeError, ValueError) as exc:
raise HTTPException(
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
detail="Recipe request is invalid.",
) from exc
return RecipeImageTransformAccepted(
job_id=job.job_id,
request_id=job.request_id,
profile=RECIPE_IMAGE_PROFILE,
status=job.status,
sequence=job.sequence,
)

@router.get("/execution/tasks", response_model=ExecutionTaskListResponse)
def execution_tasks(
request: Request,
Expand Down Expand Up @@ -500,7 +552,7 @@ def cancel_execution(
request: Request,
principal: SessionPrincipal = Depends(require_session),
) -> ExecutionStatusResponse:
coordinator = _execution_coordinator(request)
coordinator = _execution_runtime(request)
try:
coordinator.cancel(job_id, owner=_execution_owner(principal))
except ValueError as exc:
Expand Down Expand Up @@ -1388,8 +1440,8 @@ def _execution_owner(principal: SessionPrincipal) -> str:
return principal.installation_principal_id


def _execution_coordinator(request: Request) -> DurableFakeCoordinator:
"""Require the explicitly injected fake-only preview coordinator."""
def _execution_runtime(request: Request):
"""Require an explicitly enabled execution runtime for shared job routes."""
coordinator = getattr(request.app.state, "execution_coordinator", None)
if not request.app.state.preview or coordinator is None:
raise HTTPException(
Expand All @@ -1399,8 +1451,62 @@ def _execution_coordinator(request: Request) -> DurableFakeCoordinator:
return coordinator


def _fake_execution_coordinator(request: Request) -> DurableFakeCoordinator:
"""Keep the deterministic preview route separate from recipe execution."""

coordinator = _execution_runtime(request)
if not isinstance(coordinator, DurableFakeCoordinator):
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Execution preview is unavailable.",
)
return coordinator


def _recipe_coordinator(request: Request) -> RecipeExecutionCoordinator:
"""Expose recipes only from a ready, explicit qualification lifecycle."""

lifecycle = getattr(request.app.state, "execution_lifecycle", None)
if (
not request.app.state.preview
or lifecycle is None
or getattr(lifecycle, "profile", None) != "qualification"
or not getattr(getattr(lifecycle, "snapshot", None), "available", False)
):
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Recipe execution is unavailable.",
)
coordinator = getattr(lifecycle, "coordinator", None)
if not isinstance(coordinator, RecipeExecutionCoordinator):
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Recipe execution is unavailable.",
)
return coordinator


def _raise_recipe_request_error(exc: RecipeExecutionError) -> None:
"""Map internal recipe categories to stable, non-sensitive HTTP responses."""

if exc.code == "request_conflict":
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail="Recipe request conflicts with an existing request.",
) from exc
if exc.code == "input_artifact_unavailable":
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Source artifact is unavailable.",
) from exc
raise HTTPException(
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
detail="Recipe request could not be accepted safely.",
) from exc


def _execution_repository(request: Request):
return _execution_coordinator(request).repository
return _execution_runtime(request).repository


def _execution_latest_event(
Expand Down
68 changes: 67 additions & 1 deletion backend/cortex_backend/api/schemas.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,23 @@

from __future__ import annotations

from collections.abc import Mapping
from datetime import datetime
from typing import Any, Literal

from pydantic import BaseModel, ConfigDict, Field
from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator

from cortex_backend.core.generation import ConnectionResult
from cortex_backend.core.settings import CortexSettings
from cortex_backend.execution.recipe_coordinator import (
DEFAULT_RECIPE_RETENTION_SECONDS,
MAX_RECIPE_RETENTION_SECONDS,
)
from cortex_backend.execution.recipes import (
ImageTransformPlan,
RecipeValidationError,
parse_image_transform,
)


class APIModel(BaseModel):
Expand Down Expand Up @@ -196,6 +206,54 @@ class ExecutionPreviewRequest(APIModel):
step_delay_seconds: float = Field(default=0.0, ge=0.0, le=1.0)


class RecipeImageTransformRequest(APIModel):
"""Explicit request for one qualified, fixed-function image transform.

The source artifact must already have been copied into the owner-scoped
artifact store by a trusted attachment boundary. The API accepts only the
opaque artifact ID and the typed recipe plan; it never accepts a path or
executable instruction.
"""

request_id: str = Field(
min_length=1,
max_length=128,
pattern=r"^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$",
strict=True,
)
source_artifact_id: str = Field(
min_length=1,
max_length=128,
pattern=r"^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$",
strict=True,
)
plan: ImageTransformPlan
retention_seconds: int = Field(
default=DEFAULT_RECIPE_RETENTION_SECONDS,
ge=1,
le=MAX_RECIPE_RETENTION_SECONDS,
strict=True,
)

@field_validator("plan", mode="before")
@classmethod
def _parse_json_plan(cls, value: object) -> ImageTransformPlan:
if isinstance(value, ImageTransformPlan):
return value
if not isinstance(value, Mapping):
raise ValueError("typed image plan is invalid")
try:
return parse_image_transform(value)
except RecipeValidationError:
raise ValueError("typed image plan is invalid") from None

@model_validator(mode="after")
def _bind_plan_to_source(self) -> "RecipeImageTransformRequest":
if self.plan.input_artifact_id != self.source_artifact_id:
raise ValueError("source artifact and plan input must match")
return self


class ExecutionAccepted(APIModel):
job_id: str
request_id: str
Expand All @@ -204,6 +262,14 @@ class ExecutionAccepted(APIModel):
sequence: int


class RecipeImageTransformAccepted(APIModel):
job_id: str
request_id: str
profile: Literal["recipe.image.v1"]
status: ExecutionStatus
sequence: int


class ExecutionApprovalDecisionRequest(APIModel):
decision: Literal["approved", "denied"]

Expand Down
2 changes: 2 additions & 0 deletions backend/cortex_backend/execution/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
QualificationLifecycleConfig,
QualificationProfileError,
build_execution_lifecycle,
build_recipe_coordinator_factory,
parse_execution_profile,
)
from .manifest import (
Expand Down Expand Up @@ -218,6 +219,7 @@
"ExecutionLifecycle",
"ExecutionProfile",
"CoordinatorFactory",
"build_recipe_coordinator_factory",
"CalculatorPlan",
"CheckPlan",
"FakeExecutionPlan",
Expand Down
50 changes: 50 additions & 0 deletions backend/cortex_backend/execution/qualification.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,12 @@
from collections.abc import Callable
from typing import Literal

from .artifact_boundary import ArtifactBoundary
from .lifecycle import ExecutionLifecycle, LifecycleCoordinator, RuntimeHealth
from .recipe_coordinator import (
RecipeExecutionCoordinator,
RecipeWorkerAttemptFactory,
)
from .release_gate import RecipeRuntimeReleaseGate
from .repository import ExecutionRepository

Expand All @@ -47,6 +52,50 @@ def __init__(self, code: str) -> None:
CoordinatorFactory = Callable[[ExecutionRepository], LifecycleCoordinator]


def build_recipe_coordinator_factory(
worker_attempt_factory: RecipeWorkerAttemptFactory,
*,
artifact_boundary_factory: Callable[[ExecutionRepository], ArtifactBoundary] | None = None,
lease_seconds: float = 30.0,
supervisor_lease_seconds: float = 30.0,
) -> CoordinatorFactory:
"""Bind a qualified worker-attempt seam to the durable recipe coordinator.

The caller must provide the already-qualified attempt factory. This helper
does not launch a process, bind a broker, load a provider, or fall back to
host execution. It only creates a coordinator for the repository supplied
by :class:`ExecutionLifecycle`, keeping lifecycle ownership explicit.
"""

if not callable(worker_attempt_factory):
raise TypeError("worker_attempt_factory must be callable")
if artifact_boundary_factory is not None and not callable(artifact_boundary_factory):
raise TypeError("artifact_boundary_factory must be callable")
if lease_seconds <= 0 or supervisor_lease_seconds <= 0:
raise ValueError("lease durations must be positive")

def factory(repository: ExecutionRepository) -> LifecycleCoordinator:
artifact_boundary = (
artifact_boundary_factory(repository)
if artifact_boundary_factory is not None
else None
)
if artifact_boundary is not None and not isinstance(artifact_boundary, ArtifactBoundary):
raise TypeError("artifact_boundary_factory returned an invalid boundary")
if artifact_boundary is not None and artifact_boundary.repository is not repository:
raise ValueError("artifact boundary repository mismatch")
return RecipeExecutionCoordinator(
repository,
worker_attempt_factory,
artifact_boundary=artifact_boundary,
lease_seconds=lease_seconds,
supervisor_lease_seconds=supervisor_lease_seconds,
auto_recover=False,
)

return factory


@dataclass(frozen=True, slots=True)
class QualificationLifecycleConfig:
"""Controls required before the local qualification lifecycle can start."""
Expand Down Expand Up @@ -180,5 +229,6 @@ def build_execution_lifecycle(
"QualificationLifecycleConfig",
"QualificationProfileError",
"build_execution_lifecycle",
"build_recipe_coordinator_factory",
"parse_execution_profile",
]
Loading
Loading