mirror of
https://github.com/OpenHands/OpenHands.git
synced 2026-10-07 16:38:34 +08:00
feat(app-server): enrich final archive manifests and remove initial snapshots (#15058)
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Co-authored-by: Juan Michelini <juan@juan.com.uy> Co-authored-by: openhands <openhands@all-hands.dev>
This commit is contained in:
co-authored by
Claude Opus 4.8
Juan Michelini
openhands
parent
ef1a3c71a0
commit
4a607c6d94
@@ -1,4 +1,3 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import shlex
|
||||
@@ -29,7 +28,6 @@ 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.settings.settings_models import MarketplaceRegistration
|
||||
from openhands.app_server.user.user_context import UserContext
|
||||
@@ -287,28 +285,6 @@ 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)
|
||||
@@ -326,55 +302,6 @@ 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,
|
||||
|
||||
@@ -20,6 +20,7 @@ import logging
|
||||
import os
|
||||
import tempfile
|
||||
from typing import Any
|
||||
from urllib.parse import unquote
|
||||
|
||||
import httpx
|
||||
|
||||
@@ -33,6 +34,29 @@ _logger = logging.getLogger(__name__)
|
||||
# anything else 422s, so validate before issuing the request.
|
||||
_ARCHIVE_SUFFIX = {'git-delta': 'patch', 'tar.gz': 'tar.gz'}
|
||||
|
||||
_REPO_METADATA_HEADERS = {
|
||||
'repo_remote': 'X-Archive-Repo-Remote',
|
||||
'branch': 'X-Archive-Branch',
|
||||
'head_commit': 'X-Archive-Head-Commit',
|
||||
}
|
||||
_REPO_ROOT_HEADER = 'X-Archive-Repo-Root'
|
||||
|
||||
|
||||
def _extract_repo_metadata(
|
||||
headers: Any, existing: dict[str, str] | None = None
|
||||
) -> dict[str, str]:
|
||||
"""Pull repo-identity fields from an archive response's headers.
|
||||
|
||||
Falls back to ``existing`` per-key so a later response missing a header
|
||||
(e.g. a format the agent-server can't probe) doesn't clobber a value an
|
||||
earlier response in the same capture already found.
|
||||
"""
|
||||
fallback = existing or {}
|
||||
return {
|
||||
key: unquote(headers.get(header_name, '')) or fallback.get(key, '')
|
||||
for key, header_name in _REPO_METADATA_HEADERS.items()
|
||||
}
|
||||
|
||||
|
||||
def _archive_request_params(path: str, fmt: str) -> dict[str, str]:
|
||||
"""Query params for GET /api/file/archive.
|
||||
@@ -59,6 +83,11 @@ def archive_required() -> bool:
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_REQUIRED', 'false').lower() in ('true', '1')
|
||||
|
||||
|
||||
def _manifest_enrichment_enabled() -> bool:
|
||||
"""Whether to probe packages and runtime versions into the manifest."""
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_ENRICH', 'true').lower() in ('true', '1')
|
||||
|
||||
|
||||
def _archive_bucket() -> str:
|
||||
return os.getenv('RUNTIME_FILE_ARCHIVE_BUCKET', '')
|
||||
|
||||
@@ -119,8 +148,7 @@ 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().
|
||||
# delete finalizer, so a long wait here doesn't block any user request.
|
||||
return _float_env('RUNTIME_FILE_ARCHIVE_TIMEOUT', 660.0)
|
||||
|
||||
|
||||
@@ -176,6 +204,133 @@ def _write_file_to_store(store: FileStore, name: str, path: str) -> None:
|
||||
store.write_from_path(name, path)
|
||||
|
||||
|
||||
_PROBE_TIMEOUT = 15.0
|
||||
_MAX_PACKAGES_PER_MANAGER = 2000
|
||||
_ENVIRONMENT_CMD = (
|
||||
'echo "python=$(.venv/bin/python --version 2>/dev/null '
|
||||
'|| python3 --version 2>/dev/null || true)"; '
|
||||
'echo "node=$(node --version 2>/dev/null || true)"; '
|
||||
'echo "os=$(. /etc/os-release 2>/dev/null; echo "${ID:-} ${VERSION_ID:-}")"'
|
||||
)
|
||||
|
||||
|
||||
def _parse_pip_list(out: str) -> dict[str, str]:
|
||||
"""Parse ``pip list --format=json`` output into ``{name: version}``."""
|
||||
try:
|
||||
data = json.loads(out or '[]')
|
||||
except json.JSONDecodeError:
|
||||
return {}
|
||||
result: dict[str, str] = {}
|
||||
for item in data if isinstance(data, list) else []:
|
||||
if not isinstance(item, dict):
|
||||
continue
|
||||
name, version = item.get('name'), item.get('version')
|
||||
if name and version:
|
||||
result[name] = version
|
||||
if len(result) >= _MAX_PACKAGES_PER_MANAGER:
|
||||
break
|
||||
return result
|
||||
|
||||
|
||||
def _parse_npm_ls(out: str) -> dict[str, str]:
|
||||
"""Parse ``npm ls --json`` output (top-level deps) into ``{name: version}``."""
|
||||
try:
|
||||
data = json.loads(out or '{}')
|
||||
except json.JSONDecodeError:
|
||||
return {}
|
||||
deps = data.get('dependencies') if isinstance(data, dict) else None
|
||||
if not isinstance(deps, dict):
|
||||
return {}
|
||||
result: dict[str, str] = {}
|
||||
for name, meta in deps.items():
|
||||
version = meta.get('version') if isinstance(meta, dict) else None
|
||||
if name and version:
|
||||
result[name] = version
|
||||
if len(result) >= _MAX_PACKAGES_PER_MANAGER:
|
||||
break
|
||||
return result
|
||||
|
||||
|
||||
def _parse_runtime(out: str) -> dict[str, str]:
|
||||
"""Parse the `key=value` runtime lines into ``{python, node, os}``."""
|
||||
result: dict[str, str] = {}
|
||||
for line in out.splitlines():
|
||||
key, _, value = line.partition('=')
|
||||
key, value = key.strip(), value.strip()
|
||||
if key not in ('python', 'node', 'os') or not value:
|
||||
continue
|
||||
if key == 'python':
|
||||
value = value.removeprefix('Python ').strip()
|
||||
elif key == 'node':
|
||||
value = value.lstrip('v')
|
||||
if value:
|
||||
result[key] = value
|
||||
return result
|
||||
|
||||
|
||||
async def _run_probe(
|
||||
httpx_client: httpx.AsyncClient,
|
||||
agent_server_url: str,
|
||||
headers: dict[str, str],
|
||||
cwd: str,
|
||||
command: str,
|
||||
) -> str:
|
||||
"""Run one command in the workspace; '' on any failure (never raises)."""
|
||||
try:
|
||||
response = await httpx_client.post(
|
||||
f'{agent_server_url}/api/bash/execute_bash_command',
|
||||
json={'command': command, 'cwd': cwd, 'timeout': _PROBE_TIMEOUT},
|
||||
headers=headers,
|
||||
timeout=_PROBE_TIMEOUT + 1,
|
||||
)
|
||||
if response.status_code != 200:
|
||||
return ''
|
||||
data = response.json()
|
||||
except Exception as e:
|
||||
_logger.debug('Workspace probe %r failed: %s', command, e)
|
||||
return ''
|
||||
output = data.get('stdout') if isinstance(data, dict) else None
|
||||
return output if isinstance(output, str) else ''
|
||||
|
||||
|
||||
async def _probe_workspace(
|
||||
httpx_client: httpx.AsyncClient,
|
||||
agent_server_url: str,
|
||||
headers: dict[str, str],
|
||||
cwd: str,
|
||||
) -> dict[str, Any]:
|
||||
"""Collect installed package and runtime versions."""
|
||||
if not _manifest_enrichment_enabled():
|
||||
return {}
|
||||
|
||||
async def _run(command: str) -> str:
|
||||
return await _run_probe(httpx_client, agent_server_url, headers, cwd, command)
|
||||
|
||||
result: dict[str, Any] = {}
|
||||
packages: dict[str, dict[str, str]] = {}
|
||||
# Prefer a project venv / uv (where the agent's installs actually land)
|
||||
# before the system interpreter; each emits the same --format=json shape.
|
||||
pip = _parse_pip_list(
|
||||
await _run(
|
||||
'.venv/bin/python -m pip list --format=json 2>/dev/null '
|
||||
'|| uv pip list --format=json 2>/dev/null '
|
||||
'|| python3 -m pip list --format=json 2>/dev/null'
|
||||
)
|
||||
)
|
||||
if pip:
|
||||
packages['pip'] = pip
|
||||
npm = _parse_npm_ls(await _run('npm ls --json --depth=0 2>/dev/null'))
|
||||
if npm:
|
||||
packages['npm'] = npm
|
||||
if packages:
|
||||
result['packages'] = packages
|
||||
|
||||
environment = _parse_runtime(await _run(_ENVIRONMENT_CMD))
|
||||
if environment:
|
||||
result['environment'] = environment
|
||||
return result
|
||||
|
||||
|
||||
async def archive_workspace(
|
||||
httpx_client: httpx.AsyncClient,
|
||||
runtime: dict[str, Any],
|
||||
@@ -259,6 +414,13 @@ async def archive_workspace(
|
||||
# symptom this feature most needs to guard against, not "nothing to archive".
|
||||
unconfirmed_capture = False
|
||||
base_commit = ''
|
||||
# Repo identity rides the response headers (empty against an agent-server
|
||||
# image predating them — graceful) and is reused across formats so each
|
||||
# manifest is self-describing (repo / branch / captured HEAD).
|
||||
repo_metadata = dict.fromkeys(_REPO_METADATA_HEADERS, '')
|
||||
enrichment: dict[str, Any] = {}
|
||||
enrichment_probed = False
|
||||
probe_path = archive_path
|
||||
# One store per call (not per format) — building it lazily spins up a client.
|
||||
store = _get_archive_file_store()
|
||||
for fmt in formats:
|
||||
@@ -301,8 +463,25 @@ async def archive_workspace(
|
||||
header_base = response.headers.get('X-Archive-Base-Commit', '')
|
||||
if header_base:
|
||||
base_commit = header_base
|
||||
repo_metadata = _extract_repo_metadata(response.headers, repo_metadata)
|
||||
response_repo_root = response.headers.get(_REPO_ROOT_HEADER, '')
|
||||
if response_repo_root:
|
||||
probe_path = unquote(response_repo_root)
|
||||
# Stream to disk so the archive never sits whole in RAM.
|
||||
tmp_path, byte_count = await _stream_to_tempfile(response)
|
||||
if not enrichment_probed:
|
||||
enrichment_probed = True
|
||||
try:
|
||||
enrichment = await _probe_workspace(
|
||||
httpx_client,
|
||||
agent_server_url,
|
||||
headers,
|
||||
probe_path,
|
||||
)
|
||||
except Exception as e:
|
||||
_logger.debug(
|
||||
'Workspace enrichment skipped for %s: %s', sandbox_id, e
|
||||
)
|
||||
except Exception as e:
|
||||
# Network/timeout error: genuinely transient.
|
||||
_logger.warning(
|
||||
@@ -323,6 +502,9 @@ async def archive_workspace(
|
||||
'conversation_id': conversation_key,
|
||||
'phase': 'final',
|
||||
'base_commit': base_commit,
|
||||
**repo_metadata,
|
||||
'packages': enrichment.get('packages', {}),
|
||||
'environment': enrichment.get('environment', {}),
|
||||
'format': fmt,
|
||||
'source_path': archive_path,
|
||||
'byte_count': byte_count,
|
||||
@@ -355,170 +537,3 @@ async def archive_workspace(
|
||||
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)
|
||||
|
||||
@@ -106,6 +106,7 @@ def _make_stream_response(
|
||||
resp = MagicMock()
|
||||
resp.status_code = status_code
|
||||
resp.headers = headers or {}
|
||||
resp.stream_closed = False
|
||||
|
||||
async def _aiter_bytes():
|
||||
yield content
|
||||
@@ -124,9 +125,13 @@ def _stream_client(resp_or_map):
|
||||
@asynccontextmanager
|
||||
async def _stream(method, url, **kwargs):
|
||||
if isinstance(resp_or_map, dict):
|
||||
yield resp_or_map[kwargs['params']['format']]
|
||||
response = resp_or_map[kwargs['params']['format']]
|
||||
else:
|
||||
yield resp_or_map
|
||||
response = resp_or_map
|
||||
try:
|
||||
yield response
|
||||
finally:
|
||||
response.stream_closed = True
|
||||
|
||||
client.stream = MagicMock(side_effect=_stream)
|
||||
return client
|
||||
@@ -2099,11 +2104,20 @@ class TestArchiveWorkspaceHelper:
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_BUCKET', 'archive-bkt')
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_FORMAT', 'git-delta')
|
||||
|
||||
client = _stream_client(
|
||||
_make_stream_response(
|
||||
200, b'patch-bytes', {'X-Archive-Base-Commit': 'abc123'}
|
||||
)
|
||||
archive_response = _make_stream_response(
|
||||
200,
|
||||
b'patch-bytes',
|
||||
{
|
||||
'X-Archive-Base-Commit': 'abc123',
|
||||
'X-Archive-Repo-Remote': (
|
||||
'https%3A%2F%2Fgithub.com%2Fexample%2Frepo.git'
|
||||
),
|
||||
'X-Archive-Branch': 'feature-x',
|
||||
'X-Archive-Head-Commit': 'def456',
|
||||
'X-Archive-Repo-Root': '%2Fworkspace%2Fproject%2Frepo',
|
||||
},
|
||||
)
|
||||
client = _stream_client(archive_response)
|
||||
store = MagicMock()
|
||||
writes: dict[str, bytes] = {}
|
||||
store.write.side_effect = lambda path, data: writes.__setitem__(path, data)
|
||||
@@ -2113,8 +2127,18 @@ class TestArchiveWorkspaceHelper:
|
||||
path, open(src, 'rb').read()
|
||||
)
|
||||
|
||||
with patch.object(
|
||||
workspace_archive, '_get_archive_file_store', return_value=store
|
||||
probe_workspace = AsyncMock()
|
||||
|
||||
async def _probe(*args):
|
||||
assert archive_response.stream_closed
|
||||
return {'packages': {'npm': {'example': '1.0.0'}}}
|
||||
|
||||
probe_workspace.side_effect = _probe
|
||||
with (
|
||||
patch.object(
|
||||
workspace_archive, '_get_archive_file_store', return_value=store
|
||||
),
|
||||
patch.object(workspace_archive, '_probe_workspace', probe_workspace),
|
||||
):
|
||||
ok = await workspace_archive.archive_workspace(
|
||||
client,
|
||||
@@ -2139,6 +2163,17 @@ class TestArchiveWorkspaceHelper:
|
||||
assert manifest['base_commit'] == 'abc123'
|
||||
assert manifest['conversation_id'] == 'conv-1'
|
||||
assert manifest['source_path'] == '/workspace/project'
|
||||
# Repo identity from the response headers makes the blob self-describing.
|
||||
assert manifest['repo_remote'] == 'https://github.com/example/repo.git'
|
||||
assert manifest['branch'] == 'feature-x'
|
||||
assert manifest['head_commit'] == 'def456'
|
||||
assert manifest['packages'] == {'npm': {'example': '1.0.0'}}
|
||||
probe_workspace.assert_awaited_once_with(
|
||||
client,
|
||||
'https://sandbox.example.com',
|
||||
{'X-Session-API-Key': 'test-session-key'},
|
||||
'/workspace/project/repo',
|
||||
)
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_archive_missing_base_commit_header_defaults_empty(self, monkeypatch):
|
||||
@@ -2171,6 +2206,10 @@ class TestArchiveWorkspaceHelper:
|
||||
manifest_path = next(p for p in writes if p.endswith('.manifest.json'))
|
||||
manifest = json.loads(writes[manifest_path])
|
||||
assert manifest['base_commit'] == ''
|
||||
# Repo identity also degrades gracefully when the headers are absent.
|
||||
assert manifest['repo_remote'] == ''
|
||||
assert manifest['branch'] == ''
|
||||
assert manifest['head_commit'] == ''
|
||||
# Prod case: delete_sandbox passes no conversation_id, so the manifest
|
||||
# falls back to the sandbox id (which is the conversation_id.hex).
|
||||
assert manifest['conversation_id'] == 'sandbox-1'
|
||||
@@ -2519,65 +2558,6 @@ class TestArchiveWorkspaceHelper:
|
||||
# tar.gz 500 is retryable + REQUIRED -> block the delete for a retry.
|
||||
assert ok is False
|
||||
|
||||
|
||||
class TestArchiveInitialWorkspaceHelper:
|
||||
"""Unit tests for the initial-state (pre-agent) workspace snapshot."""
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_initial_archive_uploads_tar_gz_and_manifest(self, monkeypatch):
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', 'true')
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_BUCKET', 'archive-bkt')
|
||||
|
||||
client = _stream_client(_make_stream_response(200, b'tar-bytes', {}))
|
||||
store = MagicMock()
|
||||
writes: dict[str, bytes] = {}
|
||||
store.write.side_effect = lambda path, data: writes.__setitem__(path, data)
|
||||
# Archive blobs are streamed from a tempfile via write_from_path (the OOM
|
||||
# fix); record them too so blob-path assertions still see the .patch/.tar.gz.
|
||||
store.write_from_path.side_effect = lambda path, src: writes.__setitem__(
|
||||
path, open(src, 'rb').read()
|
||||
)
|
||||
|
||||
with patch.object(
|
||||
workspace_archive, '_get_archive_file_store', return_value=store
|
||||
):
|
||||
ok = await workspace_archive.archive_initial_workspace(
|
||||
client,
|
||||
agent_server_url='https://sandbox.example.com',
|
||||
session_api_key='sk',
|
||||
project_dir='/workspace/project/repo',
|
||||
sandbox_id='sandbox-1',
|
||||
conversation_id='conv-1',
|
||||
base_commit='deadbeef',
|
||||
)
|
||||
|
||||
assert ok is True
|
||||
# tar.gz requested at the (already-resolved) project dir; key forwarded.
|
||||
_, kwargs = client.stream.call_args
|
||||
assert kwargs['headers']['X-Session-API-Key'] == 'sk'
|
||||
# tar.gz is the full capture, so default excludes are disabled.
|
||||
assert kwargs['params'] == {
|
||||
'path': '/workspace/project/repo',
|
||||
'format': 'tar.gz',
|
||||
'use_default_excludes': 'false',
|
||||
}
|
||||
# Keyed by conversation and nested under /initial/ so it can never collide
|
||||
# with a sibling or with this conversation's final capture (which writes to
|
||||
# {prefix}/{sandbox_id}/{conversation_key}/{ts}).
|
||||
archive_path = next(p for p in writes if p.endswith('.tar.gz'))
|
||||
assert '/sandbox-1/conv-1/initial/' in archive_path
|
||||
manifest_path = next(p for p in writes if p.endswith('.manifest.json'))
|
||||
# Shared contract: manifest = blob + '.manifest.json' (keeps the suffix).
|
||||
assert manifest_path == archive_path + '.manifest.json'
|
||||
manifest = json.loads(writes[manifest_path])
|
||||
assert manifest['phase'] == 'initial'
|
||||
assert manifest['base_commit'] == 'deadbeef'
|
||||
assert manifest['conversation_id'] == 'conv-1'
|
||||
assert manifest['format'] == 'tar.gz'
|
||||
assert manifest['source_path'] == '/workspace/project/repo'
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_archive_sibling_conversations_distinct_keys(self, monkeypatch):
|
||||
"""Two sibling conversations on the SAME grouped sandbox, captured in the
|
||||
@@ -2624,108 +2604,6 @@ class TestArchiveInitialWorkspaceHelper:
|
||||
assert any('/shared-sandbox/conva/' in p for p in patch_blobs)
|
||||
assert any('/shared-sandbox/convb/' in p for p in patch_blobs)
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_initial_archive_sibling_conversations_distinct_keys(
|
||||
self, monkeypatch
|
||||
):
|
||||
"""Same same-second path-collision guard for the initial snapshot under
|
||||
grouping: sibling initial captures land at distinct keys."""
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', 'true')
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_BUCKET', 'archive-bkt')
|
||||
fixed = datetime(2026, 1, 1, 0, 0, 0, tzinfo=timezone.utc)
|
||||
monkeypatch.setattr(workspace_archive, 'utc_now', lambda: fixed)
|
||||
|
||||
writes: dict[str, bytes] = {}
|
||||
store = MagicMock()
|
||||
store.write.side_effect = lambda path, data: writes.__setitem__(path, data)
|
||||
store.write_from_path.side_effect = lambda path, src: writes.__setitem__(
|
||||
path, open(src, 'rb').read()
|
||||
)
|
||||
|
||||
with patch.object(
|
||||
workspace_archive, '_get_archive_file_store', return_value=store
|
||||
):
|
||||
for conv in ('conva', 'convb'):
|
||||
client = _stream_client(_make_stream_response(200, b'tar-bytes', {}))
|
||||
ok = await workspace_archive.archive_initial_workspace(
|
||||
client,
|
||||
agent_server_url='https://sandbox.example.com',
|
||||
session_api_key='sk',
|
||||
project_dir=f'/workspace/project/{conv}',
|
||||
sandbox_id='shared-sandbox',
|
||||
conversation_id=conv,
|
||||
)
|
||||
assert ok is True
|
||||
|
||||
tarballs = [p for p in writes if p.endswith('.tar.gz')]
|
||||
manifests = [p for p in writes if p.endswith('.manifest.json')]
|
||||
assert len(tarballs) == 2
|
||||
assert len(manifests) == 2
|
||||
assert any('/shared-sandbox/conva/initial/' in p for p in tarballs)
|
||||
assert any('/shared-sandbox/convb/initial/' in p for p in tarballs)
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_initial_archive_disabled_is_noop(self, monkeypatch):
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
monkeypatch.delenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', raising=False)
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_BUCKET', 'archive-bkt')
|
||||
|
||||
client = AsyncMock()
|
||||
|
||||
ok = await workspace_archive.archive_initial_workspace(
|
||||
client,
|
||||
agent_server_url='https://sandbox.example.com',
|
||||
session_api_key='sk',
|
||||
project_dir='/workspace/project/repo',
|
||||
sandbox_id='sandbox-1',
|
||||
)
|
||||
# Off by default: no request, no upload.
|
||||
assert ok is False
|
||||
client.stream.assert_not_called()
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_initial_archive_non_200_returns_false(self, monkeypatch):
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', 'true')
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_BUCKET', 'archive-bkt')
|
||||
|
||||
client = _stream_client(_make_stream_response(500))
|
||||
|
||||
ok = await workspace_archive.archive_initial_workspace(
|
||||
client,
|
||||
agent_server_url='https://sandbox.example.com',
|
||||
session_api_key='sk',
|
||||
project_dir='/workspace/project/repo',
|
||||
sandbox_id='sandbox-1',
|
||||
)
|
||||
# Best-effort: a failed capture is swallowed (never blocks startup).
|
||||
assert ok is False
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_initial_archive_no_bucket_returns_false(self, monkeypatch):
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', 'true')
|
||||
monkeypatch.delenv('RUNTIME_FILE_ARCHIVE_BUCKET', raising=False)
|
||||
|
||||
client = AsyncMock()
|
||||
|
||||
ok = await workspace_archive.archive_initial_workspace(
|
||||
client,
|
||||
agent_server_url='https://sandbox.example.com',
|
||||
session_api_key='sk',
|
||||
project_dir='/workspace/project/repo',
|
||||
sandbox_id='sandbox-1',
|
||||
)
|
||||
assert ok is False
|
||||
client.stream.assert_not_called()
|
||||
|
||||
|
||||
class TestDeleteSandboxKeyHandling:
|
||||
"""The session_api_key_hash is invalidated UP FRONT on delete (a delete is
|
||||
@@ -2868,25 +2746,16 @@ class TestArchiveEnvToggles:
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_REQUIRED', value)
|
||||
assert workspace_archive.archive_required() is True
|
||||
|
||||
@pytest.mark.parametrize('value', ['1', 'true'])
|
||||
def test_initial_archive_enabled_truthy(self, monkeypatch, value):
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED', value)
|
||||
assert workspace_archive.initial_archive_enabled() is True
|
||||
|
||||
def test_toggles_default_off(self, monkeypatch):
|
||||
from openhands.app_server.sandbox import workspace_archive
|
||||
|
||||
for var in (
|
||||
'RUNTIME_FILE_ARCHIVE_ENABLED',
|
||||
'RUNTIME_FILE_ARCHIVE_REQUIRED',
|
||||
'RUNTIME_FILE_ARCHIVE_INITIAL_ENABLED',
|
||||
):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
assert workspace_archive.archive_enabled() is False
|
||||
assert workspace_archive.archive_required() is False
|
||||
assert workspace_archive.initial_archive_enabled() is False
|
||||
|
||||
|
||||
class TestArchiveRequestParams:
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
"""Test package and runtime manifest enrichment."""
|
||||
|
||||
import json
|
||||
from unittest.mock import AsyncMock, MagicMock
|
||||
|
||||
import pytest
|
||||
|
||||
from openhands.app_server.sandbox import workspace_archive as wa
|
||||
|
||||
# --- parsers ---------------------------------------------------------------
|
||||
|
||||
|
||||
def test_parse_pip_list():
|
||||
out = json.dumps(
|
||||
[
|
||||
{'name': 'requests', 'version': '2.31.0'},
|
||||
{'name': 'flask', 'version': '3.0.0'},
|
||||
]
|
||||
)
|
||||
assert wa._parse_pip_list(out) == {'requests': '2.31.0', 'flask': '3.0.0'}
|
||||
|
||||
|
||||
def test_parse_pip_list_bad_input():
|
||||
assert wa._parse_pip_list('not json') == {}
|
||||
assert wa._parse_pip_list('') == {}
|
||||
assert wa._parse_pip_list('{"unexpected": "shape"}') == {}
|
||||
assert wa._parse_pip_list('[1]') == {}
|
||||
|
||||
|
||||
def test_parse_npm_ls_top_level():
|
||||
out = json.dumps(
|
||||
{
|
||||
'dependencies': {
|
||||
'express': {'version': '4.18.2'},
|
||||
'lodash': {'version': '4.17.21'},
|
||||
}
|
||||
}
|
||||
)
|
||||
assert wa._parse_npm_ls(out) == {'express': '4.18.2', 'lodash': '4.17.21'}
|
||||
|
||||
|
||||
def test_parse_npm_ls_bad_input():
|
||||
assert wa._parse_npm_ls('') == {}
|
||||
assert wa._parse_npm_ls('not json') == {}
|
||||
assert wa._parse_npm_ls(json.dumps({'no': 'deps'})) == {}
|
||||
assert wa._parse_npm_ls(json.dumps({'dependencies': []})) == {}
|
||||
|
||||
|
||||
def test_parse_caps_entries():
|
||||
big = json.dumps(
|
||||
[
|
||||
{'name': f'p{i}', 'version': '1'}
|
||||
for i in range(wa._MAX_PACKAGES_PER_MANAGER + 50)
|
||||
]
|
||||
)
|
||||
assert len(wa._parse_pip_list(big)) == wa._MAX_PACKAGES_PER_MANAGER
|
||||
|
||||
|
||||
def test_parse_runtime_strips_prefixes():
|
||||
out = 'python=Python 3.12.4\nnode=v20.11.0\nos=ubuntu 24.04'
|
||||
assert wa._parse_runtime(out) == {
|
||||
'python': '3.12.4',
|
||||
'node': '20.11.0',
|
||||
'os': 'ubuntu 24.04',
|
||||
}
|
||||
|
||||
|
||||
def test_parse_runtime_omits_empty():
|
||||
assert wa._parse_runtime('python=Python 3.12\nnode=\nos= ') == {'python': '3.12'}
|
||||
assert '2>&1' not in wa._ENVIRONMENT_CMD
|
||||
|
||||
|
||||
def test_extract_repo_metadata_decodes_percent_encoding():
|
||||
headers = {
|
||||
'X-Archive-Repo-Remote': (
|
||||
'https%3A%2F%2Fgithub.com%2Fexample%2Ffeature%252Frepo.git'
|
||||
),
|
||||
'X-Archive-Branch': 'caf%C3%A9%25branch',
|
||||
}
|
||||
assert wa._extract_repo_metadata(headers) == {
|
||||
'repo_remote': 'https://github.com/example/feature%2Frepo.git',
|
||||
'branch': 'café%branch',
|
||||
'head_commit': '',
|
||||
}
|
||||
|
||||
|
||||
def _response(stdout: str, status_code: int = 200):
|
||||
response = MagicMock(status_code=status_code)
|
||||
response.json.return_value = {'stdout': stdout}
|
||||
return response
|
||||
|
||||
|
||||
_ENV_OUT = 'python=Python 3.12.4\nnode=v20.11.0\nos=ubuntu 24.04\n'
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_probe_workspace_full():
|
||||
pip_json = json.dumps([{'name': 'requests', 'version': '2.31.0'}])
|
||||
npm_json = json.dumps({'dependencies': {'express': {'version': '4.18.2'}}})
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(
|
||||
side_effect=[_response(pip_json), _response(npm_json), _response(_ENV_OUT)]
|
||||
)
|
||||
result = await wa._probe_workspace(
|
||||
client, 'http://host', {'X-Session-API-Key': 'key'}, '/repo'
|
||||
)
|
||||
|
||||
assert result == {
|
||||
'packages': {'pip': {'requests': '2.31.0'}, 'npm': {'express': '4.18.2'}},
|
||||
'environment': {'python': '3.12.4', 'node': '20.11.0', 'os': 'ubuntu 24.04'},
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_probe_workspace_omits_absent_tools():
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(side_effect=[_response(''), _response(''), _response('')])
|
||||
assert await wa._probe_workspace(client, 'http://host', {}, '/repo') == {}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_probe_workspace_never_raises():
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(side_effect=RuntimeError('agent-server unreachable'))
|
||||
assert await wa._probe_workspace(client, 'http://host', {}, '/repo') == {}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_run_probe_posts_command():
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(return_value=_response('done'))
|
||||
headers = {'X-Session-API-Key': 'key'}
|
||||
|
||||
assert await wa._run_probe(client, 'http://host', headers, '/repo', 'cmd') == 'done'
|
||||
client.post.assert_awaited_once_with(
|
||||
'http://host/api/bash/execute_bash_command',
|
||||
json={'command': 'cmd', 'cwd': '/repo', 'timeout': wa._PROBE_TIMEOUT},
|
||||
headers=headers,
|
||||
timeout=wa._PROBE_TIMEOUT + 1,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_probe_disabled_by_env(monkeypatch):
|
||||
monkeypatch.setenv('RUNTIME_FILE_ARCHIVE_ENRICH', 'false')
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock()
|
||||
|
||||
assert await wa._probe_workspace(client, 'h', {}, '/r') == {}
|
||||
client.post.assert_not_awaited()
|
||||
Reference in New Issue
Block a user