mirror of
https://github.com/OpenHands/OpenHands.git
synced 2026-10-07 16:58:14 +08:00
feat(app-server): capture production workspace state — initial snapshot at start + archive before delete (APP-2403) (#14953)
Co-authored-by: Simon Rosenberg <simon@openhands.dev> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Co-authored-by: openhands <openhands@all-hands.dev>
This commit is contained in:
co-authored by
Simon Rosenberg
Claude Opus 4.8
openhands
parent
25ab318cbb
commit
77bd80d3a4
@@ -33,6 +33,12 @@ __all__ = ['SandboxGroupingStrategy']
|
||||
# The typed ``AppConversationInfo.acp_server`` field is a projection of this tag.
|
||||
ACP_SERVER_TAG_KEY = 'acpserver'
|
||||
|
||||
# Conversation-tag key pinning the resolved (grouped) workspace path at creation
|
||||
# so the delete-time archive captures the right directory without re-deriving it
|
||||
# from settings that may have changed. Must satisfy the SDK ^[a-z0-9]+$ tag-key
|
||||
# rule — no underscores.
|
||||
ARCHIVE_WORKSPACE_PATH_TAG_KEY = 'archiveworkspacepath'
|
||||
|
||||
|
||||
class ConversationTrigger(Enum):
|
||||
RESOLVER = 'resolver'
|
||||
|
||||
@@ -23,6 +23,7 @@ from openhands.app_server.app_conversation.app_conversation_info_service import
|
||||
AppConversationInfoService,
|
||||
)
|
||||
from openhands.app_server.app_conversation.app_conversation_models import (
|
||||
ARCHIVE_WORKSPACE_PATH_TAG_KEY,
|
||||
AppConversation,
|
||||
AppConversationInfo,
|
||||
AppConversationPage,
|
||||
@@ -871,17 +872,52 @@ async def _finalize_sandbox_delete(
|
||||
sandbox_id: str,
|
||||
db_session: AsyncSession,
|
||||
httpx_client: httpx.AsyncClient,
|
||||
conversation_id: UUID | None = None,
|
||||
workspace_path: str | None = None,
|
||||
) -> None:
|
||||
"""Delete sandbox if no other conversations reference it, then close connections."""
|
||||
"""Archive the conversation's workspace, then delete the sandbox if unreferenced.
|
||||
|
||||
Runs detached (background task) AFTER the delete response was already returned.
|
||||
The workspace is captured FIRST — a conversation-scoped step, while the runtime
|
||||
is still up — and only if that succeeds (or archiving is not REQUIRED) is the
|
||||
sandbox torn down, and only when no other conversation still references it
|
||||
(under grouping a sibling conversation keeps it alive). When archiving is
|
||||
REQUIRED and fails, the sandbox + running runtime are kept so the runtime-api
|
||||
idle reap captures the workspace later (the durability backstop). delete_sandbox
|
||||
is sandbox-scoped (stop + delete) and knows nothing about conversations.
|
||||
"""
|
||||
try:
|
||||
conversation_count = (
|
||||
await app_conversation_info_service.count_conversations_by_sandbox_id(
|
||||
sandbox_id
|
||||
)
|
||||
archived = await sandbox_service.archive_conversation_workspace(
|
||||
sandbox_id,
|
||||
conversation_id=conversation_id.hex if conversation_id else None,
|
||||
workspace_path=workspace_path,
|
||||
)
|
||||
if conversation_count == 0:
|
||||
await sandbox_service.delete_sandbox(sandbox_id)
|
||||
if not archived:
|
||||
# REQUIRED archive failed: keep the sandbox + running runtime for the
|
||||
# runtime-api idle reap to capture (the durability backstop).
|
||||
logger.warning(
|
||||
'Workspace archive required but failed for %s; leaving the '
|
||||
'sandbox + runtime for the idle reap',
|
||||
sandbox_id,
|
||||
)
|
||||
else:
|
||||
conversation_count = (
|
||||
await app_conversation_info_service.count_conversations_by_sandbox_id(
|
||||
sandbox_id
|
||||
)
|
||||
)
|
||||
if conversation_count == 0:
|
||||
await sandbox_service.delete_sandbox(sandbox_id)
|
||||
await db_session.commit()
|
||||
except Exception:
|
||||
# Any failure in the finalizer (a transient stop/lookup error, the count
|
||||
# query, the commit itself): do NOT commit a half-done delete, so no
|
||||
# orphaned row is left; the row + running runtime stay for the runtime-api
|
||||
# idle reap to capture + reap.
|
||||
logger.exception(
|
||||
'Deferred sandbox cleanup failed for %s; kept for retry', sandbox_id
|
||||
)
|
||||
await db_session.rollback()
|
||||
finally:
|
||||
await asyncio.gather(
|
||||
db_session.aclose(),
|
||||
@@ -965,6 +1001,8 @@ async def delete_app_conversation(
|
||||
|
||||
# Delete the sandbox in the background if no other conversations reference it
|
||||
if sandbox_id:
|
||||
# Path pinned at creation; the finalizer archives exactly this directory.
|
||||
workspace_path = app_conversation_info.tags.get(ARCHIVE_WORKSPACE_PATH_TAG_KEY)
|
||||
asyncio.create_task(
|
||||
_finalize_sandbox_delete(
|
||||
sandbox_service,
|
||||
@@ -972,6 +1010,8 @@ async def delete_app_conversation(
|
||||
sandbox_id,
|
||||
db_session,
|
||||
httpx_client,
|
||||
conversation_id=conversation_uuid,
|
||||
workspace_path=workspace_path,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
@@ -122,6 +122,7 @@ class AppConversationService(ABC):
|
||||
sandbox: SandboxInfo,
|
||||
workspace: AsyncRemoteWorkspace,
|
||||
agent_server_url: str,
|
||||
conversation_id: UUID,
|
||||
) -> AsyncGenerator[AppConversationStartTask, None]:
|
||||
"""Run the setup scripts for the project and yield status updates"""
|
||||
yield task
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import shlex
|
||||
@@ -27,6 +28,7 @@ from openhands.app_server.app_conversation.skill_loader import (
|
||||
load_skills_from_agent_server,
|
||||
)
|
||||
from openhands.app_server.integrations.service_types import ProviderType
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
from openhands.app_server.sandbox.sandbox_models import SandboxInfo
|
||||
from openhands.app_server.user.user_context import UserContext
|
||||
from openhands.app_server.utils.auth import looks_like_jwt
|
||||
@@ -46,6 +48,7 @@ from openhands.sdk.skills import Skill
|
||||
from openhands.sdk.workspace.remote.async_remote_workspace import AsyncRemoteWorkspace
|
||||
|
||||
_logger = logging.getLogger(__name__)
|
||||
|
||||
PRE_COMMIT_HOOK = '.git/hooks/pre-commit'
|
||||
PRE_COMMIT_LOCAL = '.git/hooks/pre-commit.local'
|
||||
|
||||
@@ -252,6 +255,7 @@ class AppConversationServiceBase(AppConversationService, ABC):
|
||||
sandbox: SandboxInfo,
|
||||
workspace: AsyncRemoteWorkspace,
|
||||
agent_server_url: str,
|
||||
conversation_id: UUID,
|
||||
) -> AsyncGenerator[AppConversationStartTask, None]:
|
||||
task.status = AppConversationStartTaskStatus.PREPARING_REPOSITORY
|
||||
yield task
|
||||
@@ -264,6 +268,28 @@ class AppConversationServiceBase(AppConversationService, ABC):
|
||||
workspace.working_dir, task.request.selected_repository
|
||||
)
|
||||
|
||||
# Snapshot the INITIAL workspace (repo exactly as cloned, before setup.sh)
|
||||
# for downstream dataset/eval creation. Awaited HERE — before any
|
||||
# mutating step — so
|
||||
# the capture is deterministic and can't race setup.sh writing the same
|
||||
# tree. Best-effort and bounded so it never fails or unduly delays startup.
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
self._maybe_archive_initial_state(
|
||||
task,
|
||||
sandbox,
|
||||
workspace,
|
||||
agent_server_url,
|
||||
project_dir,
|
||||
conversation_id,
|
||||
),
|
||||
timeout=workspace_archive.initial_archive_deadline(),
|
||||
)
|
||||
except Exception as e:
|
||||
_logger.warning(
|
||||
'Initial workspace archive did not complete in time (ignored): %s', e
|
||||
)
|
||||
|
||||
task.status = AppConversationStartTaskStatus.RUNNING_SETUP_SCRIPT
|
||||
yield task
|
||||
await self.maybe_run_setup_script(workspace, project_dir)
|
||||
@@ -281,6 +307,55 @@ class AppConversationServiceBase(AppConversationService, ABC):
|
||||
agent_server_url,
|
||||
)
|
||||
|
||||
async def _maybe_archive_initial_state(
|
||||
self,
|
||||
task: AppConversationStartTask,
|
||||
sandbox: SandboxInfo,
|
||||
workspace: AsyncRemoteWorkspace,
|
||||
agent_server_url: str,
|
||||
project_dir: str,
|
||||
conversation_id: UUID,
|
||||
) -> None:
|
||||
"""Snapshot the initial (as-cloned) workspace state — best-effort.
|
||||
|
||||
Awaited before setup.sh, so it captures the repo exactly as cloned with no
|
||||
concurrent mutation, keyed to the real conversation id. No-op unless
|
||||
``RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED`` is set and a repository was
|
||||
actually cloned (an empty workspace is not a useful starting state).
|
||||
Every failure is swallowed so this can never break conversation startup.
|
||||
"""
|
||||
try:
|
||||
if not workspace_archive.initial_archive_enabled():
|
||||
return
|
||||
if not task.request.selected_repository:
|
||||
# Only repo-backed conversations have a meaningful initial state.
|
||||
return
|
||||
|
||||
# Record the commit the snapshot came from — tar.gz carries no
|
||||
# base-commit header. Best-effort: a miss just leaves base_commit ''.
|
||||
base_commit = ''
|
||||
try:
|
||||
result = await workspace.execute_command(
|
||||
'git rev-parse HEAD', project_dir
|
||||
)
|
||||
if not result.exit_code:
|
||||
base_commit = result.stdout.strip()
|
||||
except Exception as e:
|
||||
_logger.debug('Initial-state base commit lookup failed: %s', e)
|
||||
|
||||
conv_hex = conversation_id.hex if conversation_id else None
|
||||
await workspace_archive.archive_initial_workspace(
|
||||
workspace.client,
|
||||
agent_server_url=agent_server_url,
|
||||
session_api_key=sandbox.session_api_key,
|
||||
project_dir=project_dir,
|
||||
sandbox_id=sandbox.id,
|
||||
conversation_id=conv_hex,
|
||||
base_commit=base_commit,
|
||||
)
|
||||
except Exception as e:
|
||||
_logger.warning('Initial workspace archive step failed (ignored): %s', e)
|
||||
|
||||
async def _configure_git_user_settings(
|
||||
self,
|
||||
workspace: AsyncRemoteWorkspace,
|
||||
|
||||
@@ -28,6 +28,7 @@ from openhands.app_server.app_conversation.app_conversation_info_service import
|
||||
)
|
||||
from openhands.app_server.app_conversation.app_conversation_models import (
|
||||
ACP_SERVER_TAG_KEY,
|
||||
ARCHIVE_WORKSPACE_PATH_TAG_KEY,
|
||||
AgentType,
|
||||
AppConversation,
|
||||
AppConversationInfo,
|
||||
@@ -97,6 +98,7 @@ from openhands.app_server.sandbox.sandbox_spec_service import (
|
||||
from openhands.app_server.services.injector import InjectorState
|
||||
from openhands.app_server.services.jwt_service import JwtService
|
||||
from openhands.app_server.settings.llm_profiles import resolve_profile_llm
|
||||
from openhands.app_server.settings.settings_models import grouped_workspace_dir
|
||||
from openhands.app_server.user.user_context import UserContext
|
||||
from openhands.app_server.user.user_models import UserInfo
|
||||
from openhands.app_server.utils.docker_utils import (
|
||||
@@ -395,10 +397,12 @@ class LiveStatusAppConversationService(AppConversationServiceBase):
|
||||
conversation_id = request.conversation_id or uuid4()
|
||||
|
||||
# Setup working dir based on grouping
|
||||
working_dir = sandbox_spec.working_dir
|
||||
sandbox_grouping_strategy = await self._get_sandbox_grouping_strategy()
|
||||
if sandbox_grouping_strategy != SandboxGroupingStrategy.NO_GROUPING:
|
||||
working_dir = f'{working_dir}/{conversation_id.hex}'
|
||||
working_dir = grouped_workspace_dir(
|
||||
sandbox_spec.working_dir,
|
||||
sandbox_grouping_strategy,
|
||||
conversation_id.hex,
|
||||
)
|
||||
|
||||
# Run setup scripts
|
||||
remote_workspace = AsyncRemoteWorkspace(
|
||||
@@ -407,7 +411,7 @@ class LiveStatusAppConversationService(AppConversationServiceBase):
|
||||
working_dir=working_dir,
|
||||
)
|
||||
async for updated_task in self.run_setup_scripts(
|
||||
task, sandbox, remote_workspace, agent_server_url
|
||||
task, sandbox, remote_workspace, agent_server_url, conversation_id
|
||||
):
|
||||
yield updated_task
|
||||
|
||||
@@ -486,6 +490,10 @@ class LiveStatusAppConversationService(AppConversationServiceBase):
|
||||
# the same agent back through the AgentBase discriminator.
|
||||
request_agent = start_conversation_request.agent
|
||||
tags: dict[str, str] = {}
|
||||
# Pin where the workspace was actually created so the delete-time
|
||||
# archive captures the right directory without re-deriving the path
|
||||
# from settings (e.g. grouping) that may change before delete.
|
||||
tags[ARCHIVE_WORKSPACE_PATH_TAG_KEY] = working_dir
|
||||
if request_agent.agent_kind == 'acp':
|
||||
llm_model = request_agent.acp_model
|
||||
agent_kind = 'acp'
|
||||
|
||||
@@ -41,3 +41,22 @@ class PermissionsError(OpenHandsError):
|
||||
|
||||
class SandboxError(OpenHandsError):
|
||||
"""Error in Sandbox."""
|
||||
|
||||
|
||||
class SandboxDeleteRetryError(OpenHandsError):
|
||||
"""The sandbox exists but its delete could not complete and was kept for retry.
|
||||
|
||||
Raised by ``delete_sandbox`` when the runtime /stop or lookup hits a transient
|
||||
failure. (Archiving never raises — it returns False from
|
||||
``archive_conversation_workspace`` to signal a REQUIRED capture should block.)
|
||||
503 (vs 404) so a client distinguishes "still here, try again" from "not
|
||||
found" and keeps retrying.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
detail: Any = None,
|
||||
headers: dict[str, str] | None = None,
|
||||
status_code: int = status.HTTP_503_SERVICE_UNAVAILABLE,
|
||||
):
|
||||
super().__init__(status_code=status_code, detail=detail, headers=headers)
|
||||
|
||||
@@ -17,6 +17,18 @@ class FileStore(DiscriminatedUnionMixin, ABC):
|
||||
def write(self, path: str, contents: str | bytes) -> None:
|
||||
pass
|
||||
|
||||
def write_from_path(self, path: str, source_path: str) -> None:
|
||||
"""Store the object at ``path`` from a local file at ``source_path``.
|
||||
|
||||
The default reads the whole file into memory and delegates to ``write``;
|
||||
backends that can stream from disk override this to bound peak memory.
|
||||
Callers uploading large blobs (e.g. multi-GB workspace archives, where
|
||||
buffering one whole copy in RAM under concurrent deletes risks
|
||||
OOM-killing the pod) should prefer this over ``write(path, f.read())``.
|
||||
"""
|
||||
with open(source_path, 'rb') as f:
|
||||
self.write(path, f.read())
|
||||
|
||||
@abstractmethod
|
||||
def read(self, path: str) -> str:
|
||||
pass
|
||||
|
||||
@@ -50,6 +50,11 @@ class GoogleCloudFileStore(FileStore):
|
||||
with blob.open(mode) as f:
|
||||
f.write(contents)
|
||||
|
||||
def write_from_path(self, path: str, source_path: str) -> None:
|
||||
# Streams the file to GCS in chunks; never buffers the whole blob in RAM.
|
||||
blob: Blob = self.bucket.blob(path)
|
||||
blob.upload_from_filename(source_path)
|
||||
|
||||
def read(self, path: str) -> str:
|
||||
blob: Blob = self.bucket.blob(path)
|
||||
try:
|
||||
|
||||
@@ -42,6 +42,20 @@ class LocalFileStore(FileStore):
|
||||
os.remove(temp_path)
|
||||
raise
|
||||
|
||||
def write_from_path(self, path: str, source_path: str) -> None:
|
||||
# shutil.copyfile streams in chunks (never the whole file in RAM); keep
|
||||
# the same write-temp-then-atomic-rename to avoid torn concurrent writes.
|
||||
full_path = self.get_full_path(path)
|
||||
os.makedirs(os.path.dirname(full_path), exist_ok=True)
|
||||
temp_path = f'{full_path}.tmp.{os.getpid()}.{threading.get_ident()}'
|
||||
try:
|
||||
shutil.copyfile(source_path, temp_path)
|
||||
os.replace(temp_path, full_path)
|
||||
except Exception:
|
||||
if os.path.exists(temp_path):
|
||||
os.remove(temp_path)
|
||||
raise
|
||||
|
||||
def read(self, path: str) -> str:
|
||||
full_path = self.get_full_path(path)
|
||||
with open(full_path, 'r') as f:
|
||||
|
||||
@@ -7,6 +7,10 @@ from openhands.app_server.utils.logger import openhands_logger as logger
|
||||
|
||||
|
||||
class InMemoryFileStore(FileStore):
|
||||
# Text-only by design: this store is part of the env-parsed FileStore config
|
||||
# union (so the value type must stay a primitive) and read() returns str, so
|
||||
# it cannot round-trip a binary archive. Not a valid RUNTIME_FILE_ARCHIVE
|
||||
# store — see workspace_archive._archive_store_type.
|
||||
files: dict[str, str] = Field(default_factory=dict)
|
||||
|
||||
def write(self, path: str, contents: str | bytes) -> None:
|
||||
|
||||
@@ -75,6 +75,24 @@ class S3FileStore(FileStore):
|
||||
f"Error: Failed to write to bucket '{self._get_bucket_name()}' at path {path}: {e}"
|
||||
)
|
||||
|
||||
def write_from_path(self, path: str, source_path: str) -> None:
|
||||
# upload_file streams the file from disk in parts; never buffers the whole
|
||||
# object in RAM (unlike put_object with Body=f.read()).
|
||||
try:
|
||||
self.client.upload_file(source_path, self._get_bucket_name(), path)
|
||||
except botocore.exceptions.ClientError as e:
|
||||
if e.response['Error']['Code'] == 'AccessDenied':
|
||||
raise FileNotFoundError(
|
||||
f"Error: Access denied to bucket '{self._get_bucket_name()}'."
|
||||
)
|
||||
elif e.response['Error']['Code'] == 'NoSuchBucket':
|
||||
raise FileNotFoundError(
|
||||
f"Error: The bucket '{self._get_bucket_name()}' does not exist."
|
||||
)
|
||||
raise FileNotFoundError(
|
||||
f"Error: Failed to write to bucket '{self._get_bucket_name()}' at path {path}: {e}"
|
||||
)
|
||||
|
||||
def read(self, path: str) -> str:
|
||||
try:
|
||||
response: GetObjectOutputDict = self.client.get_object(
|
||||
|
||||
@@ -386,7 +386,7 @@ class ProcessSandboxService(SandboxService):
|
||||
return False
|
||||
|
||||
async def delete_sandbox(self, sandbox_id: str) -> bool:
|
||||
"""Delete a sandbox."""
|
||||
"""Delete a sandbox. (No workspace archiving for local processes.)"""
|
||||
process_info = _processes.get(sandbox_id)
|
||||
if process_info is None:
|
||||
return False
|
||||
|
||||
@@ -24,7 +24,8 @@ from openhands.agent_server.utils import utc_now
|
||||
from openhands.app_server.app_conversation.app_conversation_models import (
|
||||
AppConversationInfo,
|
||||
)
|
||||
from openhands.app_server.errors import SandboxError
|
||||
from openhands.app_server.errors import SandboxDeleteRetryError, SandboxError
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
from openhands.app_server.sandbox.sandbox_models import (
|
||||
AGENT_SERVER,
|
||||
VSCODE,
|
||||
@@ -48,8 +49,12 @@ from openhands.app_server.sandbox.sandbox_spec_service import (
|
||||
resolve_sandbox_spec,
|
||||
)
|
||||
from openhands.app_server.services.injector import InjectorState
|
||||
from openhands.app_server.settings.settings_models import grouped_workspace_dir
|
||||
from openhands.app_server.user.specifiy_user_context import ADMIN, USER_CONTEXT_ATTR
|
||||
from openhands.app_server.user.user_context import UserContext
|
||||
from openhands.app_server.utils.docker_utils import (
|
||||
replace_localhost_hostname_for_docker,
|
||||
)
|
||||
from openhands.app_server.utils.sql_utils import Base, UtcDateTime
|
||||
from openhands.sdk.utils.paging import page_iterator
|
||||
|
||||
@@ -571,17 +576,49 @@ class RemoteSandboxService(SandboxService):
|
||||
async def delete_sandbox(self, sandbox_id: str) -> bool:
|
||||
"""Delete a sandbox by stopping its runtime.
|
||||
|
||||
Security: Deleting the stored_sandbox record also removes the
|
||||
session_api_key_hash, invalidating any leaked session keys.
|
||||
Purely sandbox-scoped: stop the runtime and delete the record. Workspace
|
||||
capture is a separate conversation-scoped step
|
||||
(``archive_conversation_workspace``) the conversation-delete finalizer runs
|
||||
BEFORE tearing the sandbox down — so a long archive never blocks this call
|
||||
(and the direct sandbox DELETE route can't 504 on it).
|
||||
|
||||
If the runtime is already gone (paused/reaped/double-delete, a 404 from
|
||||
the runtime API), the record is deleted directly to avoid orphaning it.
|
||||
|
||||
Returns False ONLY when the sandbox does not exist (router -> 404). A
|
||||
transient runtime /stop / lookup failure raises ``SandboxDeleteRetryError``
|
||||
(router -> 503) and keeps the row + runtime for a retry — so a live sandbox
|
||||
is never reported as 404.
|
||||
|
||||
Security: the session_api_key_hash is invalidated UP FRONT (like
|
||||
``pause_sandbox`` clears it before pausing) so a delete — commonly a
|
||||
revoke of a leaked key — kills it promptly. This goes further than pause:
|
||||
on a transient stop failure the invalidation is committed before raising,
|
||||
so the caller's rollback cannot resurrect the just-revoked key (pause does
|
||||
not commit, so its clear can still be rolled back). The row is kept for
|
||||
retry.
|
||||
"""
|
||||
had_key = False
|
||||
try:
|
||||
stored_sandbox = await self._get_stored_sandbox(sandbox_id)
|
||||
if not stored_sandbox:
|
||||
return False
|
||||
# Deleting the record also removes the session_api_key_hash,
|
||||
# which invalidates any leaked session keys.
|
||||
await self.db_session.delete(stored_sandbox)
|
||||
runtime_data = await self._get_runtime(sandbox_id)
|
||||
# Security: drop the key now, before the (fallible) runtime stop.
|
||||
had_key = stored_sandbox.session_api_key_hash is not None
|
||||
stored_sandbox.session_api_key_hash = None
|
||||
try:
|
||||
runtime_data = await self._get_runtime(sandbox_id)
|
||||
except httpx.HTTPStatusError as e:
|
||||
if e.response.status_code != 404:
|
||||
raise
|
||||
# Runtime already gone: nothing to stop. Delete the orphaned row.
|
||||
_logger.info(
|
||||
f'Runtime for sandbox {sandbox_id} already gone (404); '
|
||||
'deleting record'
|
||||
)
|
||||
await self.db_session.delete(stored_sandbox)
|
||||
return True
|
||||
|
||||
response = await self._send_runtime_api_request(
|
||||
'POST',
|
||||
'/stop',
|
||||
@@ -589,10 +626,153 @@ class RemoteSandboxService(SandboxService):
|
||||
)
|
||||
if response.status_code != 404:
|
||||
response.raise_for_status()
|
||||
await self.db_session.delete(stored_sandbox)
|
||||
return True
|
||||
except httpx.HTTPError as e:
|
||||
# Transient runtime lookup/stop failure: keep the row + runtime and
|
||||
# signal retryable (503) — never a 404. Persist the key invalidation
|
||||
# now: the caller rolls back on this raise, which would otherwise
|
||||
# restore the hash and leave a just-revoked key valid.
|
||||
_logger.error(f'Error deleting sandbox {sandbox_id}: {e}')
|
||||
return False
|
||||
if had_key:
|
||||
await self.db_session.commit()
|
||||
raise SandboxDeleteRetryError(
|
||||
f'Could not complete delete for sandbox {sandbox_id}: {e}'
|
||||
) from e
|
||||
|
||||
async def _resolve_archive_path(
|
||||
self,
|
||||
stored_sandbox: StoredRemoteSandbox,
|
||||
conversation_id: str | None,
|
||||
workspace_path: str | None,
|
||||
) -> str:
|
||||
"""Path to archive: the value pinned at conversation creation if present,
|
||||
else rebuilt from the SAME base the clone used (the sandbox spec's
|
||||
``working_dir``) plus the grouping nesting.
|
||||
|
||||
Pre-pinning conversations have no pinned path; the legacy fallback re-reads
|
||||
the live grouping strategy, which can disagree with creation if the user
|
||||
toggled it — but a resulting 404 no longer silently tears the sandbox down
|
||||
under REQUIRED (it blocks for the idle reap). Raises if the layout cannot
|
||||
be resolved, so the caller never archives to the wrong path.
|
||||
"""
|
||||
if workspace_path:
|
||||
return workspace_path
|
||||
# For cloud conversations the sandbox id is the conversation_id.hex.
|
||||
conversation_key = conversation_id or stored_sandbox.id
|
||||
sandbox_spec = await self.sandbox_spec_service.get_sandbox_spec(
|
||||
stored_sandbox.sandbox_spec_id
|
||||
)
|
||||
if sandbox_spec is None:
|
||||
raise SandboxError(
|
||||
f'No sandbox spec {stored_sandbox.sandbox_spec_id} for archive'
|
||||
)
|
||||
grouping = (await self.user_context.get_user_info()).sandbox_grouping_strategy
|
||||
return grouped_workspace_dir(
|
||||
sandbox_spec.working_dir, grouping, conversation_key
|
||||
)
|
||||
|
||||
async def _archive_workspace(
|
||||
self,
|
||||
stored_sandbox: StoredRemoteSandbox,
|
||||
conversation_id: str | None,
|
||||
runtime_data: dict,
|
||||
workspace_path: str | None,
|
||||
) -> bool:
|
||||
"""Archive one workspace via the in-pod agent-server; return may-proceed.
|
||||
|
||||
Returns True when the workspace was captured, when there was nothing to
|
||||
capture, or when archiving failed but is not REQUIRED. Returns False only
|
||||
when archiving is REQUIRED and could not confirm a capture (the caller
|
||||
decides whether to block + retry). Never raises.
|
||||
"""
|
||||
try:
|
||||
archive_path = await self._resolve_archive_path(
|
||||
stored_sandbox, conversation_id, workspace_path
|
||||
)
|
||||
# The runtime url is raw (localhost in Docker/local); transform it the
|
||||
# same way every other agent-server URL resolution does.
|
||||
runtime = dict(runtime_data)
|
||||
url = runtime.get('url')
|
||||
if url:
|
||||
runtime['url'] = replace_localhost_hostname_for_docker(url)
|
||||
return await workspace_archive.archive_workspace(
|
||||
self.httpx_client,
|
||||
runtime,
|
||||
stored_sandbox.id,
|
||||
archive_path=archive_path,
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
except Exception:
|
||||
# Could not resolve the workspace layout: never archive to the wrong
|
||||
# path. Honor REQUIRED (block + retry) vs best-effort (proceed).
|
||||
_logger.exception(
|
||||
'Could not resolve archive path for %s', stored_sandbox.id
|
||||
)
|
||||
return not workspace_archive.archive_required()
|
||||
|
||||
async def archive_conversation_workspace(
|
||||
self,
|
||||
sandbox_id: str,
|
||||
conversation_id: str | None = None,
|
||||
workspace_path: str | None = None,
|
||||
) -> bool:
|
||||
"""Archive ONE conversation's workspace; return whether delete may proceed.
|
||||
|
||||
The sole app-server capture path: the conversation-delete finalizer calls
|
||||
this for every conversation delete (while the runtime is still up), then
|
||||
tears the sandbox down only when this was its last conversation. Keying to
|
||||
the conversation lets a grouped sandbox capture the right per-conversation
|
||||
repo, and means no grouped conversation's work is lost when a sibling later
|
||||
triggers the sandbox delete.
|
||||
|
||||
``workspace_path`` is the path pinned at conversation creation; when given
|
||||
the capture uses it verbatim instead of re-deriving the layout.
|
||||
|
||||
Returns True when the workspace was captured, when there was nothing to
|
||||
capture (runtime already gone, or no repo at the path), or when archiving
|
||||
failed but is not REQUIRED. Returns False only when archiving is REQUIRED
|
||||
and could not confirm a capture, so the finalizer keeps the sandbox +
|
||||
running runtime for the runtime-api idle reap (the durability backstop).
|
||||
Never raises. No-op (returns True) unless archiving is enabled.
|
||||
"""
|
||||
if not workspace_archive.archive_enabled():
|
||||
return True
|
||||
try:
|
||||
stored_sandbox = await self._get_stored_sandbox(sandbox_id)
|
||||
if not stored_sandbox:
|
||||
return True
|
||||
runtime_data = await self._get_runtime(sandbox_id)
|
||||
except httpx.HTTPStatusError as e:
|
||||
if e.response.status_code == 404:
|
||||
# Runtime already gone: nothing to capture for this conversation.
|
||||
return True
|
||||
# Couldn't reach the runtime: honor REQUIRED (block + keep) vs
|
||||
# best-effort (let the delete proceed; delete_sandbox re-checks).
|
||||
_logger.exception(
|
||||
'Workspace archive lookup failed for %s (%s)',
|
||||
sandbox_id,
|
||||
conversation_id,
|
||||
)
|
||||
return not workspace_archive.archive_required()
|
||||
except Exception:
|
||||
_logger.exception(
|
||||
'Workspace archive lookup failed for %s (%s)',
|
||||
sandbox_id,
|
||||
conversation_id,
|
||||
)
|
||||
return not workspace_archive.archive_required()
|
||||
archived = await self._archive_workspace(
|
||||
stored_sandbox, conversation_id, runtime_data, workspace_path
|
||||
)
|
||||
if not archived:
|
||||
_logger.warning(
|
||||
'Workspace archive required but failed for %s (%s); keeping the '
|
||||
'sandbox for the idle reap to capture',
|
||||
sandbox_id,
|
||||
conversation_id,
|
||||
)
|
||||
return archived
|
||||
|
||||
async def pause_old_sandboxes(self, max_num_sandboxes: int) -> list[str]:
|
||||
"""Pause the oldest running sandboxes until at most max_num_sandboxes remain.
|
||||
|
||||
@@ -110,6 +110,9 @@ async def delete_sandbox(
|
||||
sandbox_id: str,
|
||||
sandbox_service: SandboxService = sandbox_service_dependency,
|
||||
) -> Success:
|
||||
# delete_sandbox is sandbox-scoped (stop + delete) and never archives, so this
|
||||
# request handler can't block on a minutes-long capture. Workspace capture is
|
||||
# owned by the conversation-delete finalizer and the runtime-api idle reaper.
|
||||
exists = await sandbox_service.delete_sandbox(sandbox_id)
|
||||
if not exists:
|
||||
raise HTTPException(status.HTTP_404_NOT_FOUND)
|
||||
|
||||
@@ -197,9 +197,29 @@ class SandboxService(ABC):
|
||||
async def delete_sandbox(self, sandbox_id: str) -> bool:
|
||||
"""Begin the process of deleting a sandbox (which may involve stopping it).
|
||||
|
||||
Return False if the sandbox did not exist.
|
||||
Return False if the sandbox did not exist. Purely sandbox-scoped (stop the
|
||||
runtime, delete the record); workspace capture is a separate
|
||||
conversation-scoped step (``archive_conversation_workspace``) the
|
||||
conversation-delete finalizer runs before the sandbox is torn down.
|
||||
"""
|
||||
|
||||
async def archive_conversation_workspace(
|
||||
self,
|
||||
sandbox_id: str,
|
||||
conversation_id: str | None = None,
|
||||
workspace_path: str | None = None,
|
||||
) -> bool:
|
||||
"""Archive one conversation's workspace; return whether delete may proceed.
|
||||
|
||||
Default no-op (returns True) for backends that do not archive; overridden
|
||||
by RemoteSandboxService. The conversation-delete finalizer calls this
|
||||
before ``delete_sandbox`` so the workspace is captured while the runtime is
|
||||
still up. ``workspace_path`` is the path pinned at creation. Returns False
|
||||
only when archiving is REQUIRED and failed, so the caller leaves the
|
||||
sandbox up for a later (idle-reap) capture.
|
||||
"""
|
||||
return True
|
||||
|
||||
async def pause_old_sandboxes(self, max_num_sandboxes: int) -> list[str]:
|
||||
"""Pause the oldest sandboxes if there are more than max_num_sandboxes running.
|
||||
In a multi user environment, this will pause sandboxes only for the current user.
|
||||
|
||||
@@ -0,0 +1,524 @@
|
||||
"""Archive a remote sandbox's workspace to object storage before deletion.
|
||||
|
||||
Pulls a workspace archive from the in-pod agent-server endpoint
|
||||
(``GET /api/file/archive``) and stores it, plus a small manifest, in object
|
||||
storage so the agent's work — production OpenHands Cloud workspace state —
|
||||
survives sandbox deletion and is preserved for downstream use (e.g. dataset/eval
|
||||
creation).
|
||||
|
||||
It covers the *explicit-delete-while-running* path. The dominant idle/expiry
|
||||
reap is handled separately in runtime-api at pause time, because that deletion
|
||||
never reaches the app-server.
|
||||
|
||||
Configuration is environment-driven and the feature is a no-op unless
|
||||
``RUNTIME_FILE_ARCHIVE_ENABLED`` is set.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import tempfile
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from openhands.agent_server.utils import utc_now
|
||||
from openhands.app_server.file_store import get_file_store
|
||||
from openhands.app_server.file_store.files import FileStore
|
||||
|
||||
_logger = logging.getLogger(__name__)
|
||||
|
||||
# Formats the SDK GET /api/file/archive producer accepts (git-delta | tar.gz);
|
||||
# anything else 422s, so validate before issuing the request.
|
||||
_ARCHIVE_SUFFIX = {'git-delta': 'patch', 'tar.gz': 'tar.gz'}
|
||||
|
||||
|
||||
def _archive_request_params(path: str, fmt: str) -> dict[str, str]:
|
||||
"""Query params for GET /api/file/archive.
|
||||
|
||||
The tar.gz is the SELF-CONTAINED full capture, so disable the endpoint's
|
||||
default excludes for it: otherwise agent output under dist/build/node_modules
|
||||
and the repo's .git history are dropped and it is no more complete than the
|
||||
git-delta (defeating the whole point of capturing 'both'). Credential-bearing
|
||||
git internals are still scrubbed server-side even with excludes off. git-delta
|
||||
keeps the defaults — it is the compact companion, not the full capture.
|
||||
"""
|
||||
params = {'path': path, 'format': fmt}
|
||||
if fmt == 'tar.gz':
|
||||
params['use_default_excludes'] = 'false'
|
||||
return params
|
||||
|
||||
|
||||
def archive_enabled() -> bool:
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_ENABLED', 'false').lower() in ('true', '1')
|
||||
|
||||
|
||||
def archive_required() -> bool:
|
||||
"""When true, an archive failure blocks deletion so it can be retried."""
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_REQUIRED', 'false').lower() in ('true', '1')
|
||||
|
||||
|
||||
def _archive_bucket() -> str:
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_BUCKET', '')
|
||||
|
||||
|
||||
def _archive_prefix() -> str:
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_PREFIX', 'workspace-archives')
|
||||
|
||||
|
||||
def _archive_format() -> str:
|
||||
# Default to 'both' — the compact git-delta AND a self-contained full tar.gz.
|
||||
# git-delta alone is lossy as a sole capture: it respects the repo's
|
||||
# .gitignore (so agent-authored gitignored files are dropped) and needs the
|
||||
# base tree to reconstruct, whereas the tar.gz is self-contained and captures
|
||||
# those files. Keep both until the storage cost is measured, then narrow to
|
||||
# 'git-delta' (+ bucket lifecycle) if warranted (infra#1444).
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_FORMAT', 'both')
|
||||
|
||||
|
||||
def _formats_to_capture() -> list[str] | None:
|
||||
"""Resolve RUNTIME_FILE_ARCHIVE_FORMAT to the list of formats to upload.
|
||||
|
||||
'both' captures the git-delta AND the full tar.gz; a single format captures
|
||||
just that one. Returns None for an unsupported value (a hard config error the
|
||||
SDK producer would 422), so the caller can log + skip instead of mis-reading
|
||||
it as "nothing to archive".
|
||||
"""
|
||||
fmt = _archive_format()
|
||||
if fmt == 'both':
|
||||
return ['git-delta', 'tar.gz']
|
||||
if fmt in _ARCHIVE_SUFFIX:
|
||||
return [fmt]
|
||||
return None
|
||||
|
||||
|
||||
def _float_env(name: str, default: float) -> float:
|
||||
"""Parse a float env var, falling back to default on a non-numeric value.
|
||||
|
||||
A bad override (``'120s'``, a stray newline) must not raise on every archive
|
||||
call — that would wedge every REQUIRED delete forever.
|
||||
"""
|
||||
raw = os.getenv(name)
|
||||
if not raw:
|
||||
return default
|
||||
try:
|
||||
value = float(raw)
|
||||
except ValueError:
|
||||
_logger.warning('Invalid %s=%r; using %s', name, raw, default)
|
||||
return default
|
||||
if value <= 0:
|
||||
# A non-positive timeout/deadline would make httpx raise on every archive
|
||||
# (the wedge this guard exists to prevent); fall back to the safe default.
|
||||
_logger.warning('Non-positive %s=%r; using %s', name, raw, default)
|
||||
return default
|
||||
return value
|
||||
|
||||
|
||||
def _archive_timeout() -> float:
|
||||
# Must cover the agent-server git build budget (read-tree 60 + add 300 + diff
|
||||
# 300 = up to ~660s) before the first response byte flows, or large repos
|
||||
# ReadTimeout and never capture. The final archive runs in the detached
|
||||
# delete finalizer, so a long wait here doesn't block any user request; the
|
||||
# initial snapshot is separately bounded by initial_archive_deadline().
|
||||
return _float_env('RUNTIME_FILE_ARCHIVE_TIMEOUT', 660.0)
|
||||
|
||||
|
||||
def _archive_store_type() -> str:
|
||||
# Default to GCS to preserve current behavior; local/s3 also work. NOT
|
||||
# 'memory' — it is text-only (read() returns str) and would corrupt the
|
||||
# binary archive (see InMemoryFileStore).
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_STORE_TYPE', 'google_cloud')
|
||||
|
||||
|
||||
def _get_archive_file_store() -> FileStore:
|
||||
"""Object store for archives, built via the backend-portable factory."""
|
||||
return get_file_store(_archive_store_type(), _archive_bucket())
|
||||
|
||||
|
||||
def _cleanup_tempfile(path: str | None) -> None:
|
||||
if not path:
|
||||
return
|
||||
try:
|
||||
os.unlink(path)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
async def _stream_to_tempfile(response: Any) -> tuple[str, int]:
|
||||
"""Stream a 200 response body to a temp file; return (path, byte_count).
|
||||
|
||||
Avoids buffering the whole archive in app-server RAM (OOM risk under
|
||||
concurrent large deletes). Cleans up its own file if streaming fails.
|
||||
"""
|
||||
tmp = tempfile.NamedTemporaryFile(delete=False)
|
||||
byte_count = 0
|
||||
try:
|
||||
async for chunk in response.aiter_bytes():
|
||||
tmp.write(chunk)
|
||||
byte_count += len(chunk)
|
||||
tmp.close()
|
||||
return tmp.name, byte_count
|
||||
except BaseException:
|
||||
tmp.close()
|
||||
_cleanup_tempfile(tmp.name)
|
||||
raise
|
||||
|
||||
|
||||
def _write_file_to_store(store: FileStore, name: str, path: str) -> None:
|
||||
"""Stream a temp file to the store, never buffering the whole archive in RAM.
|
||||
|
||||
The download was streamed to a tempfile precisely to avoid holding the
|
||||
archive in memory (see ``_stream_to_tempfile``); ``write_from_path`` keeps
|
||||
that guarantee on upload (GCS/local/S3 stream from disk) instead of reading
|
||||
the whole file back with ``store.write(name, f.read())``.
|
||||
"""
|
||||
store.write_from_path(name, path)
|
||||
|
||||
|
||||
async def archive_workspace(
|
||||
httpx_client: httpx.AsyncClient,
|
||||
runtime: dict[str, Any],
|
||||
sandbox_id: str,
|
||||
*,
|
||||
archive_path: str,
|
||||
conversation_id: str | None = None,
|
||||
) -> bool:
|
||||
"""Archive the workspace at ``archive_path``; return whether delete may proceed.
|
||||
|
||||
``archive_path`` is resolved by the caller — the path pinned at conversation
|
||||
creation, NOT re-derived here from live settings — so a capture can never be
|
||||
misrouted to the wrong directory. The agent-server descends from it to the
|
||||
cloned repo.
|
||||
|
||||
Returns True when the workspace was archived, when the path holds nothing to
|
||||
archive (agent-server 400: not a directory / not a git repo), or when
|
||||
archiving failed but is not required (best-effort). Returns False when
|
||||
archiving is required and either hit a transient failure (5xx / network /
|
||||
422 / 429) or could not confirm a capture (401 auth / 404 missing path), so
|
||||
the caller leaves the sandbox intact for the idle-reap retry. Never raises.
|
||||
|
||||
A pure configuration error (unsupported RUNTIME_FILE_ARCHIVE_FORMAT, or
|
||||
RUNTIME_FILE_ARCHIVE_BUCKET unset) cannot be fixed by retrying, so it is
|
||||
logged loudly and the delete is allowed to proceed rather than wedging every
|
||||
delete forever when archiving is required.
|
||||
"""
|
||||
agent_server_url = runtime.get('url')
|
||||
session_api_key = runtime.get('session_api_key')
|
||||
if not agent_server_url:
|
||||
_logger.warning(
|
||||
'Workspace archive skipped for %s: runtime has no agent-server URL',
|
||||
sandbox_id,
|
||||
)
|
||||
return not archive_required()
|
||||
if not _archive_bucket():
|
||||
# Misconfiguration, not a transient failure: no amount of retrying makes
|
||||
# a missing bucket appear. Proceed (with a loud error) so a
|
||||
# REQUIRED-without-bucket setup does not block every sandbox delete.
|
||||
_logger.error(
|
||||
'Workspace archive enabled for %s but RUNTIME_FILE_ARCHIVE_BUCKET '
|
||||
'is not set; proceeding with delete (fix the config to capture)',
|
||||
sandbox_id,
|
||||
)
|
||||
return True
|
||||
|
||||
formats = _formats_to_capture()
|
||||
if formats is None:
|
||||
# Unsupported RUNTIME_FILE_ARCHIVE_FORMAT is a pure config error, exactly
|
||||
# like the unset bucket above: no retry makes a valid format appear, so
|
||||
# proceed loudly rather than wedging every REQUIRED delete forever (the
|
||||
# app-server has no idle-reap backstop). Validated here so a bad format
|
||||
# never reaches the producer (which would 422 it).
|
||||
_logger.error(
|
||||
'Workspace archive for %s: unsupported RUNTIME_FILE_ARCHIVE_FORMAT '
|
||||
'%r (valid: %s); proceeding with delete (fix the config to capture)',
|
||||
sandbox_id,
|
||||
_archive_format(),
|
||||
['git-delta', 'tar.gz', 'both'],
|
||||
)
|
||||
return True
|
||||
|
||||
headers = {'X-Session-API-Key': session_api_key} if session_api_key else {}
|
||||
# For cloud conversations the sandbox id is the conversation_id.hex.
|
||||
conversation_key = conversation_id or sandbox_id
|
||||
ts = utc_now().strftime('%Y%m%dT%H%M%SZ')
|
||||
# Key by conversation, not just sandbox: under grouping a sandbox is shared by
|
||||
# siblings, and the 1s ts is not unique — without the conversation segment two
|
||||
# sibling captures in the same second overwrite each other at the object level.
|
||||
base_path = f'{_archive_prefix()}/{sandbox_id}/{conversation_key}/{ts}'
|
||||
|
||||
# 'both' uploads each format under its own suffix ({ts}.patch + {ts}.tar.gz),
|
||||
# each with its own manifest. base_commit only rides the git-delta response
|
||||
# header, so capture it there and reuse it for the tar.gz manifest.
|
||||
# A retry under REQUIRED re-uploads under a fresh {ts}; any blob/manifest left
|
||||
# by a partially-failed prior attempt becomes an orphan reaped by the bucket
|
||||
# lifecycle policy (we favor capture completeness over upload dedup).
|
||||
retryable_failure = False
|
||||
# A capture we could not confirm happened (401 auth / 404 missing path):
|
||||
# under REQUIRED this must NOT permit teardown — it is the misrouted-path
|
||||
# symptom this feature most needs to guard against, not "nothing to archive".
|
||||
unconfirmed_capture = False
|
||||
base_commit = ''
|
||||
# One store per call (not per format) — building it lazily spins up a client.
|
||||
store = _get_archive_file_store()
|
||||
for fmt in formats:
|
||||
suffix = _ARCHIVE_SUFFIX[fmt]
|
||||
tmp_path: str | None = None
|
||||
byte_count = 0
|
||||
try:
|
||||
async with httpx_client.stream(
|
||||
'GET',
|
||||
f'{agent_server_url}/api/file/archive',
|
||||
params=_archive_request_params(archive_path, fmt),
|
||||
headers=headers,
|
||||
timeout=_archive_timeout(),
|
||||
) as response:
|
||||
if response.status_code != 200:
|
||||
code = response.status_code
|
||||
if code == 400:
|
||||
# Path exists but holds no archivable repo (not a git
|
||||
# repo / not a directory). A positive "nothing here" —
|
||||
# safe to skip this format and proceed.
|
||||
detail = 'nothing to archive'
|
||||
elif code in (401, 404):
|
||||
# 401 auth rejected / 404 path missing: the capture did
|
||||
# NOT happen and this is not a confirmed-empty workspace,
|
||||
# so it must block a REQUIRED delete (idle reap retries).
|
||||
unconfirmed_capture = True
|
||||
detail = 'capture unconfirmed (auth/path)'
|
||||
else:
|
||||
# 422 / 429 / 5xx — transient.
|
||||
retryable_failure = True
|
||||
detail = 'retryable failure'
|
||||
_logger.warning(
|
||||
'Workspace archive (%s) for %s: agent-server returned %s; %s',
|
||||
fmt,
|
||||
sandbox_id,
|
||||
code,
|
||||
detail,
|
||||
)
|
||||
continue
|
||||
header_base = response.headers.get('X-Archive-Base-Commit', '')
|
||||
if header_base:
|
||||
base_commit = header_base
|
||||
# Stream to disk so the archive never sits whole in RAM.
|
||||
tmp_path, byte_count = await _stream_to_tempfile(response)
|
||||
except Exception as e:
|
||||
# Network/timeout error: genuinely transient.
|
||||
_logger.warning(
|
||||
'Workspace archive fetch (%s) failed for %s: %s', fmt, sandbox_id, e
|
||||
)
|
||||
retryable_failure = True
|
||||
_cleanup_tempfile(tmp_path)
|
||||
continue
|
||||
|
||||
assert tmp_path is not None # set on the 200 path above
|
||||
try:
|
||||
await asyncio.to_thread(
|
||||
_write_file_to_store, store, f'{base_path}.{suffix}', tmp_path
|
||||
)
|
||||
manifest = json.dumps(
|
||||
{
|
||||
'sandbox_id': sandbox_id,
|
||||
'conversation_id': conversation_key,
|
||||
'phase': 'final',
|
||||
'base_commit': base_commit,
|
||||
'format': fmt,
|
||||
'source_path': archive_path,
|
||||
'byte_count': byte_count,
|
||||
'created_at': ts,
|
||||
},
|
||||
sort_keys=True,
|
||||
).encode('utf-8')
|
||||
await asyncio.to_thread(
|
||||
store.write, f'{base_path}.{suffix}.manifest.json', manifest
|
||||
)
|
||||
_logger.info(
|
||||
'Archived workspace (%s) for %s (%d bytes) to %s.%s',
|
||||
fmt,
|
||||
sandbox_id,
|
||||
byte_count,
|
||||
base_path,
|
||||
suffix,
|
||||
)
|
||||
except Exception as e:
|
||||
_logger.exception(
|
||||
'Workspace archive upload (%s) failed for %s: %s', fmt, sandbox_id, e
|
||||
)
|
||||
retryable_failure = True
|
||||
finally:
|
||||
_cleanup_tempfile(tmp_path)
|
||||
|
||||
# Deletion may proceed unless archiving is REQUIRED and we either hit a
|
||||
# retryable failure or could not confirm a capture (401/404) — both leave us
|
||||
# short of the data we were meant to preserve.
|
||||
if archive_required() and (retryable_failure or unconfirmed_capture):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def initial_archive_enabled() -> bool:
|
||||
"""Whether to capture the workspace's INITIAL state (before the first step).
|
||||
|
||||
Independent of ``RUNTIME_FILE_ARCHIVE_ENABLED`` (the delete/pause capture of
|
||||
the *final* state) so the pre-agent snapshot can be toggled on its own. Off
|
||||
by default — like every other capture knob, nothing happens until enabled.
|
||||
"""
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', 'false').lower() in (
|
||||
'true',
|
||||
'1',
|
||||
)
|
||||
|
||||
|
||||
def initial_archive_deadline() -> float:
|
||||
"""Hard ceiling (seconds) on how long the inline, pre-setup initial snapshot
|
||||
may delay conversation startup.
|
||||
|
||||
The snapshot is awaited before setup.sh so it captures the repo exactly as
|
||||
cloned, with no concurrent mutation. ``wait_for`` only charges the time
|
||||
actually spent, so a fast snapshot adds no latency; this caps the worst case
|
||||
(a large repo or a hung endpoint) so it can never dominate startup. An overrun
|
||||
is logged and startup proceeds without the snapshot. A fresh clone's tar.gz is
|
||||
just the repo tree, so the default comfortably covers typical repos; raise it
|
||||
for a very large monorepo.
|
||||
"""
|
||||
return _float_env('RUNTIME_FILE_ARCHIVE_INITIAL_DEADLINE', 120.0)
|
||||
|
||||
|
||||
def _initial_archive_format() -> str:
|
||||
"""Format for the initial snapshot. Defaults to a self-contained tar.gz.
|
||||
|
||||
At conversation start the working tree has no changes yet, so a ``git-delta``
|
||||
would be empty; a full ``tar.gz`` is the only format that captures anything
|
||||
and, unlike a delta keyed to ``base_commit``, it survives the upstream repo
|
||||
or branch later disappearing (the fragile re-clone path we want to avoid).
|
||||
"""
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_INITIAL_FORMAT', 'tar.gz')
|
||||
|
||||
|
||||
async def archive_initial_workspace(
|
||||
httpx_client: httpx.AsyncClient,
|
||||
*,
|
||||
agent_server_url: str | None,
|
||||
session_api_key: str | None,
|
||||
project_dir: str,
|
||||
sandbox_id: str,
|
||||
conversation_id: str | None = None,
|
||||
base_commit: str = '',
|
||||
) -> bool:
|
||||
"""Snapshot the workspace BEFORE the agent's first step; return success.
|
||||
|
||||
Captures the repo exactly as cloned (option A — the pre- vs post-setup choice
|
||||
is the open design question tracked in All-Hands-AI/infra#1444) as a
|
||||
self-contained ``tar.gz`` plus a ``phase=initial`` manifest, so the snapshot
|
||||
records the true starting state even if the source repo later disappears.
|
||||
|
||||
This is strictly best-effort: it never raises and never blocks conversation
|
||||
startup. A failure (feature off, misconfig, agent-server hiccup) just means no
|
||||
initial snapshot for this run, logged and swallowed. Returns True only when an
|
||||
archive was actually written.
|
||||
"""
|
||||
if not initial_archive_enabled():
|
||||
return False
|
||||
if not agent_server_url:
|
||||
_logger.warning(
|
||||
'Initial workspace archive skipped for %s: no agent-server URL',
|
||||
sandbox_id,
|
||||
)
|
||||
return False
|
||||
if not _archive_bucket():
|
||||
_logger.error(
|
||||
'Initial workspace archive enabled for %s but '
|
||||
'RUNTIME_FILE_ARCHIVE_BUCKET is not set; skipping initial snapshot',
|
||||
sandbox_id,
|
||||
)
|
||||
return False
|
||||
|
||||
fmt = _initial_archive_format()
|
||||
if fmt not in _ARCHIVE_SUFFIX:
|
||||
_logger.error(
|
||||
'Initial workspace archive for %s: unsupported '
|
||||
'RUNTIME_FILE_ARCHIVE_INITIAL_FORMAT %r (valid: %s); skipping',
|
||||
sandbox_id,
|
||||
fmt,
|
||||
sorted(_ARCHIVE_SUFFIX),
|
||||
)
|
||||
return False
|
||||
suffix = _ARCHIVE_SUFFIX[fmt]
|
||||
headers = {'X-Session-API-Key': session_api_key} if session_api_key else {}
|
||||
|
||||
tmp_path: str | None = None
|
||||
byte_count = 0
|
||||
try:
|
||||
async with httpx_client.stream(
|
||||
'GET',
|
||||
f'{agent_server_url}/api/file/archive',
|
||||
params=_archive_request_params(project_dir, fmt),
|
||||
headers=headers,
|
||||
timeout=_archive_timeout(),
|
||||
) as response:
|
||||
if response.status_code != 200:
|
||||
_logger.warning(
|
||||
'Initial workspace archive for %s: agent-server returned %s; '
|
||||
'no initial snapshot',
|
||||
sandbox_id,
|
||||
response.status_code,
|
||||
)
|
||||
return False
|
||||
# tar.gz carries no base-commit header (git-delta sets it); fall back
|
||||
# to the caller-provided HEAD sha so the snapshot still records the
|
||||
# commit it came from.
|
||||
captured_base = (
|
||||
response.headers.get('X-Archive-Base-Commit', '') or base_commit
|
||||
)
|
||||
# Stream to disk so the archive never sits whole in RAM.
|
||||
tmp_path, byte_count = await _stream_to_tempfile(response)
|
||||
except Exception as e:
|
||||
_logger.warning(
|
||||
'Initial workspace archive fetch failed for %s: %s', sandbox_id, e
|
||||
)
|
||||
_cleanup_tempfile(tmp_path)
|
||||
return False
|
||||
|
||||
assert tmp_path is not None # set on the 200 path above
|
||||
try:
|
||||
store = _get_archive_file_store()
|
||||
ts = utc_now().strftime('%Y%m%dT%H%M%SZ')
|
||||
conversation_key = conversation_id or sandbox_id
|
||||
# Key by conversation (siblings share a grouped sandbox) and nest under
|
||||
# /initial/ so it never collides with that conversation's final capture
|
||||
# ({prefix}/{sandbox_id}/{conversation_key}/{ts}).
|
||||
blob_name = (
|
||||
f'{_archive_prefix()}/{sandbox_id}/{conversation_key}/initial/{ts}.{suffix}'
|
||||
)
|
||||
await asyncio.to_thread(_write_file_to_store, store, blob_name, tmp_path)
|
||||
manifest = json.dumps(
|
||||
{
|
||||
'sandbox_id': sandbox_id,
|
||||
'conversation_id': conversation_key,
|
||||
'phase': 'initial',
|
||||
'base_commit': captured_base,
|
||||
'format': fmt,
|
||||
'source_path': project_dir,
|
||||
'byte_count': byte_count,
|
||||
'created_at': ts,
|
||||
},
|
||||
sort_keys=True,
|
||||
).encode('utf-8')
|
||||
# Shared contract: manifest = blob + '.manifest.json' (was dropping the
|
||||
# format suffix, so a downstream enricher could never locate it).
|
||||
await asyncio.to_thread(store.write, f'{blob_name}.manifest.json', manifest)
|
||||
_logger.info(
|
||||
'Archived INITIAL workspace for %s (%d bytes) to %s',
|
||||
sandbox_id,
|
||||
byte_count,
|
||||
blob_name,
|
||||
)
|
||||
return True
|
||||
except Exception as e:
|
||||
_logger.exception(
|
||||
'Initial workspace archive upload failed for %s: %s', sandbox_id, e
|
||||
)
|
||||
return False
|
||||
finally:
|
||||
_cleanup_tempfile(tmp_path)
|
||||
@@ -92,6 +92,23 @@ class SandboxGroupingStrategy(str, Enum):
|
||||
ADD_TO_ANY = 'ADD_TO_ANY' # Add to any available sandbox (first found)
|
||||
|
||||
|
||||
def grouped_workspace_dir(
|
||||
base_working_dir: str,
|
||||
grouping_strategy: SandboxGroupingStrategy,
|
||||
conversation_id_hex: str,
|
||||
) -> str:
|
||||
"""Workspace dir for a conversation given the grouping strategy.
|
||||
|
||||
Single source of truth for the relocation used at conversation start and at
|
||||
archive time. Under any grouping strategy the workspace is nested under the
|
||||
conversation id so co-located conversations stay isolated; NO_GROUPING keeps
|
||||
the bare base dir.
|
||||
"""
|
||||
if grouping_strategy == SandboxGroupingStrategy.NO_GROUPING:
|
||||
return base_working_dir
|
||||
return f'{base_working_dir}/{conversation_id_hex}'
|
||||
|
||||
|
||||
# Fields the batch ``update()`` method refuses to touch:
|
||||
# - ``secrets_store`` is frozen (Pydantic would raise).
|
||||
# - ``llm_profiles`` is off-limits for the generic settings POST; profile
|
||||
|
||||
@@ -21,6 +21,7 @@ from openhands.app_server.app_conversation.app_conversation_models import (
|
||||
)
|
||||
from openhands.app_server.app_conversation.app_conversation_router import (
|
||||
AgentServerContext,
|
||||
_finalize_sandbox_delete,
|
||||
batch_get_app_conversations,
|
||||
count_app_conversations,
|
||||
get_conversation_git_changes,
|
||||
@@ -1082,3 +1083,88 @@ class TestGitProxyEndpoints:
|
||||
|
||||
assert exc_info.value.status_code == status.HTTP_409_CONFLICT
|
||||
assert 'paused' in exc_info.value.detail.lower()
|
||||
|
||||
|
||||
class TestFinalizeSandboxDelete:
|
||||
"""The detached finalizer captures THIS conversation's workspace first
|
||||
(archive_conversation_workspace, while the runtime is still up), then tears the
|
||||
sandbox down only when no other conversation still references it (count == 0).
|
||||
delete_sandbox is sandbox-scoped — no conversation_id. A REQUIRED archive
|
||||
failure (archive returns False) keeps the sandbox for the idle reap."""
|
||||
|
||||
def _deps(self, conversation_count: int, archived: bool = True):
|
||||
info_service = MagicMock()
|
||||
info_service.count_conversations_by_sandbox_id = AsyncMock(
|
||||
return_value=conversation_count
|
||||
)
|
||||
sandbox_service = MagicMock()
|
||||
sandbox_service.delete_sandbox = AsyncMock(return_value=True)
|
||||
sandbox_service.archive_conversation_workspace = AsyncMock(
|
||||
return_value=archived
|
||||
)
|
||||
db_session = AsyncMock()
|
||||
httpx_client = AsyncMock()
|
||||
return sandbox_service, info_service, db_session, httpx_client
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_archives_then_deletes_when_unreferenced(self):
|
||||
sandbox_service, info_service, db_session, httpx_client = self._deps(0)
|
||||
conv = uuid4()
|
||||
await _finalize_sandbox_delete(
|
||||
sandbox_service, info_service, 'sbx-1', db_session, httpx_client, conv
|
||||
)
|
||||
# Workspace captured first (conversation-scoped), then the sandbox torn
|
||||
# down sandbox-scoped — note delete_sandbox takes no conversation_id.
|
||||
sandbox_service.archive_conversation_workspace.assert_awaited_once_with(
|
||||
'sbx-1', conversation_id=conv.hex, workspace_path=None
|
||||
)
|
||||
sandbox_service.delete_sandbox.assert_awaited_once_with('sbx-1')
|
||||
db_session.commit.assert_awaited_once()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_archives_but_keeps_shared_sandbox(self):
|
||||
sandbox_service, info_service, db_session, httpx_client = self._deps(2)
|
||||
conv = uuid4()
|
||||
await _finalize_sandbox_delete(
|
||||
sandbox_service, info_service, 'sbx-1', db_session, httpx_client, conv
|
||||
)
|
||||
# Shared sandbox: still capture this conversation, but don't tear it down.
|
||||
sandbox_service.archive_conversation_workspace.assert_awaited_once_with(
|
||||
'sbx-1', conversation_id=conv.hex, workspace_path=None
|
||||
)
|
||||
sandbox_service.delete_sandbox.assert_not_called()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_required_archive_failure_keeps_sandbox(self):
|
||||
# archive_conversation_workspace returns False (REQUIRED archive failed).
|
||||
sandbox_service, info_service, db_session, httpx_client = self._deps(
|
||||
0, archived=False
|
||||
)
|
||||
conv = uuid4()
|
||||
await _finalize_sandbox_delete(
|
||||
sandbox_service, info_service, 'sbx-1', db_session, httpx_client, conv
|
||||
)
|
||||
# Sandbox + runtime kept for the idle reap: no count check, no delete.
|
||||
sandbox_service.archive_conversation_workspace.assert_awaited_once()
|
||||
info_service.count_conversations_by_sandbox_id.assert_not_called()
|
||||
sandbox_service.delete_sandbox.assert_not_called()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_forwards_pinned_workspace_path(self):
|
||||
sandbox_service, info_service, db_session, httpx_client = self._deps(0)
|
||||
conv = uuid4()
|
||||
await _finalize_sandbox_delete(
|
||||
sandbox_service,
|
||||
info_service,
|
||||
'sbx-1',
|
||||
db_session,
|
||||
httpx_client,
|
||||
conv,
|
||||
workspace_path='/home/openhands/workspace/' + conv.hex,
|
||||
)
|
||||
# The path pinned at creation is forwarded to the capture verbatim.
|
||||
sandbox_service.archive_conversation_workspace.assert_awaited_once_with(
|
||||
'sbx-1',
|
||||
conversation_id=conv.hex,
|
||||
workspace_path='/home/openhands/workspace/' + conv.hex,
|
||||
)
|
||||
|
||||
@@ -1893,7 +1893,9 @@ class TestLiveStatusAppConversationService:
|
||||
task.sandbox_id = self.mock_sandbox.id
|
||||
yield task
|
||||
|
||||
async def mock_run_setup_scripts(task, sandbox, workspace, agent_server_url):
|
||||
async def mock_run_setup_scripts(
|
||||
task, sandbox, workspace, agent_server_url, conversation_id
|
||||
):
|
||||
yield task
|
||||
|
||||
self.service._wait_for_sandbox_start = mock_wait_for_sandbox
|
||||
@@ -1997,7 +1999,9 @@ class TestLiveStatusAppConversationService:
|
||||
task.sandbox_id = self.mock_sandbox.id
|
||||
yield task
|
||||
|
||||
async def mock_run_setup_scripts(task, sandbox, workspace, agent_server_url):
|
||||
async def mock_run_setup_scripts(
|
||||
task, sandbox, workspace, agent_server_url, conversation_id
|
||||
):
|
||||
yield task
|
||||
|
||||
self.service._wait_for_sandbox_start = mock_wait_for_sandbox
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user