diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 06e8362b7..31ac0acb5 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -48,6 +48,13 @@ jobs: revision="$(git rev-parse HEAD)" if [[ -n "$REQUESTED_REF" && "$REQUESTED_REF" != "$revision" ]]; then exit 1; fi echo "revision=$revision" >> "$GITHUB_OUTPUT" + - name: Reserve disk space for distribution archives + run: | + # The hosted runner needs room for Docker images, tar exports and the + # microsandbox import. This job does not use these preinstalled SDKs. + df -h / + sudo rm -rf /usr/local/lib/android /usr/share/dotnet /usr/local/.ghcup + df -h / - uses: actions/setup-go@v7 with: go-version-file: go.mod @@ -79,8 +86,9 @@ jobs: CORE_DISTRIBUTION_OFFLINE: ${{ inputs.offline && '1' || '0' }} run: | inputs="$HOME/.parsar/build/release-inputs/inputs.json" - export AGENTS_RUNTIME_CODEX_PACKAGE="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["codex"])' "$inputs")" - export MCODE_HARNESS_BUILD_DIR="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["mcode"])' "$inputs")" + AGENTS_RUNTIME_CODEX_PACKAGE="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["codex"])' "$inputs")" + MCODE_HARNESS_BUILD_DIR="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1]))["mcode"])' "$inputs")" + export AGENTS_RUNTIME_CODEX_PACKAGE MCODE_HARNESS_BUILD_DIR export CORE_DISTRIBUTION_RELEASE_BASE_URL="https://github.com/$RELEASE_REPOSITORY/releases/download/$RELEASE_REVISION" bash scripts/build-core-distribution.sh mkdir -p "$HOME/.parsar/build/release-upload" diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 95c36bf4b..ee7c8ea96 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -1558,6 +1558,17 @@ microsandbox runtime/firmware hashes and executable native payloads. Release gen qualification. A release must be tested from fresh extraction with real models; no synthetic result may substitute for native execution acceptance. +Distribution `images` records each exported image's config digest; +`image_manifest_digests` records its OCI manifest/index digest. Derive and verify +both from the same archive, including its referenced config and layer bytes, and +require the build host's selected image ID to match one of them. Docker's classic +store identifies images by config, while its containerd store uses the OCI +descriptor. Core, node and self-hosted installers share one resolver for these +required identities: confirm Linux amd64 and the returned immutable local ID, +then use that ID in service/provider configuration and Runtime launches. Tags do +not replace identity verification. The microsandbox-qualified `runtime_ref` +remains independent of Docker's local store identity. + The manifest is the shared download contract for Core, node and self-hosted installers: flat versioned filenames, compressed Runtime size/hash and unpacked size/hash, with HTTPS release URLs or the explicit offline payload. Download into @@ -1567,6 +1578,11 @@ execution-only payloads. Python zipapps bundle the shared resolver with each remote bootstrap; the console publishes only fixed non-secret files and declared artifact names. Release automation builds artifacts and may create an unpublished draft, but cannot claim real execution qualification or public availability. +Qualify the exact downloaded production artifacts before publishing the draft; +keep the tested asset bytes and source identity unchanged. Never use an acceptance +image containing a private test CA or model credential as a release input. +Repository visibility is independent of publication. Do not add repository +credentials to installed node/Runtime configuration to bypass download access. Project-authenticated executor-credential extensions remain outside the upstream API namespace and reuse the existing restricted issuer. They require the exact @@ -1576,6 +1592,21 @@ explicit caller credentials on these routes and never substitutes its administra key. Self-hosted installation reuses Docker Runtime isolation, owns no sandbox node or Core allocation, and retains user-owned native history after uncertain launches. Report started, connected and real execution success separately. +Self-hosted installation confirms connection through the private daemon transport +using only its restricted executor credential. The read checks the exact live +Environment/key binding and current authenticated connection; it never enrolls, +allocates, wakes a sandbox or grants project resource access. Console forwarding +preserves this credential without replacing it with an administrator or project +key. Bounded polling and reruns retain the original container and history; +timeout is a diagnostic failure, not permission to relaunch. An explicit installer +`--public-url` supplies both the console origin and the advertised daemon `wss` +origin. Keep local managed Provider routing separate; do not return an internal +Compose hostname to a user-managed Runtime when an external origin was supplied. Bootstrap routing uses the +node bound to the authenticated device's persisted allocation, never request Host +or caller-supplied placement fields. An embedded managed node retains its internal +Core route; remote managed nodes use the selected setup/public route, while +self-hosted devices retain the deployment's advertised public address. This does +not widen sandbox network policies or change credential admission. The distribution build sets umask 022 for non-root-readable payloads; installation credentials and state retain their explicit private permissions. diff --git a/contracts/agents-api/environment-executor-credentials.md b/contracts/agents-api/environment-executor-credentials.md index e27d7438f..484a2e297 100644 --- a/contracts/agents-api/environment-executor-credentials.md +++ b/contracts/agents-api/environment-executor-credentials.md @@ -51,8 +51,32 @@ history volumes belong to the operator. Failed or uncertain launches retain their volumes and installation receipt for inspection instead of replacing history or retrying enrollment. Session deletion does not reclaim these volumes. Rerunning the installer inspects a previously started container only after its -installation and Environment labels match. It reports running separately from -Session connection, or gives a command to start the same stopped container. +installation and Environment labels match. Both first launch and rerun wait up +to 60 seconds for authenticated Core connection confirmation; a running container +alone does not establish connection. A stopped container receives a command to +start that same container before rerunning the installer. An uncertain launch without a success receipt gives label-filtered container and volume inspection commands and never creates a replacement. A cached image with the exact distribution digest and platform skips image download and import. + +## Private connection confirmation + +`GET /api/v1/agent-daemon/connection?environment_id=UUID` uses the existing +executor bearer, passed unchanged through the console. It is part of the private +daemon transport, not the public Agents API. It reads existing authorization and +binding only; it never enrolls a device, starts execution or changes resources. +The no-store response contains only the requested `environment_id` and `status` +(`connected` or `disconnected`). Connected requires the existing Environment +observation, its exact Session/device binding, current executor authority and a +live gateway socket authenticated with that same credential. A stale observation +or a socket carrying the former rotated key cannot confirm connection. + +Invalid, revoked, foreign or deleted-Session authority returns 401; a different +key for an already bound Environment returns 409. Responses do not expose the +actual binding or database diagnostics. The installer derives this HTTPS route +from the validated returned `remote_url`, rejects redirects, retries transient +read failures within its deadline and polls at two-second intervals. Permanent +rejections fail immediately. On timeout or rejection it retains the container, +volumes, credential and receipts, and prints bounded Docker log inspection and +same-command retry guidance. This confirms authenticated connectivity, not model +credentials, harness capabilities or completed execution. diff --git a/deploy/install/configuration.py b/deploy/install/configuration.py index 747b3ac10..ef3fe6233 100644 --- a/deploy/install/configuration.py +++ b/deploy/install/configuration.py @@ -1,5 +1,6 @@ """Deployment files for the existing Core, Runtime and production console.""" from pathlib import Path +from urllib.parse import urlsplit def bind(source, target, readonly=True): @@ -42,13 +43,18 @@ def core_environment(root, state, database_password): config = str(Path(root) / "config") if native else "/config" database = f'127.0.0.1:{state["database_port"]}' if native else "database:5432" daemon_host = f'host.microsandbox.internal:{state["core_port"]}' if native else "core:8091" + daemon_url = f"ws://{daemon_host}/api/v1/agent-daemon/ws" + if state.get("public_url"): + origin = urlsplit(state["public_url"]) + daemon_url = origin._replace(scheme="wss" if origin.scheme == "https" else "ws", + path="/api/v1/agent-daemon/ws").geturl() result = { "AGENTS_API_DATABASE_URL": f"postgres://agents_api:{database_password}@{database}/agents_api?sslmode=disable", "AGENTS_API_KEYS_FILE": config + "/keys.json", "AGENTS_API_CREDENTIAL_KEY_FILE": config + "/credential.key", "AGENTS_API_ADDR": f'127.0.0.1:{state["core_port"]}' if native else ":8091", "AGENTS_API_ENGINE": "codex", "AGENTS_API_HARNESSES": "codex,claude_sdk,mcode", - "AGENTS_API_DAEMON_WS_URL": f"ws://{daemon_host}/api/v1/agent-daemon/ws", + "AGENTS_API_DAEMON_WS_URL": daemon_url, } result["AGENTS_API_SANDBOX_ADMIN_DIGESTS_FILE"] = ( str(Path(root) / "admin/digests.json") if native else "/admin/digests.json") diff --git a/deploy/install/distribution.py b/deploy/install/distribution.py index 5dd9655d6..f64e8d5ca 100644 --- a/deploy/install/distribution.py +++ b/deploy/install/distribution.py @@ -9,6 +9,7 @@ import re import socket import stat +import subprocess import tempfile import time import urllib.error @@ -20,6 +21,55 @@ class DistributionError(Exception): pass +def image_identities(manifest, name): + """Both immutable IDs describe the same archive, as proven by the builder.""" + identities = [] + for field in ('images', 'image_manifest_digests'): + mapping = manifest.get(field) + value = mapping.get(name) if isinstance(mapping, dict) else None + if not isinstance(value, str) or not re.fullmatch(r'sha256:[0-9a-f]{64}', value): + raise DistributionError('Missing or invalid immutable image identity: ' + name) + identities.append(value) + return tuple(identities) + + +def docker_command(arguments, timeout=30): + try: + return subprocess.run(arguments, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, + stderr=subprocess.PIPE, text=True, timeout=timeout, check=False) + except (OSError, subprocess.SubprocessError): + raise DistributionError('Cannot inspect or load the distribution image; check Docker access and disk space') from None + + +def ensure_docker_image(manifest, name, archive, docker=('docker',)): + """Resolve a proven local ID; obtain a verified archive only on a cache miss.""" + expected = image_identities(manifest, name) + docker = list(docker) + + def inspect(): + for identity in dict.fromkeys(expected): + result = docker_command(docker + ['image', 'inspect', identity, '--format', + '{{.Id}} {{.Os}}/{{.Architecture}}']) + if result.returncode: + continue + fields = result.stdout.strip().split() + if len(fields) != 2 or fields[0] not in expected or fields[1] != 'linux/amd64': + raise DistributionError('Docker image identity or platform differs from the distribution: ' + name) + return fields[0] + return None + + identity = inspect() + if identity is not None: + return identity + result = docker_command(docker + ['load', '--input', str(archive())], timeout=1800) + if result.returncode: + raise DistributionError('Cannot load the distribution image; check Docker access and free disk space: ' + name) + identity = inspect() + if identity is None: + raise DistributionError('Cannot verify the loaded distribution image: ' + name) + return identity + + def digest(path): value = hashlib.sha256() with Path(path).open('rb') as stream: diff --git a/deploy/install/install.py b/deploy/install/install.py index 1a23df5e3..6347753a4 100644 --- a/deploy/install/install.py +++ b/deploy/install/install.py @@ -22,7 +22,7 @@ from configuration import compose_config, core_environment, managed_config import native_service -from distribution import DistributionError, artifact, obtain_artifact, runtime_archive +from distribution import DistributionError, artifact, obtain_artifact, runtime_archive, image_identities, ensure_docker_image class InstallError(Exception): @@ -70,9 +70,8 @@ def verify_bundle(bundle): if not required.issubset(covered): raise InstallError("Distribution checksum list is incomplete") manifest = json.loads((bundle / "manifest.json").read_text()) - for image in manifest["images"].values(): - if not image.startswith("sha256:") or len(image) != 71: - raise InstallError("Distribution must select immutable images") + for name in ("core", "web", "database", "runtime"): + image_identities(manifest, name) for name in ("images/runtime.tar.gz", "native/bin/parsar-sandbox-node", "native/bin/agents-api-microsandbox-provider", "native/microsandbox/msb", "native/microsandbox/libkrunfw.so.5.6.1"): @@ -198,6 +197,15 @@ def initialize(root, args, manifest): actual = (state["mode"], state["provider"], state["core_port"], state["web_port"], state.get("core_url"), state.get("public_url")) if wanted != actual or state["source_commit"] != manifest["source_commit"]: raise InstallError("Existing installation differs; preserve it and follow the upgrade/provider-change guide") + services = json.loads((root / "compose.json").read_text())["services"] + for service, config in services.items(): + name = "core" if service == "migrate" else service + if config.get("image") != manifest["images"].get(name): + raise InstallError("Retained Docker image differs; preserve the installation and inspect its configuration") + if state["provider"] == "docker": + managed = json.loads((root / "config/managed-runtimes.json").read_text()) + if managed.get("docker", {}).get("image") != manifest["images"]["runtime"]: + raise InstallError("Retained Runtime image differs; preserve the installation and inspect its configuration") return state if root.exists() and any(root.iterdir()): raise InstallError("Installation directory is not empty; refusing to overwrite existing state") @@ -327,30 +335,41 @@ def main(argv=None): manifest = verify_bundle(bundle) if args.provider and not args.web_only: print("Preparing the selected local sandbox provider...", flush=True) - runtime_archive(manifest, bundle, bundle) if args.provider == "microsandbox": + runtime_archive(manifest, bundle, bundle) for name in ("native/bin/agents-api-microsandbox-provider", "native/microsandbox/msb", "native/microsandbox/libkrunfw.so.5.6.1"): obtain_artifact(manifest, name, bundle / name, bundle) native_service.preflight(bundle) - state = initialize(root, args, manifest) + if args.web_only: + images = ["web"] + elif args.provider == "microsandbox": + images = ["database"] + else: + images = ["core", "database"] + if args.provider == "docker": + images.append("runtime") + if not args.core_only and not args.web_only: + images.append("web") + local_images = dict(manifest["images"]) + for name in images: + if name == "runtime": + archive = lambda: runtime_archive(manifest, bundle, bundle) + else: + archive = lambda name=name: bundle / f"images/{name}.tar" + local_images[name] = ensure_docker_image(manifest, name, archive) + # Deployment configuration uses Docker's local IDs; published metadata is unchanged. + deployment = dict(manifest, images=local_images) + state = initialize(root, args, deployment) prepare_node_payload(root, state, bundle) if state["provider"] == "docker": seccomp = bundle / "runtime/seccomp.json" if not (root / "config/seccomp.json").exists(): private_write(root / "config/seccomp.json", seccomp.read_text()) - images = ["web"] if state["mode"] == "web-only" else ["core", "database"] - if state["provider"] == "docker": - images.append("runtime") if native_service.is_native(state): - images = ["database"] password = (root / "config/database.password").read_text() environment = core_environment(root, state, password) native_service.prepare(root, state, bundle, environment) - if state["mode"] == "all": - images.append("web") - for name in images: - run(["docker", "load", "--input", str(bundle / f"images/{name}.tar")], stdout=subprocess.DEVNULL) import_runtime(root, state, manifest, bundle) compose(root, "up", "--detach", "--wait") if native_service.is_native(state): diff --git a/deploy/install/node_install.py b/deploy/install/node_install.py index 909f9b35d..6943db78d 100644 --- a/deploy/install/node_install.py +++ b/deploy/install/node_install.py @@ -124,6 +124,7 @@ def metadata(source): if (manifest.get("platform") != "linux/amd64" or not re.fullmatch(r"[0-9a-f]{40}", manifest.get("source_commit", "")) or not re.fullmatch(r"sha256:[0-9a-f]{64}", manifest.get("images", {}).get("runtime", ""))): raise InstallError("Unsupported node distribution") + distribution.image_identities(manifest, "runtime") for name in (COMMON[0], "images/runtime.tar.gz") + MICRO: distribution.artifact(manifest, name) if not manifest.get("artifact_base_url"): @@ -199,10 +200,10 @@ def micro_home(installation_id): return directory -def provider_config(root, args, manifest): +def provider_config(root, args, manifest, runtime_image): result = {"installation_id": args.installation_id, "provider": args.provider, "core_url": args.core_url + "/api/v1"} if args.provider == "docker": - result["docker"] = {"host": "unix:///var/run/docker.sock", "image": manifest["images"]["runtime"], + result["docker"] = {"host": "unix:///var/run/docker.sock", "image": runtime_image, "network": "parsar-node-" + args.installation_id, "seccomp_file": str(root / "runtime/seccomp.json"), "nested_sandbox": True} else: @@ -228,21 +229,13 @@ def provider_config(root, args, manifest): def prepare_runtime(root, args, manifest): if args.provider == "docker": docker = ["docker", "--host", "unix:///var/run/docker.sock"] - inspect = docker + ["image", "inspect", "--format", "{{.Id}}", manifest["images"]["runtime"]] - try: - image = checked(inspect, "Runtime image is not installed") - except InstallError: - image = None - if image != manifest["images"]["runtime"]: - archive = distribution.runtime_archive(manifest, root) - checked(docker + ["load", "--input", str(archive)], "Cannot import the Docker runtime image; check Docker access and free disk space", timeout=1800) - image = checked(inspect, "Cannot verify the imported runtime image") - if image != manifest["images"]["runtime"]: - raise InstallError("Imported runtime image identity differs") + image = distribution.ensure_docker_image( + manifest, "runtime", lambda: distribution.runtime_archive(manifest, root), docker) network = "parsar-node-" + args.installation_id networks = checked(docker + ["network", "ls", "--format", "{{.Name}}"], "Cannot inspect Docker networks").splitlines() if network not in networks: checked(docker + ["network", "create", network], "Cannot create node Docker network") + return image else: for name in MICRO: result = subprocess.run(["ldd", str(root / name)], stdout=subprocess.PIPE, stderr=subprocess.STDOUT, timeout=30) @@ -312,14 +305,18 @@ def install(args, token): distribution.obtain_artifact(manifest, name, target) os.chmod(target, 0o700) safe_directory(root / "state/node") + print("Checking the sandbox runtime...", flush=True) + runtime_image = prepare_runtime(root, args, manifest) # Retain the original network policy when recovering a partial installation. if not existing_file(root / "provider.json"): - write_once(root / "provider.json", json_text(provider_config(root, args, manifest))) + write_once(root / "provider.json", json_text(provider_config(root, args, manifest, runtime_image))) + elif args.provider == "docker": + stored = json.loads((root / "provider.json").read_text()) + if stored.get("docker", {}).get("image") != runtime_image: + raise InstallError("Retained Docker image differs; preserve the node and inspect its configuration") unit = root / ("parsar-node-" + args.installation_id + ".service") write_once(unit, service_unit(root)) marker = root / "registered.json" - print("Checking the sandbox runtime...", flush=True) - prepare_runtime(root, args, manifest) if not existing_file(marker): # The one-time credential is never passed through process arguments or service environments. descriptor, secret_path = tempfile.mkstemp(prefix=".enrollment-", dir=root) diff --git a/deploy/install/self_hosted_install.py b/deploy/install/self_hosted_install.py index cdb36d4de..599454e6a 100644 --- a/deploy/install/self_hosted_install.py +++ b/deploy/install/self_hosted_install.py @@ -3,18 +3,23 @@ import argparse import fcntl import getpass +import http.client import json import os from pathlib import Path import platform import re +import signal import stat import subprocess import sys -from urllib.parse import urlsplit +import time +import urllib.error +import urllib.request +from urllib.parse import urlencode, urlsplit import uuid -from distribution import DistributionError, load_manifest, obtain_artifact, runtime_archive +from distribution import DistributionError, load_manifest, obtain_artifact, runtime_archive, image_identities, ensure_docker_image class InstallError(Exception): @@ -104,6 +109,83 @@ def preflight(): 'Docker access through /var/run/docker.sock is required') +class NoRedirect(urllib.request.HTTPRedirectHandler): + def redirect_request(self, request, fp, code, message, headers, url): + raise InstallError('Core connection redirects are not supported; check the returned remote_url') + + +def open_connection(request, timeout): + return urllib.request.build_opener(NoRedirect()).open(request, timeout=timeout) + + +class _ConnectionDeadline(Exception): + pass + + +def wait_connected(remote, environment, key, container, timeout=60): + # Only the validated daemon origin receives the restricted credential. + identity(environment, remote) + address = urlsplit(remote) + endpoint = 'https://' + address.netloc + '/api/v1/agent-daemon/connection?' + request = urllib.request.Request(endpoint + urlencode({'environment_id': environment}), + headers={'Authorization': 'Bearer ' + key['executor_token']}) + guidance = (' Inspect with: docker --host unix:///var/run/docker.sock logs --tail 100 ' + container + + '. Check the Runtime network, TLS and executor credential, then rerun the same installation command.' + + ' Keep the existing container, volumes and installation state; do not replace history.') + started_at = time.monotonic() + deadline = started_at + timeout + detail = 'Core has not confirmed this Environment connection' + + def deadline_expired(_signal, _frame): + raise _ConnectionDeadline() + + # Socket timeouts only limit inactivity. The Linux CLI needs a process timer + # as well so a slow response cannot keep the overall deadline alive. + previous_handler = signal.signal(signal.SIGALRM, deadline_expired) + previous_timer = signal.setitimer(signal.ITIMER_REAL, max(0.001, timeout)) + try: + while time.monotonic() < deadline: + try: + with open_connection(request, min(10, max(0.1, deadline - time.monotonic()))) as response: + raw = response.read(4097) + if len(raw) > 4096: + raise ValueError() + result = json.loads(raw) + if (not isinstance(result, dict) or set(result) != {'environment_id', 'status'} + or result['environment_id'] != environment + or result['status'] not in ('connected', 'disconnected')): + raise ValueError() + if time.monotonic() >= deadline: + raise _ConnectionDeadline() + if result['status'] == 'connected': + print('Runtime connected to Environment ' + environment + ': ' + container) + return + detail = 'Core reports this Environment disconnected' + except urllib.error.HTTPError as error: + if error.code not in (408, 429, 500, 502, 503, 504): + raise InstallError('Core connection check rejected (HTTP ' + str(error.code) + + '); verify the exact Environment and active executor key.' + guidance) from None + detail = 'Core connection check unavailable (HTTP ' + str(error.code) + ')' + except (urllib.error.URLError, TimeoutError, ConnectionError, http.client.IncompleteRead): + detail = 'Cannot reach the Core connection endpoint; check DNS, TLS and network access' + except (ValueError, TypeError, UnicodeError): + raise InstallError('Core returned an invalid connection response.' + guidance) from None + except InstallError as error: + raise InstallError(str(error) + guidance) from None + remaining = deadline - time.monotonic() + if remaining > 0: + time.sleep(min(2, remaining)) + except _ConnectionDeadline: + pass + finally: + signal.setitimer(signal.ITIMER_REAL, 0) + signal.signal(signal.SIGALRM, previous_handler) + if previous_timer[0] > 0: + remaining_timer = max(0.001, previous_timer[0] - (time.monotonic() - started_at)) + signal.setitimer(signal.ITIMER_REAL, remaining_timer, previous_timer[1]) + raise InstallError('Runtime connection timed out: ' + detail + '.' + guidance) + + def inspect_prior_launch(root, state): docker = 'docker --host unix:///var/run/docker.sock' filters = (' --filter label=io.parsar.agents-api.installation=' + state['installation_id'] @@ -120,8 +202,8 @@ def inspect_prior_launch(root, state): raise InstallError('Invalid retained Runtime startup receipt.' + guidance) try: raw = checked(['docker', '--host', 'unix:///var/run/docker.sock', 'container', 'inspect', name, - '--format', '{{json .Config.Labels}} {{.State.Status}}'], 'Cannot inspect retained Runtime') - labels, status = raw.strip().rsplit(' ', 1) + '--format', '{{json .Config.Labels}} {{.State.Status}} {{.Image}}'], 'Cannot inspect retained Runtime') + labels, status, image = raw.strip().rsplit(' ', 2) labels = json.loads(labels) except (InstallError, ValueError): raise InstallError('The prior Runtime cannot be confirmed.' + guidance) from None @@ -130,13 +212,14 @@ def inspect_prior_launch(root, state): 'io.parsar.agents-api.user-owned': 'true'} if not isinstance(labels, dict) or any(labels.get(key) != value for key, value in expected.items()): raise InstallError('The retained container does not match this installation and Environment.' + guidance) + if image not in (state['runtime_image'], state['runtime_manifest']): + raise InstallError('The retained container image does not match this distribution.' + guidance) if status == 'running': print('Runtime already running: ' + name) - print('Read the Session to confirm connection; a running container does not establish it.') - return + return name if status == 'exited': raise InstallError('The existing Runtime is stopped. Preserve its history and resume that same container with: ' - + docker + ' start ' + name + '. Then read the Session to confirm connection.') + + docker + ' start ' + name + '. Then rerun the same installation command to confirm connection.') raise InstallError('The retained Runtime requires inspection before continuing.' + guidance) @@ -144,11 +227,11 @@ def install(args, root): target = identity(args.environment_id, args.remote) manifest = load_manifest(source_url=args.source_url, offline_root=args.offline_root) revision = manifest.get('source_commit', '') - runtime_image = manifest.get('images', {}).get('runtime', '') + runtime_image, runtime_manifest = image_identities(manifest, 'runtime') if (manifest.get('platform') != 'linux/amd64' or not re.fullmatch(r'[0-9a-f]{40}', revision) or not re.fullmatch(r'sha256:[0-9a-f]{64}', runtime_image)): raise InstallError('The distribution does not contain a matched Linux amd64 Runtime') - target.update(source_commit=revision, runtime_image=runtime_image) + target.update(source_commit=revision, runtime_image=runtime_image, runtime_manifest=runtime_manifest) state_file = root / 'installation.json' if state_file.exists(): state = json.loads(private_read(state_file)) @@ -160,7 +243,11 @@ def install(args, root): state = dict(target, installation_id=str(uuid.uuid4())) write_private(state_file, state) if (root / 'launch.json').exists() or (root / 'launch.json').is_symlink(): - inspect_prior_launch(root, state) + name = inspect_prior_launch(root, state) + key = credential(private_read(root / 'executor-key.json'), args.environment_id) + if args.credential_file and credential(private_read(Path(args.credential_file)), args.environment_id) != key: + raise InstallError('Stored executor credential differs; inspect the existing installation') + wait_connected(args.remote, args.environment_id, key, name) return key_file = root / 'executor-key.json' if key_file.exists(): @@ -175,19 +262,9 @@ def install(args, root): write_private(key_file, key) launcher = obtain_artifact(manifest, 'native/bin/parsar-runtime', root / 'native/bin/parsar-runtime', args.offline_root) seccomp = obtain_artifact(manifest, 'runtime/seccomp.json', root / 'runtime/seccomp.json', args.offline_root) - inspect = ['docker', '--host', 'unix:///var/run/docker.sock', 'image', 'inspect', runtime_image, - '--format', '{{.Id}} {{.Os}}/{{.Architecture}}'] - try: - image = checked(inspect, 'Runtime image is not installed').strip() - except InstallError: - image = None - if image != runtime_image + ' linux/amd64': - archive = runtime_archive(manifest, root, args.offline_root) - checked(['docker', '--host', 'unix:///var/run/docker.sock', 'image', 'load', '--input', str(archive)], - 'Cannot load the matched Runtime image; retry after checking Docker', timeout=600) - image = checked(inspect, 'Cannot verify the loaded Runtime image').strip() - if image != runtime_image + ' linux/amd64': - raise InstallError('Loaded Runtime image does not match the distribution') + runtime_image = ensure_docker_image( + manifest, 'runtime', lambda: runtime_archive(manifest, root, args.offline_root), + ('docker', '--host', 'unix:///var/run/docker.sock')) command = [str(launcher), '--installation-id', state['installation_id'], '--environment-id', args.environment_id, '--remote', args.remote, '--image', runtime_image, '--seccomp-file', str(seccomp), '--credential-file', str(key_file)] @@ -200,7 +277,8 @@ def install(args, root): raise InstallError('Runtime launcher returned an invalid result; inspect the retained container') write_private(root / 'started.json', result) print('Runtime started: ' + result['container']) - print('Read the Session to confirm connection. Stop this user-owned Runtime with: docker stop ' + result['container']) + wait_connected(args.remote, args.environment_id, key, result['container']) + print('Stop this user-owned Runtime with: docker --host unix:///var/run/docker.sock stop ' + result['container']) def main(): diff --git a/deploy/install/test_distribution.py b/deploy/install/test_distribution.py index 6d785d476..a55cd3cea 100644 --- a/deploy/install/test_distribution.py +++ b/deploy/install/test_distribution.py @@ -6,8 +6,9 @@ from pathlib import Path import tempfile import threading +from types import SimpleNamespace import unittest -from unittest.mock import patch +from unittest.mock import Mock, patch import distribution @@ -119,5 +120,68 @@ def test_url_and_path_boundaries(self): distribution.obtain_artifact(self.manifest, 'native/bin/node', self.root / 'link') +class DockerIdentityTests(unittest.TestCase): + def setUp(self): + self.config = 'sha256:' + 'a' * 64 + self.oci = 'sha256:' + 'b' * 64 + self.manifest = {'images': {'runtime': self.config}, + 'image_manifest_digests': {'runtime': self.oci}} + self.archive = Mock(return_value=Path('/verified/runtime.tar')) + + def result(self, identity=None, platform='linux/amd64', code=0): + return SimpleNamespace(returncode=code, stdout=(identity + ' ' + platform) if identity else '') + + def ensure(self): + return distribution.ensure_docker_image(self.manifest, 'runtime', self.archive) + + def test_both_stores_use_proven_immutable_cache_without_archive(self): + for responses, identity in (([self.result(self.config)], self.config), + ([self.result(code=1), self.result(self.oci)], self.oci), + ([self.result(self.oci)], self.oci)): + with self.subTest(identity=identity), patch.object(distribution, 'docker_command', side_effect=responses) as command: + self.assertEqual(self.ensure(), identity) + self.assertTrue(all('load' not in call.args[0] for call in command.call_args_list)) + self.archive.assert_not_called() + + def test_load_is_verified_by_either_digest_on_both_stores(self): + for identity in (self.config, self.oci): + replies = [self.result(code=1), self.result(code=1), self.result()] + if identity == self.oci: + replies.append(self.result(code=1)) + replies.append(self.result(identity)) + with self.subTest(identity=identity), patch.object(distribution, 'docker_command', side_effect=replies) as command: + self.assertEqual(self.ensure(), identity) + self.assertEqual(command.call_args_list[2].args[0], ['docker', 'load', '--input', '/verified/runtime.tar']) + self.assertEqual(self.archive.call_count, 2) + + def test_wrong_id_platform_or_malformed_inspection_is_never_trusted(self): + invalid = [self.result('sha256:' + 'c' * 64), self.result(self.oci, 'linux/arm64'), + self.result(self.config, 'windows/amd64'), self.result()] + for reply in invalid: + for loaded in (False, True): + responses = ([self.result(code=1), self.result(code=1), self.result()] if loaded else []) + [reply] + with self.subTest(reply=reply, loaded=loaded), patch.object(distribution, 'docker_command', side_effect=responses), \ + self.assertRaisesRegex(distribution.DistributionError, 'identity or platform'): + self.ensure() + + def test_successful_load_without_inspectable_identity_fails(self): + with patch.object(distribution, 'docker_command', side_effect=[self.result(code=1), self.result(code=1), + self.result(), self.result(code=1), self.result(code=1)]), \ + self.assertRaisesRegex(distribution.DistributionError, 'Cannot verify'): + self.ensure() + + def test_both_metadata_identities_are_required_before_docker_or_archive(self): + for field in ('images', 'image_manifest_digests'): + original = self.manifest[field] + for invalid in (None, {}, {'runtime': 'mutable:tag'}, {'runtime': 'sha256:' + 'g' * 64}): + self.manifest[field] = invalid + with self.subTest(field=field, invalid=invalid), patch.object(distribution, 'docker_command') as command, \ + self.assertRaisesRegex(distribution.DistributionError, 'immutable image identity'): + self.ensure() + command.assert_not_called() + self.archive.assert_not_called() + self.manifest[field] = original + + if __name__ == '__main__': unittest.main() diff --git a/deploy/install/test_install.py b/deploy/install/test_install.py index 9370b3a89..8cb415633 100644 --- a/deploy/install/test_install.py +++ b/deploy/install/test_install.py @@ -20,6 +20,7 @@ from unittest import mock import install +import distribution class InstallerTests(unittest.TestCase): @@ -34,9 +35,15 @@ def setUp(self): "source_commit": "a" * 40, "images": {name: "sha256:" + digit * 64 for name, digit in ( ("core", "1"), ("runtime", "2"), ("database", "3"), ("web", "4"))}, + "image_manifest_digests": {name: "sha256:" + digit * 64 for name, digit in ( + ("core", "a"), ("runtime", "b"), ("database", "c"), ("web", "d"))}, "runtime_ref": "localhost/parsar-runtime:test-install", "microsandbox": {"runtime_sha256": "5" * 64, "firmware_sha256": "6" * 64}, } + self.loaded_images = set() + self.containerd = False + self.invalid_image = None + self.patched(mock.patch.object(distribution, "docker_command", side_effect=self.docker_command)) self.ports = self.patched(mock.patch.object(install, "free_port")) self.device_probes = [] original_stat = os.stat @@ -51,6 +58,19 @@ def controlled_device_stat(name, *args, **kwargs): self.patched(mock.patch.object(install.os, "getuid", return_value=1000)) self.patched(mock.patch.object(install.os, "getgid", return_value=1000)) + def docker_command(self, arguments, **kwargs): + if 'load' in arguments: + self.loaded_images.add(Path(arguments[-1]).stem) + install.run(arguments) + return SimpleNamespace(returncode=0, stdout='') + identity = arguments[3] + name = next(name for name in self.manifest['images'] if identity in distribution.image_identities(self.manifest, name)) + if name not in self.loaded_images: + return SimpleNamespace(returncode=1, stdout='') + if self.containerd and identity == self.manifest['images'][name]: + return SimpleNamespace(returncode=1, stdout='') + return SimpleNamespace(returncode=0, stdout=self.invalid_image or identity + ' linux/amd64') + def patched(self, patcher): result = patcher.start() self.addCleanup(patcher.stop) @@ -173,6 +193,20 @@ def test_repeat_installation_preserves_execution_identity_and_all_secrets(self): self.assertEqual(before, self.snapshot()) self.assertEqual(keys, self.document("config/keys.json")) + def test_retained_config_cannot_launch_an_unresolved_image(self): + self.initialize('--sandbox-provider', 'true', '--provider', 'docker') + for path, select in (('compose.json', lambda value: value['services']['web']), + ('config/managed-runtimes.json', lambda value: value['docker'])): + original = (self.root / path).read_bytes() + config = self.document(path) + select(config)['image'] = 'sha256:' + 'f' * 64 + (self.root / path).write_text(json.dumps(config)) + before = self.snapshot() + with self.subTest(path=path), self.assertRaisesRegex(install.InstallError, 'image differs'): + self.initialize('--sandbox-provider', 'true', '--provider', 'docker') + self.assertEqual(before, self.snapshot()) + (self.root / path).write_bytes(original) + def test_configuration_changes_refuse_without_mutating_existing_deployment(self): self.initialize() before = self.snapshot() @@ -415,9 +449,35 @@ def test_public_origin_is_explicit_and_preserved(self): self.assertEqual(web["environment"]["CORE_CONSOLE_ORIGIN"], "https://core.example") self.assertEqual(web["ports"], ["127.0.0.1:8080:8080"]) self.assertEqual(state["public_url"], "https://core.example") + self.assertEqual(self.document("compose.json")["services"]["core"]["environment"]["AGENTS_API_DAEMON_WS_URL"], + "wss://core.example/api/v1/agent-daemon/ws") with self.assertRaises(install.InstallError): self.initialize("--public-url", "https://other.example") + def test_public_daemon_address_is_shared_across_placement_modes(self): + state = self.initialize("--public-url", "https://core.example:8443") + for provider in (None, "docker", "microsandbox"): + with self.subTest(provider=provider): + configured = dict(state, provider=provider, database_port=15432) + env = install.core_environment(self.root, configured, "fixture-password") + self.assertEqual(env["AGENTS_API_DAEMON_WS_URL"], + "wss://core.example:8443/api/v1/agent-daemon/ws") + configured["public_url"] = None + local = install.core_environment(self.root, configured, "fixture-password") + expected_host = "host.microsandbox.internal:8091" if provider == "microsandbox" else "core:8091" + self.assertEqual(local["AGENTS_API_DAEMON_WS_URL"], "ws://" + expected_host + "/api/v1/agent-daemon/ws") + + def test_accepted_public_origin_schemes_generate_websocket_urls(self): + state = self.initialize() + for origin, expected in (("HTTPS://core.example", "wss://core.example"), + ("http://localhost:8080", "ws://localhost:8080"), + ("http://127.0.0.1:8080", "ws://127.0.0.1:8080")): + with self.subTest(origin=origin): + args = self.args("--public-url", origin) + configured = dict(state, public_url=args.public_url) + env = install.core_environment(self.root, configured, "fixture-password") + self.assertEqual(env["AGENTS_API_DAEMON_WS_URL"], expected + "/api/v1/agent-daemon/ws") + def test_provider_requires_explicit_enablement_and_cannot_belong_to_web_only(self): self.assertIsNone(self.args("--sandbox-provider", "false").provider) self.assertEqual(self.args("--sandbox-provider").provider, "microsandbox") @@ -518,6 +578,53 @@ def external_command(args, **_kwargs): for path in (self.root / "config").iterdir(): self.assertNotIn(path.read_text(), output.getvalue()) + def test_containerd_core_configuration_uses_local_ids_and_keeps_published_manifest(self): + self.containerd = True + self.loaded_images.update(self.manifest['images']) + for provider in (None, 'docker'): + with self.subTest(provider=provider): + bundle = self.bundle() + (bundle / 'images/runtime.tar').unlink() + self.write_checksums(bundle) + published = (bundle / 'manifest.json').read_bytes() + with mock.patch.object(install, '__file__', str(bundle / 'install.py')), \ + mock.patch.object(install.platform, 'system', return_value='Linux'), \ + mock.patch.object(install.platform, 'machine', return_value='x86_64'), \ + mock.patch.object(install, 'run'), mock.patch.object(install, 'wait_http', return_value=True), \ + mock.patch.object(install, 'runtime_archive', side_effect=AssertionError('cache must avoid Runtime download')), \ + contextlib.redirect_stdout(io.StringIO()): + flags = ['--sandbox-provider', 'true', '--provider', provider] if provider else [] + install.main(['--install-dir', str(self.root), *flags]) + before = self.snapshot() + install.main(['--install-dir', str(self.root), *flags]) + self.assertEqual(self.snapshot(), before) + for service, config in self.document('compose.json')['services'].items(): + self.assertEqual(config['image'], self.manifest['image_manifest_digests']['core' if service == 'migrate' else service]) + if provider: + self.assertEqual(self.document('config/managed-runtimes.json')['docker']['image'], + self.manifest['image_manifest_digests']['runtime']) + else: + self.assertFalse((self.root / 'state').exists()) + self.assertEqual((self.root / 'node-payload/manifest.json').read_bytes(), published) + self.assertEqual((bundle / 'manifest.json').read_bytes(), published) + shutil.rmtree(bundle) + shutil.rmtree(self.root) + + def test_postload_wrong_identity_or_platform_cannot_create_deployment(self): + bundle = self.bundle() + for observed in ('sha256:' + 'f' * 64 + ' linux/amd64', self.manifest['images']['core'] + ' linux/arm64'): + self.invalid_image = observed + self.loaded_images.clear() + with self.subTest(observed=observed), \ + mock.patch.object(install, '__file__', str(bundle / 'install.py')), \ + mock.patch.object(install.platform, 'system', return_value='Linux'), \ + mock.patch.object(install.platform, 'machine', return_value='x86_64'), \ + mock.patch.object(install, 'run') as command, \ + self.assertRaisesRegex(distribution.DistributionError, 'identity or platform'): + install.main(['--install-dir', str(self.root)]) + self.assertFalse(self.root.exists()) + self.assertFalse(any('up' in call.args[0] for call in command.call_args_list)) + def test_cli_failure_does_not_print_external_command_secrets(self): secret = "synthetic-sensitive-command-value" failure = subprocess.CalledProcessError(1, ["docker", secret], output=secret, stderr=secret) diff --git a/deploy/install/test_node_install.py b/deploy/install/test_node_install.py index bb8aaee87..b88752321 100644 --- a/deploy/install/test_node_install.py +++ b/deploy/install/test_node_install.py @@ -27,11 +27,14 @@ def setUp(self): provider="docker", installation_id="94be54a1-138c-4f30-bc87-b13686272dbe") self.root = self.home / ".parsar/nodes" / self.args.installation_id self.manifest = {"platform": "linux/amd64", "source_commit": "a" * 40, "images": {"runtime": "sha256:" + "b" * 64}, + "image_manifest_digests": {"runtime": "sha256:" + "c" * 64}, "runtime_ref": "parsar-core-runtime@sha256:" + "c" * 64, "microsandbox": {"runtime_sha256": "d" * 64, "firmware_sha256": "e" * 64}} self.payloads = {name: b"fixture-payload-" + name.encode() for name in installer.COMMON + installer.MICRO} self.payloads["images/runtime.tar.gz"] = gzip.compress(b"runtime archive") self.refresh_manifest() + self.containerd = False + self.invalid_image = None self.image_present = False self.calls = [] self.fail_service = False @@ -43,6 +46,7 @@ def setUp(self): mock.patch.object(installer, "micro_home", return_value=self.home / "m"), mock.patch.object(installer, "fetch", side_effect=lambda source, name: io.BytesIO(self.payloads[name])), mock.patch.object(installer, "checked", side_effect=self.checked), + mock.patch.object(installer.distribution, "docker_command", side_effect=self.docker_command), mock.patch.object(installer.subprocess, "run", return_value=subprocess.CompletedProcess([], 0, b"statically linked", b""))): patch.start() self.addCleanup(patch.stop) @@ -81,10 +85,19 @@ def checked(self, arguments, failure, **kwargs): if not self.image_present: raise installer.InstallError(failure) if arguments[0] == "docker": - return self.manifest["images"]["runtime"] + if self.containerd and arguments[arguments.index("inspect") + 1] == self.manifest["images"]["runtime"]: + raise installer.InstallError(failure) + identity = self.manifest["image_manifest_digests" if self.containerd else "images"]["runtime"] + return self.invalid_image or identity + " linux/amd64" return json.dumps({"digest": self.manifest["runtime_ref"].split("@", 1)[1], "os": "linux", "architecture": "amd64"}) return "" + def docker_command(self, arguments, **kwargs): + try: + return subprocess.CompletedProcess(arguments, 0, self.checked(arguments, "image missing", **kwargs), "") + except installer.InstallError: + return subprocess.CompletedProcess(arguments, 1, "", "") + def install(self): installer.install(self.args, "synthetic-once-token") @@ -103,6 +116,35 @@ def test_docker_installs_matched_payload_registers_and_starts_persistent_service self.assertTrue(any("register" in call for call, _ in self.calls)) self.assertTrue(any("is-active" in call for call, _ in self.calls)) + def test_containerd_node_persists_actual_id_and_warm_retry_avoids_archive(self): + self.containerd = True + self.install() + expected = self.manifest['image_manifest_digests']['runtime'] + self.assertEqual(json.loads((self.root / 'provider.json').read_text())['docker']['image'], expected) + before = (self.root / 'provider.json').read_bytes() + with mock.patch.object(installer.distribution, 'runtime_archive', side_effect=AssertionError('warm Runtime download')): + self.install() + self.assertEqual((self.root / 'provider.json').read_bytes(), before) + + def test_retained_provider_image_cannot_bypass_verified_selection(self): + self.install() + config = json.loads((self.root / 'provider.json').read_text()) + config['docker']['image'] = 'sha256:' + 'f' * 64 + (self.root / 'provider.json').write_text(json.dumps(config)) + self.calls.clear() + with self.assertRaisesRegex(installer.InstallError, 'Retained Docker image differs'): + self.install() + self.assertFalse(any('register' in call or 'enable' in call for call, _ in self.calls)) + + def test_wrong_loaded_runtime_cannot_register_or_write_provider_config(self): + for observed in ('sha256:' + 'f' * 64 + ' linux/amd64', self.manifest['images']['runtime'] + ' linux/arm64'): + self.invalid_image = observed + self.image_present = False + with self.subTest(observed=observed), self.assertRaisesRegex(installer.distribution.DistributionError, 'identity or platform'): + self.install() + self.assertFalse((self.root / 'provider.json').exists()) + self.assertFalse(any('register' in call or 'enable' in call for call, _ in self.calls)) + def test_microsandbox_imports_image_and_allows_only_explicit_private_core_endpoint(self): self.args.provider = "microsandbox" self.install() diff --git a/deploy/install/test_self_hosted_connection.py b/deploy/install/test_self_hosted_connection.py new file mode 100644 index 000000000..9fba7e140 --- /dev/null +++ b/deploy/install/test_self_hosted_connection.py @@ -0,0 +1,135 @@ +"""Core confirmation is independent of Docker launch and never recreates history.""" +import contextlib +import io +import http.server +import json +import signal +import threading +import time +import unittest +from unittest.mock import patch +import urllib.error +import urllib.request + +import self_hosted_install as installer + + +class ConnectionTests(unittest.TestCase): + environment = '55b5311c-4df9-43ce-b875-faf901e6d10f' + remote = 'wss://core.example:8443/api/v1/agent-daemon/ws' + key = {'executor_token': 'synthetic-private-token'} + container = 'parsar-selfhost-' + 'c' * 32 + + def response(self, status='connected', environment=None): + return io.BytesIO(json.dumps({'status': status, 'environment_id': environment or self.environment}).encode()) + + def wait(self): + installer.wait_connected(self.remote, self.environment, self.key, self.container) + + def test_waits_for_core_and_recovers_transient_failure(self): + failure = urllib.error.HTTPError('url', 503, 'private body', {}, None) + with patch.object(installer, 'open_connection', side_effect=[failure, self.response('disconnected'), self.response()]) as request, \ + patch.object(installer.time, 'sleep') as sleep, contextlib.redirect_stdout(io.StringIO()) as output: + self.wait() + self.assertEqual(request.call_count, 3) + self.assertEqual(sleep.call_count, 2) + sent = request.call_args.args[0] + self.assertEqual(sent.full_url, 'https://core.example:8443/api/v1/agent-daemon/connection?environment_id=' + self.environment) + self.assertEqual(sent.get_header('Authorization'), 'Bearer ' + self.key['executor_token']) + self.assertIn('Runtime connected to Environment', output.getvalue()) + self.assertNotIn(self.key['executor_token'], output.getvalue()) + + def test_permanent_denial_does_not_retry_or_echo_server_body(self): + for status in (401, 403, 404, 409): + with patch.object(installer, 'open_connection', side_effect=urllib.error.HTTPError('url', status, 'private body', {}, None)) as request: + with self.assertRaisesRegex(installer.InstallError, 'HTTP ' + str(status)) as failure: + self.wait() + self.assertEqual(request.call_count, 1) + self.assertIn('logs --tail 100 ' + self.container, str(failure.exception)) + self.assertNotIn('private body', str(failure.exception)) + + def test_deadline_retains_container_and_has_retry_guidance(self): + with patch.object(installer, 'open_connection', return_value=self.response('disconnected')), \ + patch.object(installer.time, 'monotonic', side_effect=[0, 0, 0, 0, 60, 60]), \ + patch.object(installer, 'checked') as mutation: + with self.assertRaisesRegex(installer.InstallError, 'timed out.*disconnected') as failure: + self.wait() + mutation.assert_not_called() + self.assertIn('rerun the same installation command', str(failure.exception)) + self.assertIn('do not replace history', str(failure.exception)) + + def test_wrong_or_malformed_response_cannot_confirm_connection(self): + for response in (self.response(environment='other'), self.response(status='running'), io.BytesIO(b'{}'), io.BytesIO(b'x' * 4097)): + with patch.object(installer, 'open_connection', return_value=response): + with self.assertRaisesRegex(installer.InstallError, 'invalid connection response'): + self.wait() + + def test_slow_response_cannot_outlive_deadline_or_report_success(self): + body = self.response().getvalue() + class SlowResponse(http.server.BaseHTTPRequestHandler): + def do_GET(self): + self.send_response(200) + self.send_header('Content-Length', str(len(body))) + self.end_headers() + try: + for offset in range(0, len(body), 10): + self.wfile.write(body[offset:offset + 10]) + self.wfile.flush() + time.sleep(0.07) + except (BrokenPipeError, ConnectionResetError): + pass + + def log_message(self, *_args): + pass + + server = http.server.HTTPServer(('127.0.0.1', 0), SlowResponse) + worker = threading.Thread(target=server.serve_forever, daemon=True) + worker.start() + previous_handler = signal.getsignal(signal.SIGALRM) + try: + def open_local(request, timeout): + local = urllib.request.Request('http://127.0.0.1:' + str(server.server_port), + headers=dict(request.header_items())) + return urllib.request.urlopen(local, timeout=timeout) + with patch.object(installer, 'open_connection', side_effect=open_local), \ + contextlib.redirect_stdout(io.StringIO()) as output: + started = time.monotonic() + with self.assertRaisesRegex(installer.InstallError, 'connection timed out'): + installer.wait_connected(self.remote, self.environment, self.key, self.container, timeout=0.2) + elapsed = time.monotonic() - started + self.assertLess(elapsed, 0.5) + self.assertNotIn('Runtime connected', output.getvalue()) + self.assertEqual(signal.getitimer(signal.ITIMER_REAL), (0, 0)) + self.assertEqual(signal.getsignal(signal.SIGALRM), previous_handler) + finally: + server.shutdown() + server.server_close() + worker.join(timeout=2) + + def test_timer_and_handler_are_restored_on_success_and_rejection(self): + previous_handler = signal.getsignal(signal.SIGALRM) + try: + signal.setitimer(signal.ITIMER_REAL, 30) + for response in (self.response(), io.BytesIO(b'{}')): + with patch.object(installer, 'open_connection', return_value=response), \ + contextlib.redirect_stdout(io.StringIO()): + try: + self.wait() + except installer.InstallError: + pass + remaining, interval = signal.getitimer(signal.ITIMER_REAL) + self.assertGreater(remaining, 20) + self.assertLessEqual(remaining, 30) + self.assertEqual(interval, 0) + self.assertEqual(signal.getsignal(signal.SIGALRM), previous_handler) + finally: + signal.setitimer(signal.ITIMER_REAL, 0) + + def test_redirect_never_forwards_bearer(self): + for target in ('https://other.example/connection', 'https://core.example/new'): + with self.assertRaisesRegex(installer.InstallError, 'redirects'): + installer.NoRedirect().redirect_request(None, None, 302, '', {}, target) + + +if __name__ == '__main__': + unittest.main() diff --git a/deploy/install/test_self_hosted_install.py b/deploy/install/test_self_hosted_install.py index 7c233fc59..c44ea5d2a 100644 --- a/deploy/install/test_self_hosted_install.py +++ b/deploy/install/test_self_hosted_install.py @@ -5,6 +5,9 @@ import os from pathlib import Path import tempfile +from types import SimpleNamespace + +import distribution import unittest from unittest.mock import patch import uuid @@ -14,7 +17,9 @@ class SelfHostedInstallTests(unittest.TestCase): def setUp(self): - self.temporary = tempfile.TemporaryDirectory() + base = Path.home() / '.parsar/tests/selfhost-install' + base.mkdir(parents=True, exist_ok=True) + self.temporary = tempfile.TemporaryDirectory(dir=base) self.addCleanup(self.temporary.cleanup) self.root = Path(self.temporary.name).resolve() self.environment = str(uuid.uuid4()) @@ -25,8 +30,18 @@ def setUp(self): installer.write_private(self.source_key, self.key) self.args = argparse.Namespace(source_url='https://core.example', offline_root=None, environment_id=self.environment, remote=self.remote, credential_file=str(self.source_key)) + self.wait = patch.object(installer, 'wait_connected').start() + self.addCleanup(patch.stopall) self.manifest = {'source_commit': 'a' * 40, 'platform': 'linux/amd64', - 'images': {'runtime': 'sha256:' + 'b' * 64}} + 'images': {'runtime': 'sha256:' + 'b' * 64}, + 'image_manifest_digests': {'runtime': 'sha256:' + 'd' * 64}} + patch.object(distribution, 'docker_command', side_effect=self.docker_command).start() + + def docker_command(self, arguments, timeout=30): + try: + return SimpleNamespace(returncode=0, stdout=installer.checked(arguments, 'image missing', timeout=timeout)) + except installer.InstallError: + return SimpleNamespace(returncode=1, stdout='') def test_exact_restricted_credential_and_tls_identity(self): self.assertEqual(installer.credential(json.dumps(self.key), self.environment), self.key) @@ -70,6 +85,32 @@ def checked(command, message, timeout=30): with self.assertRaises(installer.InstallError): installer.install(self.args, self.root) self.assertEqual(len(commands), 3) + def test_connection_failure_retry_preserves_same_container_and_credential(self): + commands = [] + name = 'parsar-selfhost-'+'c'*32 + def checked(command, message, timeout=30): + commands.append(command) + if 'container' in command: + state = json.loads(installer.private_read(self.root/'installation.json')) + labels = {'io.parsar.agents-api.installation': state['installation_id'], + 'io.parsar.agents-api.environment': self.environment, + 'io.parsar.agents-api.user-owned': 'true'} + return json.dumps(labels)+' running '+state['runtime_image'] + if 'inspect' in command: return self.manifest['images']['runtime'] + ' linux/amd64' + return json.dumps({'container': name, 'status': 'started'}) + self.wait.side_effect = [installer.InstallError('connection timed out'), None] + with patch.object(installer, 'load_manifest', return_value=self.manifest), \ + patch.object(installer, 'obtain_artifact', side_effect=lambda m, n, dest, o: dest), \ + patch.object(installer, 'checked', side_effect=checked), contextlib.redirect_stdout(io.StringIO()): + with self.assertRaisesRegex(installer.InstallError, 'timed out'): + installer.install(self.args, self.root) + original = {name: (self.root/name).read_bytes() for name in ('installation.json', 'launch.json', 'started.json', 'executor-key.json')} + installer.install(self.args, self.root) + self.assertEqual(len([c for c in commands if c[0].endswith('parsar-runtime')]), 1) + self.assertEqual(self.wait.call_count, 2) + self.wait.assert_called_with(self.remote, self.environment, self.key, name) + self.assertEqual(original, {name: (self.root/name).read_bytes() for name in original}) + def test_cached_runtime_avoids_archive_download_and_load(self): def checked(command, message, timeout=30): self.assertNotIn('load', command) @@ -86,7 +127,8 @@ def test_missing_runtime_loads_and_verifies_before_launch(self): commands = [] def checked(command, message, timeout=30): commands.append(command) - if len(commands) == 1: raise installer.InstallError('image missing') + if 'inspect' in command and not any('load' in item for item in commands): + raise installer.InstallError('image missing') if 'inspect' in command: return self.manifest['images']['runtime'] + ' linux/amd64' if command[0].endswith('parsar-runtime'): return json.dumps({'container': 'parsar-selfhost-'+'c'*32, 'status': 'started'}) @@ -97,32 +139,83 @@ def checked(command, message, timeout=30): patch.object(installer, 'checked', side_effect=checked), contextlib.redirect_stdout(io.StringIO()): installer.install(self.args, self.root) archive.assert_called_once() - self.assertIn('load', commands[1]) - self.assertIn('inspect', commands[2]) - self.assertTrue(commands[3][0].endswith('parsar-runtime')) + self.assertIn('load', commands[2]) + self.assertIn('inspect', commands[3]) + self.assertTrue(commands[4][0].endswith('parsar-runtime')) + + def test_containerd_cache_passes_actual_id_to_launcher_without_changing_manifest(self): + expected = self.manifest['image_manifest_digests']['runtime'] + published = json.dumps(self.manifest, sort_keys=True) + def checked(command, message, timeout=30): + if 'inspect' in command: + if command[command.index('inspect') + 1] == self.manifest['images']['runtime']: + raise installer.InstallError('image missing') + return expected + ' linux/amd64' + self.assertEqual(command[command.index('--image') + 1], expected) + return json.dumps({'container': 'parsar-selfhost-'+'c'*32, 'status': 'started'}) + with patch.object(installer, 'load_manifest', return_value=self.manifest), \ + patch.object(installer, 'obtain_artifact', side_effect=lambda m, n, dest, o: dest), \ + patch.object(installer, 'runtime_archive', side_effect=AssertionError('warm Runtime download')), \ + patch.object(installer, 'checked', side_effect=checked), contextlib.redirect_stdout(io.StringIO()): + installer.install(self.args, self.root) + self.assertEqual(json.dumps(self.manifest, sort_keys=True), published) + state = json.loads(installer.private_read(self.root/'installation.json')) + self.assertEqual(state['runtime_image'], self.manifest['images']['runtime']) + self.assertEqual(state['runtime_manifest'], expected) + + def test_wrong_loaded_id_or_platform_prevents_launch_receipt(self): + for observed in ('sha256:' + 'f' * 64 + ' linux/amd64', self.manifest['images']['runtime'] + ' linux/arm64'): + loaded = False + def checked(command, message, timeout=30): + nonlocal loaded + if 'load' in command: + loaded = True + return '' + if 'inspect' in command: + if not loaded: + raise installer.InstallError('image missing') + return observed + self.fail('Runtime launcher must not run') + with self.subTest(observed=observed), patch.object(installer, 'load_manifest', return_value=self.manifest), \ + patch.object(installer, 'obtain_artifact', side_effect=lambda m, n, dest, o: dest), \ + patch.object(installer, 'runtime_archive', return_value=self.root/'runtime.tar'), \ + patch.object(installer, 'checked', side_effect=checked), \ + self.assertRaisesRegex(distribution.DistributionError, 'identity or platform'): + installer.install(self.args, self.root) + self.assertFalse((self.root/'launch.json').exists()) def test_existing_launch_verifies_labels_and_reports_state(self): name = 'parsar-selfhost-'+'c'*32 - state = {'installation_id': str(uuid.uuid4()), 'environment_id': self.environment} + state = {'installation_id': str(uuid.uuid4()), 'environment_id': self.environment, + 'runtime_image': self.manifest['images']['runtime'], + 'runtime_manifest': self.manifest['image_manifest_digests']['runtime']} labels = {'io.parsar.agents-api.installation': state['installation_id'], 'io.parsar.agents-api.environment': self.environment, 'io.parsar.agents-api.user-owned': 'true'} installer.write_private(self.root/'started.json', {'container': name, 'status': 'started'}) - with patch.object(installer, 'checked', return_value=json.dumps(labels)+' running'), \ + with patch.object(installer, 'checked', return_value=json.dumps(labels)+' running '+state['runtime_image']), \ contextlib.redirect_stdout(io.StringIO()) as output: installer.inspect_prior_launch(self.root, state) self.assertIn('already running', output.getvalue()) - self.assertIn('Read the Session to confirm connection', output.getvalue()) - with patch.object(installer, 'checked', return_value=json.dumps(labels)+' exited'): + self.assertNotIn('connected to Environment', output.getvalue()) + with patch.object(installer, 'checked', return_value=json.dumps(labels)+' running '+state['runtime_manifest']), \ + contextlib.redirect_stdout(io.StringIO()): + self.assertEqual(installer.inspect_prior_launch(self.root, state), name) + with patch.object(installer, 'checked', return_value=json.dumps(labels)+' exited '+state['runtime_image']): with self.assertRaisesRegex(installer.InstallError, 'start '+name): installer.inspect_prior_launch(self.root, state) + with patch.object(installer, 'checked', return_value=json.dumps(labels)+' running sha256:'+'f'*64): + with self.assertRaisesRegex(installer.InstallError, 'container image does not match'): + installer.inspect_prior_launch(self.root, state) labels['io.parsar.agents-api.environment'] = str(uuid.uuid4()) - with patch.object(installer, 'checked', return_value=json.dumps(labels)+' running'): + with patch.object(installer, 'checked', return_value=json.dumps(labels)+' running '+state['runtime_image']): with self.assertRaisesRegex(installer.InstallError, 'does not match'): installer.inspect_prior_launch(self.root, state) def test_uncertain_launch_has_specific_safe_inspection_commands(self): - state = {'installation_id': str(uuid.uuid4()), 'environment_id': self.environment} + state = {'installation_id': str(uuid.uuid4()), 'environment_id': self.environment, + 'runtime_image': self.manifest['images']['runtime'], + 'runtime_manifest': self.manifest['image_manifest_digests']['runtime']} with self.assertRaises(installer.InstallError) as failure: installer.inspect_prior_launch(self.root, state) message = str(failure.exception) diff --git a/docs/getting-started/install.md b/docs/getting-started/install.md index e3d2ab160..ca21dd827 100644 --- a/docs/getting-started/install.md +++ b/docs/getting-started/install.md @@ -23,11 +23,14 @@ in the [add-node steps](#add-nodes-after-a-default-installation). ## Verify, extract and install -Obtain the archive and checksum from a trusted distributor. Until a release is -published, obtain a verified bundle from your deployment administrator or build -an archive using [Build a distribution](#build-a-distribution); -a source checkout alone is not an installable binary bundle. Do not substitute an -unpublished download URL. +Download the matching Linux amd64 archive and its `.sha256` file from +[GitHub Releases](https://github.com/MiniMax-AI/parsar-core/releases). Use one +release for the entire installation. If GitHub requires sign-in, use an authenticated +browser or `gh release download RELEASE --repo MiniMax-AI/parsar-core`. +Choose the ordinary `.tar.gz` for a zero-node installation, or `-offline.tar.gz` +when you also need all execution assets locally. A source checkout alone is not +an installable binary bundle; [build a distribution](#build-a-distribution) for +unreleased changes. For the recommended node workflow, choose an HTTPS address that both node hosts and their sandbox guests can reach, such as `https://core.example`. Configure @@ -175,8 +178,13 @@ The installer prepares the matched daemon, native harnesses and local workspace inside the same isolated Runtime used for hosted execution. No shared node, model credential or source build is needed. Model access remains execution input. -A started Runtime is not proof of a connected Environment or a successful model -request. Read the Session to confirm its connection, then submit your task. +The command waits for Core to confirm that this Environment and its restricted +credential are connected. It distinguishes a running container from a connected +Environment. Connection failure or timeout exits with diagnostic and retry +instructions, preserving the same container, credentials and native history. +Rerun the command after correcting the reported problem; it does not create +replacement history. A connected Environment does not prove model availability. +Submit your task through the Session API to test execution. Installation state stays under `~/.parsar/self-hosted/ENVIRONMENT_UUID`; retain its credentials, volumes and native history. An uncertain launch gives inspection instructions instead of creating replacement history. Session deletion does not @@ -260,7 +268,9 @@ For nodes added through Web, install with the intended shared HTTPS endpoint: ./install.sh --public-url https://core.example ``` -Configure your TLS reverse proxy to forward that origin to the loopback Web port, +The installer also uses this origin for the `wss` connection URL returned by +self-hosted Sessions, so remote Runtime hosts never receive a Compose-only +hostname. Configure your TLS reverse proxy to forward that origin to the loopback Web port, preserve Host, and support WebSocket upgrades. The bundled Web forwards the fixed node and daemon transport routes to Core using their own credentials. Both the node host and its sandbox guests must reach this address. Installation does not @@ -313,3 +323,35 @@ An explicit manual option can create an unpublished draft release. Neither a successful build nor a draft makes a private repository anonymously downloadable; publish qualified assets through your chosen distribution channel before sharing installation instructions with external users. + +## Produce and qualify a release + +The `core-release` GitHub Actions workflow builds production assets from a full +committed source SHA. It uses the existing pinned Runtime builders; acceptance +credentials and private test certificate authorities must never enter its inputs. +Run the workflow from the repository's Actions page, or use: + +```sh +revision=$(git rev-parse HEAD) +gh workflow run core-release --repo MiniMax-AI/parsar-core --ref main \ + -f ref="$revision" -f offline=true -f draft_release=true +``` + +The workflow uploads the matched files as an Actions artifact and creates a draft +Release whose tag is that full SHA. The manifest records the same tag in every +asset URL. Do not mix files across releases or resolve individual components +through `latest`. The node command comes from its connected Core, which selects +the matching release automatically. + +Download the draft assets using repository access, verify their checksums, and +qualify a fresh installation plus the node/self-hosted connection paths before +publishing the draft. A workflow build alone is not live acceptance. Retain the +exact tested assets when publishing; do not rebuild or replace files under the +same release identity. Publishing a Release does not change repository visibility. + +For an offline installation, provide the extracted matching archive through the +existing `--offline-root` option where supported. Remote node commands use the +manifest's release URL; use the explicitly configured console-hosted offline +build described above when node hosts cannot access that URL. Download access +errors should be fixed at the distribution source, without passing repository +credentials into Runtime or changing its executor authorization. diff --git a/internal/agentdaemon/device/credential.go b/internal/agentdaemon/device/credential.go index fc9a26be1..e2a0377c9 100644 --- a/internal/agentdaemon/device/credential.go +++ b/internal/agentdaemon/device/credential.go @@ -14,6 +14,8 @@ type Credential struct { Name string Type string CredentialHash string + // RuntimeNodeID is the persisted managed allocation binding, never caller input. + RuntimeNodeID string } // HashCredential preserves the paired runtime bearer format, including trimming diff --git a/internal/agentdaemon/gateway/auth.go b/internal/agentdaemon/gateway/auth.go index 9709f64e1..bb9c53819 100644 --- a/internal/agentdaemon/gateway/auth.go +++ b/internal/agentdaemon/gateway/auth.go @@ -36,9 +36,10 @@ type RuntimeStore interface { // AuthenticatedRuntime is the result of a successful credential check. type AuthenticatedRuntime struct { - DeviceID string - WorkspaceID string - Name string + DeviceID string + WorkspaceID string + Name string + RuntimeNodeID string } // Authenticator validates the (device_id, token, version) trio that @@ -84,8 +85,9 @@ func (a *Authenticator) AuthenticateBearer(ctx context.Context, deviceID, bearer return AuthenticatedRuntime{}, ErrAuthBadCredential } return AuthenticatedRuntime{ - DeviceID: rt.ID, - WorkspaceID: rt.WorkspaceID, - Name: rt.Name, + DeviceID: rt.ID, + WorkspaceID: rt.WorkspaceID, + Name: rt.Name, + RuntimeNodeID: rt.RuntimeNodeID, }, nil } diff --git a/internal/agentdaemon/gateway/bootstrap_url_test.go b/internal/agentdaemon/gateway/bootstrap_url_test.go index db4b22c65..43a88d0b5 100644 --- a/internal/agentdaemon/gateway/bootstrap_url_test.go +++ b/internal/agentdaemon/gateway/bootstrap_url_test.go @@ -26,8 +26,8 @@ func TestBootstrapResolvesPublicURLOnlyAfterAuthentication(t *testing.T) { t.Run(tc.name, func(t *testing.T) { calls := 0 h := NewHandler(HandlerConfig{Registry: NewRegistry(), PublicWSURL: "ws://private-core:8091/api/v1/agent-daemon/ws", - Authenticator: NewAuthenticator(&stubRuntimeStore{ok: true, row: device.Credential{ID: "device", WorkspaceID: "tenant", Type: RuntimeTypeAgentDaemon, CredentialHash: device.HashCredential(token)}}), - ResolvePublicWSURL: func(context.Context) (string, error) { calls++; return tc.result, tc.err }, + Authenticator: NewAuthenticator(&stubRuntimeStore{ok: true, row: device.Credential{ID: "device", WorkspaceID: "tenant", Type: RuntimeTypeAgentDaemon, CredentialHash: device.HashCredential(token)}}), + ResolveWSURL: func(context.Context, AuthenticatedRuntime) (string, error) { calls++; return tc.result, tc.err }, }) r := httptest.NewRequest(http.MethodPost, "/agent-daemon/bootstrap", strings.NewReader(`{"device_id":"device"}`)) r.Header.Set("Authorization", "Bearer "+tc.bearer) diff --git a/internal/agentdaemon/gateway/handler.go b/internal/agentdaemon/gateway/handler.go index b1db3f6f6..cbfb56f1e 100644 --- a/internal/agentdaemon/gateway/handler.go +++ b/internal/agentdaemon/gateway/handler.go @@ -41,9 +41,9 @@ type HandlerConfig struct { // response so deployments behind a TLS terminator can advertise // the externally-reachable URL. PublicWSURL string - // ResolvePublicWSURL reads a deployment's configured public endpoint after - // authentication. A failure must not fall back to an unreachable private URL. - ResolvePublicWSURL func(context.Context) (string, error) + // ResolveWSURL chooses a configured endpoint from the authenticated device's + // persisted allocation binding. Resolution failures never use a fallback URL. + ResolveWSURL func(context.Context, AuthenticatedRuntime) (string, error) // OwnerStore enables multi-pod WebSocket ownership. When set, every // successful daemon WS dial-in claims device_id -> owner_pod_id in @@ -240,8 +240,8 @@ func (h *Handler) Bootstrap(w http.ResponseWriter, r *http.Request) { return } wsURL := h.cfg.PublicWSURL - if h.cfg.ResolvePublicWSURL != nil { - wsURL, err = h.cfg.ResolvePublicWSURL(r.Context()) + if h.cfg.ResolveWSURL != nil { + wsURL, err = h.cfg.ResolveWSURL(r.Context(), auth) if err != nil || wsURL == "" { writeAuthError(w, http.StatusServiceUnavailable, "bootstrap_unavailable", "Runtime connection address is unavailable") return diff --git a/scripts/core-distribution-manifest.py b/scripts/core-distribution-manifest.py index 7488db298..7e5f5b3e0 100644 --- a/scripts/core-distribution-manifest.py +++ b/scripts/core-distribution-manifest.py @@ -47,6 +47,71 @@ def verify_image(image): return details +def image_identities(archive, build_id): + """Bind both Docker store identities to one exported Linux amd64 image.""" + if not DIGEST.fullmatch(build_id): + raise ValueError("Missing immutable distribution image identity") + with tarfile.open(archive, "r:") as contents: + members = contents.getmembers() + + def member(name): + matches = [entry for entry in members if entry.name == name] + if len(matches) != 1 or not matches[0].isfile(): + raise ValueError("Image archive must contain one regular " + name) + return matches[0] + + def blob(descriptor, parse=False): + digest = descriptor.get("digest", "") + size = descriptor.get("size") + if (not isinstance(digest, str) or not DIGEST.fullmatch(digest) + or type(size) is not int or size <= 0): + raise ValueError("Invalid image archive descriptor") + entry = member("blobs/sha256/" + digest.removeprefix("sha256:")) + if entry.size != size or parse and size > 1024 * 1024: + raise ValueError("Image archive descriptor size mismatch") + checksum, chunks = hashlib.sha256(), [] + with contents.extractfile(entry) as stream: + for block in iter(lambda: stream.read(1024 * 1024), b""): + checksum.update(block) + if parse: + chunks.append(block) + if "sha256:" + checksum.hexdigest() != digest: + raise ValueError("Image archive blob checksum mismatch") + return json.loads(b"".join(chunks)) if parse else None + + index = member("index.json") + if index.size > 1024 * 1024: + raise ValueError("Image archive index exceeds size limit") + with contents.extractfile(index) as stream: + descriptors = json.load(stream).get("manifests", []) + if len(descriptors) != 1: + raise ValueError("Image archive must select exactly one image") + descriptor = descriptors[0] + manifest_digest = descriptor.get("digest", "") + image = blob(descriptor, parse=True) + # A containerd export may wrap its single platform in an image index. + if descriptor.get("mediaType") in ("application/vnd.oci.image.index.v1+json", + "application/vnd.docker.distribution.manifest.list.v2+json"): + descriptors = image.get("manifests", []) + if len(descriptors) != 1: + raise ValueError("Image archive must select exactly one platform") + descriptor = descriptors[0] + image = blob(descriptor, parse=True) + if descriptor.get("mediaType") not in ("application/vnd.oci.image.manifest.v1+json", + "application/vnd.docker.distribution.manifest.v2+json"): + raise ValueError("Image archive must select an image manifest") + config_descriptor = image.get("config", {}) + config = blob(config_descriptor, parse=True) + config_digest = config_descriptor["digest"] + if config.get("os") != "linux" or config.get("architecture") != "amd64": + raise ValueError("Image archive contains an unexpected platform") + for layer in image.get("layers", []): + blob(layer) + if build_id not in (config_digest, manifest_digest): + raise ValueError("Image archive does not match the selected build image") + return config_digest, manifest_digest + + def verify_runtime(image, daemon, helpers, source): details = verify_image(image) helpers, source = pathlib.Path(helpers), pathlib.Path(source) @@ -152,16 +217,17 @@ def manifest(bundle, stage, revision, source_tree, artifact_base_url="", offline raise ValueError("msb did not return an immutable OCI manifest digest") if inspected.get("architecture") != "amd64" or inspected.get("os") != "linux": raise ValueError("msb imported an unexpected Runtime platform") - images = {name: (stage / (name + ".id")).read_text().strip() for name in ("core", "web", "runtime", "database")} - if any(not DIGEST.fullmatch(image) for image in images.values()): - raise ValueError("Missing immutable distribution image identity") + identities = {name: image_identities(bundle / "images" / (name + ".tar"), + (stage / (name + ".id")).read_text().strip()) + for name in ("core", "web", "runtime", "database")} metadata = { "source_commit": revision, "source_tree": source_tree, "platform": "linux/amd64", "artifact_base_url": artifact_base_url, "artifacts": package_artifacts(bundle, stage, revision), - "images": images, + "images": {name: identity[0] for name, identity in identities.items()}, + "image_manifest_digests": {name: identity[1] for name, identity in identities.items()}, "runtime_ref": "parsar-core-runtime@" + digest, "microsandbox": { "version": "0.7.2", diff --git a/scripts/core-distribution-manifest.test.py b/scripts/core-distribution-manifest.test.py index e8dd709cb..2e6cbea45 100644 --- a/scripts/core-distribution-manifest.test.py +++ b/scripts/core-distribution-manifest.test.py @@ -3,6 +3,7 @@ import gzip import hashlib import importlib.util +import io import json import pathlib import tarfile @@ -21,6 +22,32 @@ RELEASE_BASE = "https://example.com/releases/" + REVISION +def image_archive(path, name, *, nested=False, architecture="amd64", corrupt=None, multiple=False): + blobs = {} + + def descriptor(data, media_type, label): + raw = json.dumps(data).encode() if isinstance(data, dict) else data + digest = hashlib.sha256(raw).hexdigest() + blobs["blobs/sha256/" + digest] = raw if corrupt != label else b"!" + raw[1:] + return {"digest": "sha256:" + digest, "size": len(raw), "mediaType": media_type} + + config = descriptor({"os": "linux", "architecture": architecture, "image": name}, + "application/vnd.oci.image.config.v1+json", "config") + layer = descriptor(b"layer:" + name.encode(), "application/vnd.oci.image.layer.v1.tar", "layer") + image = descriptor({"schemaVersion": 2, "config": config, "layers": [layer]}, + "application/vnd.oci.image.manifest.v1+json", "manifest") + if nested: + image = descriptor({"schemaVersion": 2, "manifests": [image]}, + "application/vnd.oci.image.index.v1+json", "index") + blobs["index.json"] = json.dumps({"schemaVersion": 2, "manifests": [image] * (2 if multiple else 1)}).encode() + with tarfile.open(path, "w") as archive: + for filename, raw in blobs.items(): + entry = tarfile.TarInfo(filename) + entry.size = len(raw) + archive.addfile(entry, io.BytesIO(raw)) + return config["digest"], image["digest"] + + class DistributionTests(unittest.TestCase): def setUp(self): self.temporary = tempfile.TemporaryDirectory() @@ -34,8 +61,6 @@ def setUp(self): runtime.mkdir(parents=True) (runtime / "msb").write_bytes(b"runtime") (runtime / "libkrunfw.so.5.6.1").write_bytes(b"firmware") - for number, name in enumerate(("core", "web", "runtime", "database"), 1): - (self.stage / (name + ".id")).write_text("sha256:" + str(number) * 64 + "\n") self.inspection = {"digest": "sha256:" + "a" * 64, "architecture": "amd64", "os": "linux"} self.write_inspection() for logical in distribution.ARTIFACTS: @@ -44,6 +69,11 @@ def setUp(self): path.write_bytes(b"payload:" + logical.encode()) if logical.startswith("native/"): path.chmod(0o555) + self.identities = {} + for name in ("core", "web", "runtime", "database"): + self.identities[name] = image_archive(self.bundle / "images" / (name + ".tar"), name) + (self.stage / (name + ".id")).write_text(self.identities[name][0] + "\n") + self.runtime_bytes = (self.bundle / "images/runtime.tar").read_bytes() def manifest(self, base=RELEASE_BASE, offline="0"): distribution.manifest(self.bundle, self.stage, REVISION, TREE, base, offline) @@ -55,12 +85,44 @@ def test_oci_manifest_identity_is_distinct_from_docker_config_identity(self): self.manifest() metadata = json.loads((self.bundle / "manifest.json").read_text()) self.assertEqual(metadata["runtime_ref"], "parsar-core-runtime@sha256:" + "a" * 64) - self.assertEqual(metadata["images"]["runtime"], "sha256:" + "3" * 64) + self.assertEqual(metadata["images"]["runtime"], self.identities["runtime"][0]) + self.assertEqual(metadata["image_manifest_digests"]["runtime"], self.identities["runtime"][1]) self.assertEqual(metadata["microsandbox"]["runtime_sha256"], hashlib.sha256(b"runtime").hexdigest()) for line in (self.bundle / "SHA256SUMS").read_text().splitlines(): digest, name = line.split(" ", 1) self.assertEqual(digest, distribution.sha256(self.bundle / name)) + def test_containerd_build_ids_still_publish_archive_config_ids(self): + for name, identity in self.identities.items(): + (self.stage / (name + ".id")).write_text(identity[1]) + self.manifest() + metadata = json.loads((self.bundle / "manifest.json").read_text()) + self.assertEqual(metadata["images"], {name: identity[0] for name, identity in self.identities.items()}) + self.assertEqual(metadata["image_manifest_digests"], {name: identity[1] for name, identity in self.identities.items()}) + + def test_nested_single_platform_index_retains_its_containerd_identity(self): + path = self.stage / "nested.tar" + config, index = image_archive(path, "nested", nested=True) + self.assertEqual(distribution.image_identities(path, index), (config, index)) + self.assertEqual(distribution.image_identities(path, config), (config, index)) + + def test_archive_identity_rejects_unrelated_build_or_ambiguous_platform(self): + path = self.stage / "bad-image.tar" + for options, expected in (({}, "selected build"), ({"multiple": True}, "exactly one image"), + ({"architecture": "arm64"}, "unexpected platform")): + with self.subTest(options=options): + image_archive(path, "bad", **options) + with self.assertRaisesRegex(ValueError, expected): + distribution.image_identities(path, "sha256:" + "f" * 64) + + def test_archive_identity_verifies_manifest_config_and_layer_bytes(self): + path = self.stage / "corrupt.tar" + for corrupt in ("manifest", "config", "layer", "index"): + with self.subTest(corrupt=corrupt): + config, _ = image_archive(path, "corrupt", nested=True, corrupt=corrupt) + with self.assertRaisesRegex(ValueError, "checksum mismatch"): + distribution.image_identities(path, config) + def test_missing_manifest_digest_does_not_fall_back_to_config_id(self): self.inspection.pop("digest") self.inspection["config"] = {"digest": "sha256:" + "3" * 64} @@ -112,10 +174,10 @@ def test_thin_archive_and_detached_payload_share_one_manifest(self): self.assertEqual(artifact["size"], path.stat().st_size) self.assertEqual(artifact["sha256"], distribution.sha256(path)) runtime = self.stage / "artifacts" / metadata["artifacts"]["images/runtime.tar.gz"]["filename"] - self.assertEqual(gzip.decompress(runtime.read_bytes()), b"payload:images/runtime.tar.gz") + self.assertEqual(gzip.decompress(runtime.read_bytes()), self.runtime_bytes) runtime_entry = metadata["artifacts"]["images/runtime.tar.gz"] - self.assertEqual(runtime_entry["unpacked_size"], len(b"payload:images/runtime.tar.gz")) - self.assertEqual(runtime_entry["unpacked_sha256"], hashlib.sha256(b"payload:images/runtime.tar.gz").hexdigest()) + self.assertEqual(runtime_entry["unpacked_size"], len(self.runtime_bytes)) + self.assertEqual(runtime_entry["unpacked_sha256"], hashlib.sha256(self.runtime_bytes).hexdigest()) self.assertFalse((self.bundle / "images/runtime.tar").exists()) distribution.archive(self.bundle, "1700000000") thin = self.bundle.with_name(self.bundle.name + ".tar.gz") diff --git a/services/agents-api/cmd/server/daemon_bootstrap.go b/services/agents-api/cmd/server/daemon_bootstrap.go new file mode 100644 index 000000000..29669f0b5 --- /dev/null +++ b/services/agents-api/cmd/server/daemon_bootstrap.go @@ -0,0 +1,41 @@ +package main + +import ( + "context" + "errors" + "net/url" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" +) + +func (m *managedNodes) webSocketURL(publicURL string) func(context.Context, gateway.AuthenticatedRuntime) (string, error) { + return func(ctx context.Context, auth gateway.AuthenticatedRuntime) (string, error) { + if m == nil || auth.RuntimeNodeID == "" { + return publicURL, nil + } + if m.runtime != nil && auth.RuntimeNodeID == m.runtime.LocalNodeID { + return runtimeWebSocketURL(m.runtime.CoreURL) + } + if m.setup != nil { + return m.setup.webSocketURL(publicURL)(ctx) + } + return publicURL, nil + } +} + +func runtimeWebSocketURL(coreURL string) (string, error) { + u, err := url.Parse(coreURL) + if err != nil || u.Hostname() == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" { + return "", errors.New("managed Runtime Core address is unavailable") + } + switch u.Scheme { + case "https": + u.Scheme = "wss" + case "http": + u.Scheme = "ws" + default: + return "", errors.New("managed Runtime Core address is unavailable") + } + u.Path = "/api/v1/agent-daemon/ws" + return u.String(), nil +} diff --git a/services/agents-api/cmd/server/daemon_bootstrap_test.go b/services/agents-api/cmd/server/daemon_bootstrap_test.go new file mode 100644 index 000000000..06118bd16 --- /dev/null +++ b/services/agents-api/cmd/server/daemon_bootstrap_test.go @@ -0,0 +1,75 @@ +package main + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" +) + +type bootstrapCredentialStore struct{ nodeID string } + +func (s bootstrapCredentialStore) GetDeviceCredential(context.Context, string) (device.Credential, bool, error) { + return device.Credential{ID: "runtime", Type: gateway.RuntimeTypeAgentDaemon, + CredentialHash: device.HashCredential("synthetic-token"), RuntimeNodeID: s.nodeID}, true, nil +} + +func TestBootstrapAddressFollowsAuthenticatedAllocation(t *testing.T) { + const publicURL = "wss://private-proxy.example/api/v1/agent-daemon/ws" + const selectedURL = "wss://selected-node-entry.example/api/v1/agent-daemon/ws" + local := &managedNodes{runtime: &execution.RuntimeProvider{LocalNodeID: "local-node", CoreURL: "http://host.microsandbox.internal:8091/api/v1"}} + remote := &managedNodes{setup: &managedSetup{}} + remote.setup.selected.Store(&execution.RuntimeProvider{CoreURL: "https://selected-node-entry.example/api/v1"}) + zero := &managedNodes{setup: &managedSetup{}} + for _, tc := range []struct { + name, node, want string + managed *managedNodes + }{ + {"embedded managed", "local-node", "ws://host.microsandbox.internal:8091/api/v1/agent-daemon/ws", local}, + {"remote managed beside embedded node", "remote-node", publicURL, local}, + {"self-hosted beside embedded node", "", publicURL, local}, + {"remote managed after Web setup", "remote-node", selectedURL, remote}, + {"self-hosted after Web setup", "", publicURL, remote}, + {"self-hosted before Web setup", "", publicURL, zero}, + {"standalone self-hosted", "", publicURL, nil}, + } { + t.Run(tc.name, func(t *testing.T) { + h := gateway.NewHandler(gateway.HandlerConfig{Registry: gateway.NewRegistry(), PublicWSURL: publicURL, + Authenticator: gateway.NewAuthenticator(bootstrapCredentialStore{nodeID: tc.node}), + ResolveWSURL: tc.managed.webSocketURL(publicURL)}) + request := httptest.NewRequest(http.MethodPost, "https://forged.example/api/v1/agent-daemon/bootstrap", + strings.NewReader(`{"device_id":"runtime","node_id":"local-node","runtime_node_id":"local-node"}`)) + request.Header.Set("Authorization", "Bearer synthetic-token") + request.Header.Set("X-Forwarded-Host", "forged-proxy.example") + response := httptest.NewRecorder() + h.Bootstrap(response, request) + var body map[string]any + if response.Code != http.StatusOK || json.Unmarshal(response.Body.Bytes(), &body) != nil || body["ws_url"] != tc.want { + t.Fatalf("bootstrap status=%d body=%s", response.Code, response.Body.String()) + } + if len(body) != 5 { + t.Fatal("bootstrap exposed private allocation metadata") + } + }) + } +} + +func TestEmbeddedBootstrapDoesNotFallbackWhenInternalAddressIsInvalid(t *testing.T) { + managed := &managedNodes{runtime: &execution.RuntimeProvider{LocalNodeID: "local-node", CoreURL: "invalid"}} + h := gateway.NewHandler(gateway.HandlerConfig{Registry: gateway.NewRegistry(), PublicWSURL: "wss://public.example/api/v1/agent-daemon/ws", + Authenticator: gateway.NewAuthenticator(bootstrapCredentialStore{nodeID: "local-node"}), + ResolveWSURL: managed.webSocketURL("wss://public.example/api/v1/agent-daemon/ws")}) + request := httptest.NewRequest(http.MethodPost, "/api/v1/agent-daemon/bootstrap", strings.NewReader(`{"device_id":"runtime"}`)) + request.Header.Set("Authorization", "Bearer synthetic-token") + response := httptest.NewRecorder() + h.Bootstrap(response, request) + if response.Code != http.StatusServiceUnavailable || strings.Contains(response.Body.String(), "public.example") { + t.Fatalf("bootstrap fell back after route failure: %d %s", response.Code, response.Body.String()) + } +} diff --git a/services/agents-api/cmd/server/main.go b/services/agents-api/cmd/server/main.go index 3993c485d..506fab8cd 100644 --- a/services/agents-api/cmd/server/main.go +++ b/services/agents-api/cmd/server/main.go @@ -180,11 +180,7 @@ func run() error { var daemonHandler http.Handler var registry *gateway.Registry if wsURL := os.Getenv("AGENTS_API_DAEMON_WS_URL"); wsURL != "" { - if managedNodes != nil && managedNodes.setup != nil { - daemonHandler, registry, err = runtime.NewGatewayWithURLResolver(executionStore, wsURL, managedNodes.setup.webSocketURL(wsURL)) - } else { - daemonHandler, registry, err = runtime.NewGateway(executionStore, wsURL) - } + daemonHandler, registry, err = runtime.NewGatewayWithURLResolver(executionStore, wsURL, managedNodes.webSocketURL(wsURL)) if err != nil { return err } @@ -255,6 +251,7 @@ func run() error { mux := http.NewServeMux() mux.Handle("/api/v1/agent-daemon/", daemonHandler) mux.Handle("/api/v1/agent-daemon/enroll", runtimeenrollment.EnrollmentHandler(executionStore)) + mux.Handle("/api/v1/agent-daemon/connection", runtimeenrollment.ConnectionHandler(executionStore, registry)) if managedNodes != nil { mux.Handle("/core/v1/sandbox/node/connect", managedNodes.hub) diff --git a/services/agents-api/cmd/server/managed_setup.go b/services/agents-api/cmd/server/managed_setup.go index a08ea601b..fcc9aba61 100644 --- a/services/agents-api/cmd/server/managed_setup.go +++ b/services/agents-api/cmd/server/managed_setup.go @@ -3,7 +3,6 @@ package main import ( "context" "errors" - "net/url" "sync/atomic" "time" @@ -59,17 +58,7 @@ func (s *managedSetup) webSocketURL(fallback string) func(context.Context) (stri if selected == nil { return fallback, nil } - u, err := url.Parse(selected.CoreURL) - if err != nil { - return "", err - } - if u.Scheme == "https" { - u.Scheme = "wss" - } else { - u.Scheme = "ws" - } - u.Path = "/api/v1/agent-daemon/ws" - return u.String(), nil + return runtimeWebSocketURL(selected.CoreURL) } } diff --git a/services/agents-api/internal/db/queries/devices.sql b/services/agents-api/internal/db/queries/devices.sql index 698c5ecfe..7d9dc0af3 100644 --- a/services/agents-api/internal/db/queries/devices.sql +++ b/services/agents-api/internal/db/queries/devices.sql @@ -6,7 +6,10 @@ VALUES ($1, $2, $3, $4) RETURNING id; SELECT id, name FROM devices WHERE tenant_id = $1 AND id = $2 AND revoked_at IS NULL; -- name: GetDeviceCredential :one -SELECT id, name, credential_hash FROM runtime_device_authority WHERE id = $1; +SELECT d.id, d.name, d.credential_hash, COALESCE(a.node_id::text, '')::text AS runtime_node_id +FROM runtime_device_authority d +LEFT JOIN runtime_allocations a ON a.device_id = d.id +WHERE d.id = $1; -- name: RevokeDevice :execrows UPDATE devices SET revoked_at = COALESCE(revoked_at, clock_timestamp()) diff --git a/services/agents-api/internal/db/sqlc/devices.sql.go b/services/agents-api/internal/db/sqlc/devices.sql.go index 8cc805f5d..9f5c0d13b 100644 --- a/services/agents-api/internal/db/sqlc/devices.sql.go +++ b/services/agents-api/internal/db/sqlc/devices.sql.go @@ -82,19 +82,28 @@ func (q *Queries) GetDevice(ctx context.Context, arg GetDeviceParams) (GetDevice } const getDeviceCredential = `-- name: GetDeviceCredential :one -SELECT id, name, credential_hash FROM runtime_device_authority WHERE id = $1 +SELECT d.id, d.name, d.credential_hash, COALESCE(a.node_id::text, '')::text AS runtime_node_id +FROM runtime_device_authority d +LEFT JOIN runtime_allocations a ON a.device_id = d.id +WHERE d.id = $1 ` type GetDeviceCredentialRow struct { ID pgtype.UUID `json:"id"` Name string `json:"name"` CredentialHash string `json:"credential_hash"` + RuntimeNodeID string `json:"runtime_node_id"` } func (q *Queries) GetDeviceCredential(ctx context.Context, id pgtype.UUID) (GetDeviceCredentialRow, error) { row := q.db.QueryRow(ctx, getDeviceCredential, id) var i GetDeviceCredentialRow - err := row.Scan(&i.ID, &i.Name, &i.CredentialHash) + err := row.Scan( + &i.ID, + &i.Name, + &i.CredentialHash, + &i.RuntimeNodeID, + ) return i, err } diff --git a/services/agents-api/internal/runtime/gateway.go b/services/agents-api/internal/runtime/gateway.go index 987e80cf4..053935a15 100644 --- a/services/agents-api/internal/runtime/gateway.go +++ b/services/agents-api/internal/runtime/gateway.go @@ -23,9 +23,9 @@ func NewGateway(s DeviceStore, publicWSURL string) (http.Handler, *gateway.Regis return NewGatewayWithURLResolver(s, publicWSURL, nil) } -// NewGatewayWithURLResolver allows a zero-node deployment to publish its chosen -// public address after setup without replacing its gateway or live connections. -func NewGatewayWithURLResolver(s DeviceStore, publicWSURL string, resolve func(context.Context) (string, error)) (http.Handler, *gateway.Registry, error) { +// NewGatewayWithURLResolver chooses the bootstrap address using authenticated +// allocation bindings without replacing the gateway or its live connections. +func NewGatewayWithURLResolver(s DeviceStore, publicWSURL string, resolve func(context.Context, gateway.AuthenticatedRuntime) (string, error)) (http.Handler, *gateway.Registry, error) { u, err := url.Parse(publicWSURL) if err != nil || s == nil || (u.Scheme != "ws" && u.Scheme != "wss") || u.Hostname() == "" || u.User != nil || u.RawQuery != "" || u.Fragment != "" || u.Path != "/api/v1/agent-daemon/ws" { return nil, nil, errors.New("daemon URL must be an absolute ws(s) URL ending in /api/v1/agent-daemon/ws") @@ -33,7 +33,7 @@ func NewGatewayWithURLResolver(s DeviceStore, publicWSURL string, resolve func(c registry := gateway.NewRegistry() h := gateway.NewHandler(gateway.HandlerConfig{ Authenticator: gateway.NewAuthenticator(s), Registry: registry, - Heartbeat: s, PublicWSURL: publicWSURL, ResolvePublicWSURL: resolve, + Heartbeat: s, PublicWSURL: publicWSURL, ResolveWSURL: resolve, }) r := chi.NewRouter() r.Route("/api/v1", func(r chi.Router) { gateway.RegisterRoutes(r, h) }) diff --git a/services/agents-api/internal/runtimeenrollment/connection.go b/services/agents-api/internal/runtimeenrollment/connection.go new file mode 100644 index 000000000..07563337c --- /dev/null +++ b/services/agents-api/internal/runtimeenrollment/connection.go @@ -0,0 +1,119 @@ +package runtimeenrollment + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/url" + "strings" + "time" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" +) + +type ConnectionStore interface { + AuthenticateEnvironmentExecutor(context.Context, string, string) (string, error) + GetEnvironment(context.Context, string, string) (store.Environment, error) + GetSessionDevice(context.Context, string, string) (store.ExecutionDevice, error) + GetDeviceCredential(context.Context, string) (device.Credential, bool, error) +} + +// ConnectionHandler observes an existing binding without enrollment or execution. +// Executor authority never grants access to the public Session API. +func ConnectionHandler(s ConnectionStore, registry *gateway.Registry) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Cache-Control", "no-store") + fail := func(status int) { http.Error(w, http.StatusText(status), status) } + if r.Method != http.MethodGet { + w.Header().Set("Allow", http.MethodGet) + fail(http.StatusMethodNotAllowed) + return + } + authorization := strings.Fields(r.Header.Get("Authorization")) + if len(authorization) != 2 || !strings.EqualFold(authorization[0], "Bearer") { + fail(http.StatusUnauthorized) + return + } + query, err := url.ParseQuery(r.URL.RawQuery) + if err != nil || len(query) != 1 || len(query["environment_id"]) != 1 || query.Get("environment_id") == "" { + fail(http.StatusBadRequest) + return + } + environment := query.Get("environment_id") + digest := device.HashCredential(authorization[1]) + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + connected, err := runtimeConnected(ctx, s, registry, environment, digest) + switch { + case errors.Is(err, store.ErrNotFound): + fail(http.StatusUnauthorized) + case errors.Is(err, store.ErrDeviceBindingConflict): + fail(http.StatusConflict) + case err != nil: + fail(http.StatusServiceUnavailable) + default: + status := "disconnected" + if connected { + status = "connected" + } + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(struct { + EnvironmentID string `json:"environment_id"` + Status string `json:"status"` + }{environment, status}) + } + }) +} + +func runtimeConnected(ctx context.Context, s ConnectionStore, registry *gateway.Registry, environment, digest string) (bool, error) { + tenant, err := s.AuthenticateEnvironmentExecutor(ctx, environment, digest) + if err != nil { + return false, err + } + current, err := s.GetEnvironment(ctx, tenant, environment) + if err != nil { + return false, err + } + if current.Status == "failed" || current.Status == "expired" { + return false, store.ErrNotFound + } + bound, err := s.GetSessionDevice(ctx, tenant, current.SessionID) + if errors.Is(err, store.ErrNotFound) { + return false, nil + } + if err != nil { + return false, err + } + if bound.EnvironmentID != environment { + return false, store.ErrDeviceBindingConflict + } + credential, found, err := s.GetDeviceCredential(ctx, bound.ID) + if err != nil { + return false, err + } + if !found { + return false, store.ErrNotFound + } + if credential.CredentialHash != digest { + return false, store.ErrDeviceBindingConflict + } + peer, err := registry.LookupDevice(bound.ID) + if errors.Is(err, gateway.ErrDeviceNotRegistered) { + return false, nil + } + if err != nil { + return false, err + } + if current.Status != "connected" || peer.IsClosed() || !peer.AuthenticatedWith(digest) { + return false, nil + } + // Recheck authority after reading the socket; rotation/revocation never inherits + // the connected observation of a socket authenticated with the former key. + if _, err = s.AuthenticateEnvironmentExecutor(ctx, environment, digest); err != nil { + return false, err + } + return !peer.IsClosed(), nil +} diff --git a/services/agents-api/internal/runtimeenrollment/connection_test.go b/services/agents-api/internal/runtimeenrollment/connection_test.go new file mode 100644 index 000000000..15ea4fee2 --- /dev/null +++ b/services/agents-api/internal/runtimeenrollment/connection_test.go @@ -0,0 +1,67 @@ +package runtimeenrollment + +import ( + "context" + "errors" + "net/http/httptest" + "strings" + "testing" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" +) + +type connectionStub struct { + err error + calls int +} + +func (s *connectionStub) AuthenticateEnvironmentExecutor(_ context.Context, environment, digest string) (string, error) { + s.calls++ + if environment != "environment" || digest != device.HashCredential("test-key") { + return "", store.ErrNotFound + } + return "tenant", s.err +} +func (*connectionStub) GetEnvironment(context.Context, string, string) (store.Environment, error) { + return store.Environment{ID: "environment", SessionID: "session", Status: "pending"}, nil +} +func (*connectionStub) GetSessionDevice(context.Context, string, string) (store.ExecutionDevice, error) { + return store.ExecutionDevice{}, store.ErrNotFound +} +func (*connectionStub) GetDeviceCredential(context.Context, string) (device.Credential, bool, error) { + panic("unbound lookup") +} + +func TestConnectionReadContract(t *testing.T) { + for _, tc := range []struct { + method, query, bearer string + err error + code, calls int + }{ + {"GET", "environment_id=environment", "Bearer test-key", nil, 200, 1}, + {"GET", "environment_id=environment", "", nil, 401, 0}, + {"POST", "environment_id=environment", "Bearer test-key", nil, 405, 0}, + {"GET", "environment_id=environment&environment_id=other", "Bearer test-key", nil, 400, 0}, + {"GET", "environment_id=environment&other=1", "Bearer test-key", nil, 400, 0}, + {"GET", "environment_id=%zz", "Bearer test-key", nil, 400, 0}, + {"GET", "environment_id=environment", "Bearer test-key", store.ErrNotFound, 401, 1}, + {"GET", "environment_id=environment", "Bearer test-key", errors.New("private detail"), 503, 1}, + } { + s := &connectionStub{err: tc.err} + req := httptest.NewRequest(tc.method, "/api/v1/agent-daemon/connection?"+tc.query, nil) + req.Header.Set("Authorization", tc.bearer) + res := httptest.NewRecorder() + ConnectionHandler(s, gateway.NewRegistry()).ServeHTTP(res, req) + if res.Code != tc.code || s.calls != tc.calls || res.Header().Get("Cache-Control") != "no-store" { + t.Fatalf("%s %s: %d, %d calls", tc.method, tc.query, res.Code, s.calls) + } + if strings.Contains(res.Body.String(), "private") || strings.Contains(res.Body.String(), "test-key") { + t.Fatal("private data exposed") + } + if res.Code == 200 && res.Body.String() != `{"environment_id":"environment","status":"disconnected"}`+"\n" { + t.Fatal("unexpected response", res.Body.String()) + } + } +} diff --git a/services/agents-api/internal/store/device_bootstrap_binding_test.go b/services/agents-api/internal/store/device_bootstrap_binding_test.go new file mode 100644 index 000000000..c68aac5b6 --- /dev/null +++ b/services/agents-api/internal/store/device_bootstrap_binding_test.go @@ -0,0 +1,89 @@ +package store + +import ( + "errors" + "strings" + "testing" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" + "github.com/google/uuid" +) + +func TestDeviceCredentialCarriesPersistedAllocationNode(t *testing.T) { + s, writer, deployment := managerFixture(t, 4, 8) + token, _, err := s.CreateRuntimeEnrollment(t.Context()) + if err != nil { + t.Fatal(err) + } + remote := uuid.NewString() + _, err = s.EnrollRuntimeNode(t.Context(), token, RuntimeNodeEnrollment{NodeID: remote, Credential: strings.Repeat("x", 64), + Name: "remote", Provider: "docker", BackendFingerprint: strings.Repeat("b", 64), MaxActive: 4, MaxRetained: 8}) + if err != nil { + t.Fatal(err) + } + onlineManagerNode(t, s, remote) + for _, nodeID := range []string{deployment.LocalNodeID, remote} { + t.Run(nodeID, func(t *testing.T) { + tenant, bearer := uuid.NewString(), uuid.NewString() + session, err := s.CreateSession(t.Context(), tenant, managerSessionInput(uuid.NewString(), nodeID)) + if err != nil { + t.Fatal(err) + } + environment, err := s.GetSessionEnvironment(t.Context(), tenant, session.ID) + if err != nil { + t.Fatal(err) + } + allocation, err := writer.ReserveRuntimeAllocation(t.Context(), tenant, environment.ID, deployment.InstallationID, device.HashCredential(bearer)) + if err != nil { + t.Fatal(err) + } + authenticator := gateway.NewAuthenticator(s) + auth, err := authenticator.AuthenticateBearer(t.Context(), allocation.DeviceID, bearer) + if err != nil || auth.RuntimeNodeID != nodeID { + t.Fatalf("authenticated node=%s want=%s error=%v", auth.RuntimeNodeID, nodeID, err) + } + if _, err := authenticator.AuthenticateBearer(t.Context(), allocation.DeviceID, "wrong-token"); !errors.Is(err, gateway.ErrAuthBadCredential) { + t.Fatal("binding bypassed credential check", err) + } + if err := s.RevokeDevice(t.Context(), tenant, allocation.DeviceID); err != nil { + t.Fatal(err) + } + if _, err := authenticator.AuthenticateBearer(t.Context(), allocation.DeviceID, bearer); !errors.Is(err, gateway.ErrAuthUnknownDevice) { + t.Fatal("allocation revived revoked credential", err) + } + }) + } +} + +func TestDeviceCredentialWithoutManagedNodeRetainsPublicRouteIdentity(t *testing.T) { + s, _ := testStore(t) + tenant := uuid.NewString() + ordinary, err := s.CreateDevice(t.Context(), tenant, "ordinary", device.HashCredential("ordinary-token")) + if err != nil { + t.Fatal(err) + } + _, environment := localEnvironment(t, s, tenant) + allocation, err := executionLease(t, s).Store().ReserveRuntimeAllocation(t.Context(), tenant, environment.ID, uuid.NewString(), device.HashCredential("allocation-token")) + if err != nil { + t.Fatal(err) + } + principal := FixtureExecutorPrincipal(t, s, uuid.NewString()) + _, selfhost, key := runtimeEnrollmentFixture(t, s, principal) + enrolled, err := s.EnrollRuntime(t.Context(), selfhost.ID, executorDigest(key.Token)) + if err != nil { + t.Fatal(err) + } + for _, id := range []string{ordinary.ID, allocation.DeviceID, enrolled.DeviceID} { + credential, found, err := s.GetDeviceCredential(t.Context(), id) + if err != nil || !found || credential.RuntimeNodeID != "" { + t.Fatalf("non-node credential acquired allocation route: found=%v node=%s error=%v", found, credential.RuntimeNodeID, err) + } + } + if err := s.RevokeExecutorCredential(t.Context(), principal, key.KeyID); err != nil { + t.Fatal(err) + } + if _, found, err := s.GetDeviceCredential(t.Context(), enrolled.DeviceID); err != nil || found { + t.Fatal("LEFT JOIN revived revoked executor key", err) + } +} diff --git a/services/agents-api/internal/store/devices.go b/services/agents-api/internal/store/devices.go index 7ac67e0ad..b13687c4e 100644 --- a/services/agents-api/internal/store/devices.go +++ b/services/agents-api/internal/store/devices.go @@ -74,7 +74,7 @@ func (s *Store) GetDeviceCredential(ctx context.Context, deviceID string) (devic return device.Credential{}, false, err } return device.Credential{ID: uuid.UUID(row.ID.Bytes).String(), Name: row.Name, - Type: gateway.RuntimeTypeAgentDaemon, CredentialHash: row.CredentialHash}, true, nil + Type: gateway.RuntimeTypeAgentDaemon, CredentialHash: row.CredentialHash, RuntimeNodeID: row.RuntimeNodeID}, true, nil } func (s *Store) RevokeDevice(ctx context.Context, tenantID, deviceID string) error { diff --git a/services/agents-api/internal/store/runtime_enrollment_connection_test.go b/services/agents-api/internal/store/runtime_enrollment_connection_test.go index 5021e3712..2e2c6c7bb 100644 --- a/services/agents-api/internal/store/runtime_enrollment_connection_test.go +++ b/services/agents-api/internal/store/runtime_enrollment_connection_test.go @@ -14,6 +14,7 @@ import ( "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/proto" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtime" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeenrollment" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" "github.com/google/uuid" "github.com/gorilla/websocket" @@ -48,6 +49,36 @@ func TestEnrolledDaemonConnectionRevocationAndRestart(t *testing.T) { if err != nil { t.Fatal(err) } + connection := runtimeenrollment.ConnectionHandler(s, registry) + assertConnection := func(target, token, status string, code int) { + t.Helper() + request := httptest.NewRequest("GET", "/api/v1/agent-daemon/connection?environment_id="+target, nil) + request.Header.Set("Authorization", "Bearer "+token) + response := httptest.NewRecorder() + connection.ServeHTTP(response, request) + if response.Code != code || response.Header().Get("Cache-Control") != "no-store" { + t.Fatalf("connection status %d, want %d", response.Code, code) + } + if code == 200 { + var got map[string]string + if json.Unmarshal(response.Body.Bytes(), &got) != nil || len(got) != 2 || got["environment_id"] != target || got["status"] != status { + t.Fatalf("connection response: %s", response.Body.String()) + } + } + } + assertConnection(environment.ID, key.Token, "disconnected", 200) + assertConnection(uuid.NewString(), key.Token, "", 401) + otherKey, err := s.IssueExecutorCredential(t.Context(), principal, uuid.NewString(), environment.ID) + if err != nil { + t.Fatal(err) + } + assertConnection(environment.ID, otherKey.Token, "", 409) + foreign := store.FixtureExecutorPrincipal(t, s, uuid.NewString()) + foreignKey, err := s.IssueExecutorCredential(t.Context(), foreign, uuid.NewString(), "") + if err != nil { + t.Fatal(err) + } + assertConnection(environment.ID, foreignKey.Token, "", 401) server.Config.Handler = handler server.Start() t.Cleanup(func() { server.Close(); runtime.CloseConnections(registry) }) @@ -101,10 +132,13 @@ func TestEnrolledDaemonConnectionRevocationAndRestart(t *testing.T) { } first := connect(key.Token) await("connected") + assertConnection(environment.ID, key.Token, "connected", 200) rotated, err := s.RotateExecutorCredential(t.Context(), principal, key.KeyID) if err != nil { t.Fatal(err) } + assertConnection(environment.ID, key.Token, "", 401) + assertConnection(environment.ID, rotated.Token, "disconnected", 200) // No heartbeat is sent: the Worker's authority check must fence the old socket. await("disconnected") _ = first.SetReadDeadline(time.Now().Add(time.Second)) @@ -113,6 +147,7 @@ func TestEnrolledDaemonConnectionRevocationAndRestart(t *testing.T) { } second := connect(rotated.Token) await("connected") + assertConnection(environment.ID, rotated.Token, "connected", 200) stop() stop = nil // A new Core owner clears prior transport evidence, then observes the same @@ -122,6 +157,7 @@ func TestEnrolledDaemonConnectionRevocationAndRestart(t *testing.T) { if err = s.RevokeExecutorCredential(t.Context(), principal, key.KeyID); err != nil { t.Fatal(err) } + assertConnection(environment.ID, rotated.Token, "", 401) await("disconnected") _ = second.SetReadDeadline(time.Now().Add(time.Second)) if _, _, err = second.ReadMessage(); err == nil { diff --git a/services/core-console/distribution_test.go b/services/core-console/distribution_test.go index 595ea9851..8785d0241 100644 --- a/services/core-console/distribution_test.go +++ b/services/core-console/distribution_test.go @@ -15,7 +15,7 @@ import ( func TestEnvironmentConnectionCredentialsStayScoped(t *testing.T) { upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { want := "Bearer project-token" - if r.URL.Path == "/api/v1/agent-daemon/enroll" { + if r.URL.Path == "/api/v1/agent-daemon/enroll" || r.URL.Path == "/api/v1/agent-daemon/connection" { want = "Bearer executor-key" } if r.Header.Get("Authorization") != want { @@ -46,6 +46,9 @@ func TestEnvironmentConnectionCredentialsStayScoped(t *testing.T) { {"GET", "/core/v1/environments/env/executor-credentials", "", 404}, {"POST", "/api/v1/agent-daemon/enroll", "executor-key", 200}, {"POST", "/api/v1/agent-daemon/enroll", "", 403}, + {"GET", "/api/v1/agent-daemon/connection?environment_id=env", "executor-key", 200}, + {"GET", "/api/v1/agent-daemon/connection?environment_id=env", "", 403}, + {"POST", "/api/v1/agent-daemon/connection", "executor-key", 403}, } { req := consoleRequest(t, server, tc.method, tc.path) if tc.bearer != "" { diff --git a/services/core-console/node_installation.go b/services/core-console/node_installation.go index ea0bbbef7..ea13c4a6b 100644 --- a/services/core-console/node_installation.go +++ b/services/core-console/node_installation.go @@ -112,7 +112,7 @@ func nodeTransportRequest(r *http.Request) bool { switch r.URL.Path { case "/core/v1/sandbox/enroll", "/api/v1/agent-daemon/enroll", "/api/v1/agent-daemon/bootstrap": return r.Method == http.MethodPost && r.Header.Get("Upgrade") == "" - case "/core/v1/sandbox/node/identity", "/api/v1/agent-daemon/device-status": + case "/core/v1/sandbox/node/identity", "/api/v1/agent-daemon/device-status", "/api/v1/agent-daemon/connection": return r.Method == http.MethodGet && r.Header.Get("Upgrade") == "" case "/core/v1/sandbox/node/connect", "/api/v1/agent-daemon/ws": return r.Method == http.MethodGet && strings.EqualFold(r.Header.Get("Upgrade"), "websocket")