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
1 change: 1 addition & 0 deletions rock/actions/sandbox/response.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ class SandboxStatusResponse(BaseModel):
host_ip: str | None = None
is_alive: bool = True
image: str | None = None
metadata: dict[str, str] | None = None
gateway_version: str | None = None
swe_rex_version: str | None = None
user_id: str | None = None
Expand Down
8 changes: 6 additions & 2 deletions rock/admin/core/sandbox_table.py
Original file line number Diff line number Diff line change
Expand Up @@ -254,7 +254,11 @@ def _merge_status_blob(raw: dict[str, Any]) -> dict[str, Any]:

The ``status`` column stores the full SandboxInfo snapshot. Fields that
only exist in the blob (e.g. ``state_history``) become top-level keys.
Scalar columns take priority over blob values.
Scalar columns take priority over blob values. The database ``labels``
column is exposed through the business-facing ``metadata`` field.
"""
status_blob = raw.pop("status", None) or {}
return {**status_blob, **raw}
result = {**status_blob, **raw}
if "labels" in result:
result["metadata"] = result["labels"]
return result
41 changes: 14 additions & 27 deletions rock/admin/entrypoints/e2b_api.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,14 @@
import math
from typing import Any
from typing import Annotated

from fastapi import APIRouter, Request
from fastapi import APIRouter, Depends, Request
from fastapi.exceptions import RequestValidationError
from fastapi.responses import JSONResponse
from fastapi.routing import APIRoute
from pydantic import BaseModel, ConfigDict, Field

from rock.common.validation import NonBlankStr
from rock.admin.proto.request import E2BCreateSandboxRequest, StartHeaders
from rock.admin.proto.response import E2BCreateSandboxResponse
from rock.common.constants import AP_SANDBOX_ID_METADATA_KEY, E2B_CLIENT_ID, E2B_ENVD_VERSION
from rock.deployments.config import DockerDeploymentConfig
from rock.logger import init_logger
from rock.sandbox.sandbox_manager import SandboxManager
Expand All @@ -16,26 +17,6 @@
logger = init_logger(__name__)


class E2BCreateSandboxRequest(BaseModel):
model_config = ConfigDict(extra="forbid", populate_by_name=True)

template_id: NonBlankStr = Field(alias="templateID")
timeout: int = Field(gt=0, strict=True)
metadata: dict[str, str]
secure: bool | None = None
allow_internet_access: bool | None = None
env_vars: dict[str, str] = Field(default_factory=dict, alias="envVars")
auto_pause: bool | None = Field(default=None, alias="autoPause")
auto_resume: dict[str, Any] | None = Field(default=None, alias="autoResume")


class E2BCreateSandboxResponse(BaseModel):
sandbox_id: str = Field(alias="sandboxID")
envd_version: str = Field(alias="envdVersion")
client_id: str = Field(alias="clientID")
template_id: str = Field(alias="templateID")


class E2BAPIRoute(APIRoute):
def get_route_handler(self):
route_handler = super().get_route_handler()
Expand Down Expand Up @@ -79,19 +60,25 @@ def _error_response(status_code: int, message: str) -> JSONResponse:
)
async def create_sandbox(
request: E2BCreateSandboxRequest,
headers: Annotated[StartHeaders, Depends()],
) -> E2BCreateSandboxResponse:
# ROCK stores lifecycle TTLs in whole minutes. Round up so an E2B timeout
# never expires a sandbox earlier than the caller requested.
config = DockerDeploymentConfig(
image=request.template_id,
auto_clear_time_minutes=math.ceil(request.timeout / 60),
container_name=request.metadata.get(AP_SANDBOX_ID_METADATA_KEY),
metadata=request.metadata,
env_vars=request.env_vars,
)
result = await e2b_sandbox_manager.start(config)
result = await e2b_sandbox_manager.start(
config,
user_info=headers.user_info,
cluster_info=headers.cluster_info,
)
return E2BCreateSandboxResponse(
sandboxID=result.sandbox_id,
envdVersion="0.1.0",
clientID="rock",
envdVersion=E2B_ENVD_VERSION,
clientID=E2B_CLIENT_ID,
templateID=request.template_id,
)
60 changes: 60 additions & 0 deletions rock/admin/entrypoints/e2b_proxy_api.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
from fastapi import APIRouter, Request
from fastapi.exceptions import RequestValidationError
from fastapi.responses import JSONResponse
from fastapi.routing import APIRoute

from rock.admin.proto.response import E2BSandboxDetail
from rock.logger import init_logger
from rock.sandbox.service.sandbox_proxy_service import SandboxProxyService
from rock.sdk.common.exceptions import BadRequestRockError, E2BSandboxNotFoundError

logger = init_logger(__name__)


class E2BProxyAPIRoute(APIRoute):
def get_route_handler(self):
route_handler = super().get_route_handler()

async def handler(request: Request):
try:
return await route_handler(request)
except RequestValidationError as error:
message = "; ".join(
f"{'.'.join(str(part) for part in item['loc'])}: {item['msg']}" for item in error.errors()
)
return _error_response(400, message)
except E2BSandboxNotFoundError as error:
return _error_response(404, str(error))
except Exception:
logger.exception("E2B get sandbox failed")
return _error_response(500, "Internal server error")

return handler


e2b_proxy_router = APIRouter(route_class=E2BProxyAPIRoute)
e2b_proxy_service: SandboxProxyService


def set_e2b_proxy_service(service: SandboxProxyService) -> None:
global e2b_proxy_service
e2b_proxy_service = service


def _error_response(status_code: int, message: str) -> JSONResponse:
return JSONResponse(status_code=status_code, content={"code": status_code, "message": message})


@e2b_proxy_router.get(
"/sandboxes/{sandboxID}",
response_model=E2BSandboxDetail,
response_model_by_alias=True,
)
async def get_sandbox(sandboxID: str) -> E2BSandboxDetail:
try:
sandbox_status = await e2b_proxy_service.get_status(sandboxID, include_all_states=True)
except BadRequestRockError as error:
if str(error) == f"Sandbox {sandboxID} not found":
raise E2BSandboxNotFoundError(str(error)) from None
raise
return E2BSandboxDetail.from_sandbox_status(sandboxID, sandbox_status)
3 changes: 3 additions & 0 deletions rock/admin/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
from rock.admin.core.template_table import TemplateTable
from rock.admin.entrypoints.admin_ops_api import admin_ops_router, set_ops_service
from rock.admin.entrypoints.e2b_api import e2b_router, set_e2b_sandbox_manager
from rock.admin.entrypoints.e2b_proxy_api import e2b_proxy_router, set_e2b_proxy_service
from rock.admin.entrypoints.sandbox_api import sandbox_router, set_sandbox_manager
from rock.admin.entrypoints.sandbox_proxy_api import sandbox_proxy_router, set_sandbox_proxy_service
from rock.admin.entrypoints.warmup_api import set_warmup_service, warmup_router
Expand Down Expand Up @@ -262,6 +263,7 @@ async def lifespan(app: FastAPI):

else:
sandbox_manager = create_sandbox_proxy_service(rock_config=rock_config, meta_store=meta_store)
set_e2b_proxy_service(sandbox_manager)
set_sandbox_proxy_service(sandbox_manager)
proxy_service_ref = sandbox_manager

Expand Down Expand Up @@ -351,6 +353,7 @@ def _include_routers(app: FastAPI, role: str) -> None:
app.include_router(sandbox_router, prefix="/apis/envs/sandbox/v1", tags=["sandbox"])
app.include_router(admin_ops_router, prefix="/apis/envs/sandbox/v1/ops", tags=["admin-ops"])
else:
app.include_router(e2b_proxy_router, tags=["e2b"])
app.include_router(sandbox_proxy_router, prefix="/apis/envs/sandbox/v1", tags=["sandbox"])
app.include_router(warmup_router, prefix="/apis/envs/sandbox/v1", tags=["warmup"])
app.include_router(gem_router, prefix="/apis/v1/envs/gem", tags=["gem"])
Expand Down
21 changes: 19 additions & 2 deletions rock/admin/proto/request.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from typing import Annotated, Literal, TypedDict
from typing import Annotated, Any, Literal, TypedDict

from fastapi import Header
from pydantic import BaseModel, Field, field_validator
from pydantic import BaseModel, ConfigDict, Field, field_validator

from rock import env_vars
from rock.actions import (
Expand All @@ -12,9 +12,23 @@
ReadFileRequest,
WriteFileRequest,
)
from rock.common.constants import BEARER_AUTHORIZATION_PREFIX
from rock.common.validation import NonBlankStr


class E2BCreateSandboxRequest(BaseModel):
model_config = ConfigDict(extra="forbid", populate_by_name=True)

template_id: NonBlankStr = Field(alias="templateID")
timeout: int = Field(gt=0, strict=True)
metadata: dict[str, str]
secure: bool | None = None
allow_internet_access: bool | None = None
env_vars: dict[str, str] = Field(default_factory=dict, alias="envVars")
auto_pause: bool | None = Field(default=None, alias="autoPause")
auto_resume: dict[str, Any] | None = Field(default=None, alias="autoResume")


class SandboxStartRequest(BaseModel):
image: NonBlankStr
"""image"""
Expand Down Expand Up @@ -190,9 +204,12 @@ def __init__(
x_user_id: str | None = Header(default="default", alias="X-User-Id"),
x_experiment_id: str | None = Header(default="default", alias="X-Experiment-Id"),
rock_authorization: str | None = Header(default="default", alias="X-Key"),
x_api_key: str | None = Header(default=None, alias="X-API-Key"),
x_namespace: str | None = Header(default="default", alias="X-Namespace"),
x_cluster: str | None = Header(default="default", alias="X-Cluster"),
):
if x_api_key is not None:
rock_authorization = f"{BEARER_AUTHORIZATION_PREFIX}{x_api_key}"
self.user_info: UserInfo = {
"user_id": x_user_id,
"experiment_id": x_experiment_id,
Expand Down
97 changes: 96 additions & 1 deletion rock/admin/proto/response.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,25 @@
from pydantic import BaseModel, Field
import datetime
import math
from ipaddress import ip_address
from typing import Literal

from pydantic import BaseModel, ConfigDict, Field

from rock.actions import SandboxResponse
from rock.actions.sandbox.response import State, StateTransitionRecord
from rock.actions.sandbox.sandbox_info import SandboxInfo
from rock.admin.proto.request import TaskSetSpec
from rock.common.constants import E2B_CLIENT_ID, E2B_ENVD_VERSION, E2B_SANDBOX_IP_METADATA_KEY, E2B_STATE_BY_ROCK_STATE
from rock.sandbox.utils.timeout import SandboxTimeoutHelper
from rock.sdk.common.exceptions import E2BSandboxNotFoundError
from rock.utils.format import parse_size_to_bytes


class E2BCreateSandboxResponse(BaseModel):
sandbox_id: str = Field(alias="sandboxID")
envd_version: str = Field(alias="envdVersion")
client_id: str = Field(alias="clientID")
template_id: str = Field(alias="templateID")


class SandboxStartResponse(SandboxResponse):
Expand All @@ -27,6 +42,7 @@ class SandboxStatusResponse(BaseModel):
host_ip: str | None = None
is_alive: bool = True
image: str | None = None
metadata: dict[str, str] | None = None
gateway_version: str | None = None
swe_rex_version: str | None = None
user_id: str | None = None
Expand Down Expand Up @@ -62,6 +78,7 @@ def from_sandbox_info(cls, sandbox_info: "SandboxInfo") -> "SandboxStatusRespons
host_ip=sandbox_info.get("host_ip"),
host_name=sandbox_info.get("host_name"),
image=sandbox_info.get("image"),
metadata=sandbox_info.get("metadata"),
user_id=sandbox_info.get("user_id"),
experiment_id=sandbox_info.get("experiment_id"),
namespace=sandbox_info.get("namespace"),
Expand All @@ -83,6 +100,84 @@ def from_sandbox_info(cls, sandbox_info: "SandboxInfo") -> "SandboxStatusRespons
)


class E2BSandboxDetail(BaseModel):
model_config = ConfigDict(populate_by_name=True)

sandbox_id: str = Field(alias="sandboxID")
metadata: dict[str, str]
state: Literal["running", "paused"]
client_id: str = Field(alias="clientID")
template_id: str = Field(alias="templateID")
envd_version: str = Field(alias="envdVersion")
cpu_count: int = Field(alias="cpuCount")
memory_mb: int = Field(alias="memoryMB")
disk_size_mb: int = Field(alias="diskSizeMB")
started_at: str = Field(alias="startedAt")
end_at: str = Field(alias="endAt")

@staticmethod
def _state(sandbox_id: str, state: State | str | None) -> Literal["running", "paused"]:
try:
rock_state = state if isinstance(state, State) else State(state)
return E2B_STATE_BY_ROCK_STATE[rock_state.value]
except (KeyError, TypeError, ValueError):
raise E2BSandboxNotFoundError(f"Sandbox {sandbox_id} not found") from None

@staticmethod
def _iso8601_timestamp(sandbox_id: str, field: str, value: object) -> str:
if not isinstance(value, str):
raise ValueError(f"Sandbox {sandbox_id} {field} is invalid")
try:
parsed = datetime.datetime.fromisoformat(value.replace("Z", "+00:00"))
except ValueError:
raise ValueError(f"Sandbox {sandbox_id} {field} is invalid") from None
if parsed.tzinfo is None:
raise ValueError(f"Sandbox {sandbox_id} {field} must include a timezone")
return parsed.isoformat(timespec="seconds")

@classmethod
def from_sandbox_status(
cls,
sandbox_id: str,
sandbox_status: SandboxStatusResponse,
) -> "E2BSandboxDetail":
state = cls._state(sandbox_id, sandbox_status.state)
end_at = (
sandbox_status.auto_stop_time
if state == "running"
else sandbox_status.auto_delete_time or sandbox_status.archive_time
)

metadata = sandbox_status.metadata
if not isinstance(metadata, dict) or not all(
isinstance(key, str) and isinstance(value, str) for key, value in metadata.items()
):
raise ValueError(f"Sandbox {sandbox_id} metadata is invalid")

host_ip = sandbox_status.host_ip
if not isinstance(host_ip, str) or not host_ip.strip():
raise ValueError(f"Sandbox {sandbox_id} IP is missing")
ip_address(host_ip)

return cls(
sandboxID=sandbox_id,
metadata={**metadata, E2B_SANDBOX_IP_METADATA_KEY: host_ip},
state=state,
clientID=E2B_CLIENT_ID,
templateID=str(sandbox_status.image),
envdVersion=E2B_ENVD_VERSION,
cpuCount=max(1, math.ceil(float(sandbox_status.cpus))),
memoryMB=parse_size_to_bytes(str(sandbox_status.memory)) // (1024**2),
diskSizeMB=parse_size_to_bytes(str(sandbox_status.disk)) // (1024**2),
startedAt=cls._iso8601_timestamp(
sandbox_id,
"start time",
sandbox_status.start_time or sandbox_status.create_time,
),
endAt=cls._iso8601_timestamp(sandbox_id, "end time", end_at),
)


class SandboxListStatusResponse(SandboxStatusResponse):
rock_authorization_encrypted: str | None = None

Expand Down
10 changes: 10 additions & 0 deletions rock/common/constants.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from enum import Enum
from typing import Literal

KATA_RUNTIME_SWITCH = "use_kata_enabled"
SUPPORT_KATA_SWITCH = "support_kata_enabled"
Expand All @@ -11,6 +12,15 @@
PID_PREFIX = "PIDSTART"
PID_SUFFIX = "PIDEND"
SCHEDULER_LOG_NAME = "scheduler.log"
BEARER_AUTHORIZATION_PREFIX = "Bearer "
AP_SANDBOX_ID_METADATA_KEY = "ap-sandbox-id"
E2B_CLIENT_ID = "rock"
E2B_ENVD_VERSION = "0.1.0"
E2B_SANDBOX_IP_METADATA_KEY = "e2b.agents.kruise.io/sandbox-ip"
E2B_STATE_BY_ROCK_STATE: dict[str, Literal["running", "paused"]] = {
"running": "running",
"archived": "paused",
}


class DeploymentHookStep(str, Enum):
Expand Down
Loading
Loading