From 2b5dcc09a5ff411fd9ad987f5bf4d0d8547195f1 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 25 Sep 2026 18:29:51 +0000 Subject: [PATCH 01/14] Give every process setting one home in config.json, applied with parsar A new installation writes /config.json, the only file an operator edits, with every setting that applies to its mode. The installer flags only seed it. deploy/install/config.schema.json describes it; config_model.py validates the small keyword subset the schema uses and a test keeps the schema inside that subset. The installation directory now holds secrets/ (one copy of each secret), state.json (identity and install facts) and generated/, which parsar apply derives from config.json: compose.json and core.env without secrets, the Core key digest file, the native unit, the settings snapshot Core serves at GET /core/v1/installation and runtime-history.json when set. Core reads the new AGENTS_API_PUBLIC_URL, AGENTS_API_DATABASE_PASSWORD_FILE and AGENTS_API_SETTINGS_FILE; the retired AGENTS_API_DAEMON_WS_URL and AGENTS_API_CONFIG_FILE are no longer generated. The parsar zipapp in the installation directory runs status, start, stop, apply and rotate-core-key without the bundle. apply refuses hand-edited generated files, changed fixed fields and changed secrets, confirms a public URL change that strands bound nodes, recreates only the services whose inputs changed (a Compose label carries each service's input digest) and restores the previous files when Core rejects the new ones. rotate-core-key cuts over at once. install.sh reruns read config.json, reject flags and repair; --status and --stop name their replacements. install.sh --convert moves an installation made before config.json to this layout and release. Its preflight changes nothing and stops on anything it can't convert; secrets are renamed, not copied; an interrupted conversion resumes. An AGENTS_API_EXECUTION_OPTIONS_FILE stays as it is until the model settings import, which hooks into parsar apply. scripts/config-reference.py renders the settings table in docs/configuration.md from the schema; make check-distribution checks it. --- Makefile | 1 + deploy/install/config.schema.json | 254 ++++++ deploy/install/config_model.py | 287 +++++++ deploy/install/configuration.py | 294 ++++--- deploy/install/convert.py | 479 +++++++++++ deploy/install/install.py | 514 ++++++------ deploy/install/installer_fakes.py | 236 ++++++ deploy/install/local_node.py | 2 +- deploy/install/native_service.py | 88 +- deploy/install/parsar_cli.py | 676 ++++++++++++++++ deploy/install/test_config_model.py | 106 +++ deploy/install/test_convert.py | 213 +++++ deploy/install/test_install.py | 897 ++++----------------- deploy/install/test_local_node.py | 4 +- deploy/install/test_native_service.py | 61 +- deploy/install/test_parsar.py | 192 +++++ docs/configuration.md | 189 +++-- scripts/build-core-distribution.sh | 3 +- scripts/config-reference.py | 99 +++ scripts/core-distribution-manifest.py | 12 + scripts/core-distribution-manifest.test.py | 9 +- 21 files changed, 3383 insertions(+), 1233 deletions(-) create mode 100644 deploy/install/config.schema.json create mode 100644 deploy/install/config_model.py create mode 100644 deploy/install/convert.py create mode 100644 deploy/install/installer_fakes.py create mode 100644 deploy/install/parsar_cli.py create mode 100644 deploy/install/test_config_model.py create mode 100644 deploy/install/test_convert.py create mode 100644 deploy/install/test_parsar.py create mode 100755 scripts/config-reference.py diff --git a/Makefile b/Makefile index b0964a7eb..ac54cf3a1 100644 --- a/Makefile +++ b/Makefile @@ -122,6 +122,7 @@ check-distribution: go test ./services/core-console -count=1 PYTHONDONTWRITEBYTECODE=1 python3 -m unittest discover -s deploy/install -p 'test_*.py' PYTHONDONTWRITEBYTECODE=1 python3 scripts/core-distribution-manifest.test.py + PYTHONDONTWRITEBYTECODE=1 python3 scripts/config-reference.py --check bash -n deploy/install/install.sh scripts/build-core-console.sh scripts/build-core-distribution.sh scripts/prepare-release-runtimes.sh ./scripts/build-core-console.sh diff --git a/deploy/install/config.schema.json b/deploy/install/config.schema.json new file mode 100644 index 000000000..14b6b866d --- /dev/null +++ b/deploy/install/config.schema.json @@ -0,0 +1,254 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "title": "Parsar Core installation configuration", + "description": "Process settings of one installation. Edit config.json, then run parsar apply.", + "type": "object", + "additionalProperties": false, + "required": ["format", "mode"], + "properties": { + "$schema": { + "type": "string", + "description": "Editor hint that points at the installed copy of this schema. Ignored.", + "x-parsar": {"setting": false} + }, + "format": { + "const": 1, + "description": "Configuration format. Only an upgrade changes it.", + "x-parsar": {"setting": false, "changeable": false} + }, + "mode": { + "enum": ["all", "core-only", "web-only"], + "default": "all", + "description": "Which services this installation runs.", + "x-parsar": {"changeable": false, "install_flag": "--core-only or --web-only"} + }, + "native_core": { + "type": "boolean", + "default": false, + "description": "Run Core as a systemd user service instead of a container.", + "x-parsar": {"changeable": false, "modes": ["all", "core-only"], "install_flag": "--native-core"} + }, + "public_url": { + "type": ["string", "null"], + "default": null, + "description": "Public origin of Core and Web behind your TLS reverse proxy, such as https://core.example. Nodes, sandboxes and self-hosted executors use it. null means local access only through http://127.0.0.1.", + "x-parsar": { + "check": "origin", + "restarts": ["core", "web"], + "derives": ["AGENTS_API_PUBLIC_URL", "CORE_CONSOLE_ORIGIN"], + "install_flag": "--public-url" + } + }, + "ports": { + "type": "object", + "additionalProperties": false, + "properties": { + "core": { + "type": "integer", + "minimum": 1024, + "maximum": 65535, + "default": 8091, + "description": "Loopback port of the Core API. With native Core, Web follows it.", + "x-parsar": { + "modes": ["all", "core-only"], + "restarts": ["core"], + "derives": ["Core port mapping or AGENTS_API_ADDR"], + "install_flag": "--core-port" + } + }, + "web": { + "type": "integer", + "minimum": 1024, + "maximum": 65535, + "default": 8080, + "description": "Loopback port of Web.", + "x-parsar": { + "modes": ["all", "web-only"], + "restarts": ["web"], + "derives": ["Web port mapping or CORE_CONSOLE_ADDR"], + "install_flag": "--web-port" + } + }, + "database": { + "type": "integer", + "minimum": 1024, + "maximum": 65535, + "description": "Loopback port of PostgreSQL. Present exactly when native_core is true; the installer picks a free port.", + "x-parsar": { + "modes": ["all", "core-only"], + "restarts": ["database", "core"], + "derives": ["database port mapping", "AGENTS_API_DATABASE_URL"] + } + } + } + }, + "web": { + "type": "object", + "additionalProperties": false, + "x-parsar": {"modes": ["web-only"]}, + "properties": { + "core_url": { + "type": "string", + "description": "Origin of the Core that this Web connects to: HTTPS, or HTTP on a loopback host.", + "x-parsar": { + "check": "origin", + "restarts": ["web"], + "derives": ["CORE_CONSOLE_UPSTREAM"], + "install_flag": "--core-url" + } + } + } + }, + "log": { + "type": "object", + "additionalProperties": false, + "properties": { + "level": { + "enum": ["debug", "info", "warn", "error"], + "default": "info", + "description": "Minimum log level of Core and Web.", + "x-parsar": {"restarts": ["core", "web"], "derives": ["PARSAR_LOG_LEVEL"]} + }, + "format": { + "enum": ["auto", "text", "json"], + "default": "auto", + "description": "Log format. auto writes text to a terminal and JSON otherwise.", + "x-parsar": {"restarts": ["core", "web"], "derives": ["PARSAR_LOG_FORMAT"]} + }, + "add_source": { + "type": "boolean", + "default": false, + "description": "Add the source file and line to each log record.", + "x-parsar": {"restarts": ["core", "web"], "derives": ["PARSAR_LOG_ADD_SOURCE"]} + } + } + }, + "core": { + "type": "object", + "additionalProperties": false, + "x-parsar": {"modes": ["all", "core-only"]}, + "properties": { + "execution_concurrency": { + "type": "integer", + "minimum": 1, + "maximum": 1024, + "default": 4, + "description": "Concurrent execution work units in Core. Unrelated to node sandbox capacity.", + "x-parsar": {"restarts": ["core"], "derives": ["AGENTS_API_EXECUTION_CONCURRENCY"]} + }, + "harnesses": { + "type": "array", + "minItems": 1, + "uniqueItems": true, + "items": {"enum": ["claude_sdk", "codex", "mcode"]}, + "default": ["claude_sdk", "codex", "mcode"], + "description": "Harnesses that Sessions may select.", + "x-parsar": {"restarts": ["core"], "derives": ["AGENTS_API_HARNESSES"]} + }, + "default_harness": { + "enum": ["claude_sdk", "codex", "mcode"], + "default": "codex", + "description": "Harness used when a Session names none. It must be listed in core.harnesses.", + "x-parsar": {"restarts": ["core"], "derives": ["AGENTS_API_ENGINE"]} + }, + "write_audit_retention": { + "type": "string", + "default": "2160h", + "description": "How long non-creation write history is kept, as a Go duration of at least 1h.", + "x-parsar": {"check": "go_duration_min_1h", "restarts": ["core"], "derives": ["AGENTS_API_WRITE_AUDIT_RETENTION"]} + }, + "oauth_trusted_origins": { + "type": "array", + "uniqueItems": true, + "items": {"type": "string", "x-parsar": {"check": "https_origin"}}, + "default": [], + "description": "Extra HTTPS origins trusted as private OAuth issuers.", + "x-parsar": {"restarts": ["core"], "derives": ["AGENTS_API_OAUTH_TRUSTED_ORIGINS"]} + }, + "database_pool": { + "type": "object", + "additionalProperties": false, + "properties": { + "max_conns": { + "type": ["integer", "null"], + "minimum": 1, + "default": null, + "description": "Maximum database connections. null keeps the driver default, max(4, CPU count).", + "x-parsar": {"restarts": ["core"], "derives": ["pool_max_conns in AGENTS_API_DATABASE_URL"]} + }, + "min_conns": { + "type": ["integer", "null"], + "minimum": 0, + "default": null, + "description": "Minimum idle database connections. null keeps the driver default, 0.", + "x-parsar": {"restarts": ["core"], "derives": ["pool_min_conns in AGENTS_API_DATABASE_URL"]} + }, + "max_conn_lifetime": { + "type": ["string", "null"], + "default": null, + "description": "Go duration. null keeps the driver default, 1h.", + "x-parsar": {"check": "go_duration", "restarts": ["core"], "derives": ["pool_max_conn_lifetime in AGENTS_API_DATABASE_URL"]} + }, + "max_conn_idle_time": { + "type": ["string", "null"], + "default": null, + "description": "Go duration. null keeps the driver default, 30m.", + "x-parsar": {"check": "go_duration", "restarts": ["core"], "derives": ["pool_max_conn_idle_time in AGENTS_API_DATABASE_URL"]} + }, + "health_check_period": { + "type": ["string", "null"], + "default": null, + "description": "Go duration. null keeps the driver default, 1m.", + "x-parsar": {"check": "go_duration", "restarts": ["core"], "derives": ["pool_health_check_period in AGENTS_API_DATABASE_URL"]} + } + } + }, + "runtime_history": { + "type": ["object", "null"], + "default": null, + "additionalProperties": false, + "description": "Runtime history collection and OTLP export. null keeps local collection with Core's defaults. Core checks the values at startup.", + "x-parsar": {"restarts": ["core"], "derives": ["generated/runtime-history.json", "AGENTS_API_RUNTIME_HISTORY_FILE"]}, + "properties": { + "transport": { + "type": "string", + "description": "OTLP export transport.", + "x-parsar": {"restarts": ["core"]} + }, + "endpoint": { + "type": "string", + "description": "OTLP collector endpoint. Omit it to keep history local.", + "x-parsar": {"restarts": ["core"]} + }, + "insecure": { + "type": "boolean", + "description": "Export without TLS.", + "x-parsar": {"restarts": ["core"]} + }, + "headers": { + "type": "object", + "additionalProperties": {"type": "string"}, + "description": "Headers sent with each export, such as credentials. Never shown by parsar or Core.", + "x-parsar": {"sensitive": true, "restarts": ["core"]} + }, + "queue_capacity": { + "type": "integer", + "description": "Export queue capacity.", + "x-parsar": {"restarts": ["core"]} + }, + "timeout_seconds": { + "type": "integer", + "description": "Export and query timeout in seconds.", + "x-parsar": {"restarts": ["core"]} + }, + "sample_interval_seconds": { + "type": "integer", + "description": "Periodic sampling interval in seconds.", + "x-parsar": {"restarts": ["core"]} + } + } + } + } + } + } +} diff --git a/deploy/install/config_model.py b/deploy/install/config_model.py new file mode 100644 index 000000000..8a9604a04 --- /dev/null +++ b/deploy/install/config_model.py @@ -0,0 +1,287 @@ +"""The config.json model: a small JSON Schema subset, defaults and the settings snapshot. + +This deliberately implements only the keywords config.schema.json uses. It is not a +general JSON Schema engine; a test fails if the schema uses any other keyword. +""" +import copy +import ipaddress +import json +import os +import re + +MODES = ("all", "core-only", "web-only") +SERVICES = {"all": ("core", "web", "database"), "core-only": ("core", "database"), "web-only": ("web",)} +KEYWORDS = {"$schema", "title", "type", "enum", "const", "default", "description", "minimum", "maximum", + "pattern", "items", "minItems", "uniqueItems", "properties", "required", + "additionalProperties", "x-parsar"} +ANNOTATIONS = {"changeable", "modes", "restarts", "sensitive", "derives", "install_flag", "check", "setting"} + + +def _schema_text(): + # The loader reads from a directory or from inside the parsar zipapp. + path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.schema.json") + return __loader__.get_data(path).decode("utf-8") + + +SCHEMA_TEXT = _schema_text() +SCHEMA = json.loads(SCHEMA_TEXT) + + +class ConfigError(Exception): + """Every problem found in config.json, each naming its key. Values are never included.""" + + def __init__(self, problems): + self.problems = list(problems) + super().__init__("config.json is not valid:\n" + "\n".join(" - " + item for item in self.problems)) + + +def annotation(node, name, default=None): + return node.get("x-parsar", {}).get(name, default) + + +def leaves(node=None, prefix="", modes=MODES): + """Yield (dotted key, schema node, modes) for every leaf, in schema order.""" + node = SCHEMA if node is None else node + for name, child in node["properties"].items(): + key = prefix + name + child_modes = tuple(annotation(child, "modes", modes)) + if "properties" in child: + yield from leaves(child, key + ".", child_modes) + else: + yield key, child, child_modes + + +def lookup(config, key): + value = config + for part in key.split("."): + if not isinstance(value, dict) or part not in value: + return None + value = value[part] + return value + + +# Checks named by x-parsar.check. Core stays the authority for its own semantic rules. +def _origin(value, https_only=False): + """Mirror of Core's ValidateSandboxCoreURL: a canonical origin, HTTP only on loopback.""" + match = re.fullmatch(r"(https?)://([^/?#@\\\s%]+)", value) + if not match or match[2] != match[2].lower() or match[2].endswith(":"): + return False + scheme, host = match[1], match[2] + port = None + if host.startswith("["): + bracket = re.fullmatch(r"\[([^\]]+)\](?::([0-9]+))?", host) + if not bracket: + return False + name, port = bracket[1], bracket[2] + else: + name, _, port = host.partition(":") + port = port or None + if port is not None and (not port.isdigit() or str(int(port)) != port or not 1 <= int(port) <= 65535): + return False + try: + loopback = ipaddress.ip_address(name).is_loopback + except ValueError: + if host.startswith("[") or len(name) > 253: + return False + for label in name.split("."): + if not re.fullmatch(r"[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?", label): + return False + loopback = name == "localhost" + if https_only: + return scheme == "https" + return scheme == "https" or loopback + + +_DURATION_UNITS = {"ns": 1e-9, "us": 1e-6, "µs": 1e-6, "μs": 1e-6, "ms": 1e-3, "s": 1, "m": 60, "h": 3600} +_DURATION_PART = re.compile(r"([0-9]*(?:\.[0-9]*)?)(ns|us|µs|μs|ms|s|m|h)") + + +def duration_seconds(text): + """Seconds in a Go duration string, or None when Go would reject it.""" + body = text[1:] if text[:1] in "+-" else text + if body == "0": + return 0.0 + total, position = 0.0, 0 + while position < len(body): + part = _DURATION_PART.match(body, position) + if not part or part[1] in ("", "."): + return None + total += float(part[1]) * _DURATION_UNITS[part[2]] + position = part.end() + if not body: + return None + return -total if text.startswith("-") else total + + +CHECKS = { + "origin": (_origin, "must be a canonical origin such as https://core.example: lowercase, no path or " + "trailing slash, and HTTP only for a loopback host"), + "https_origin": (lambda value: _origin(value, https_only=True), "must be a canonical HTTPS origin"), + "go_duration": (lambda value: duration_seconds(value) is not None, "must be a Go duration such as 30m or 1h"), + "go_duration_min_1h": (lambda value: (duration_seconds(value) or 0) >= 3600, + "must be a Go duration of at least 1h"), +} + + +def _type_ok(value, name): + if name == "integer": + return isinstance(value, int) and not isinstance(value, bool) + return isinstance(value, {"string": str, "boolean": bool, "object": dict, "array": list, + "null": type(None)}[name]) + + +def _validate(node, value, key, mode, problems): + label = key or "config.json" + types = node.get("type") + if types is not None: + types = [types] if isinstance(types, str) else types + if not any(_type_ok(value, name) for name in types): + problems.append(f"{label}: must be {' or '.join(types)}") + return + if "const" in node and value != node["const"]: + problems.append(f"{label}: must be {json.dumps(node['const'])}") + return + if "enum" in node and value not in node["enum"]: + problems.append(f"{label}: must be one of {', '.join(json.dumps(item) for item in node['enum'])}") + return + if _type_ok(value, "integer"): + if "minimum" in node and value < node["minimum"]: + problems.append(f"{label}: must be at least {node['minimum']}") + if "maximum" in node and value > node["maximum"]: + problems.append(f"{label}: must be at most {node['maximum']}") + if isinstance(value, str) and "pattern" in node and not re.search(node["pattern"], value): + problems.append(f"{label}: has an invalid format") + check = annotation(node, "check") + if check and isinstance(value, str) and not CHECKS[check][0](value): + problems.append(f"{label}: {CHECKS[check][1]}") + if isinstance(value, list): + if len(value) < node.get("minItems", 0): + problems.append(f"{label}: needs at least {node['minItems']} item(s)") + if node.get("uniqueItems") and len({json.dumps(item, sort_keys=True) for item in value}) != len(value): + problems.append(f"{label}: lists an item twice") + for index, item in enumerate(value): + _validate(node.get("items", {}), item, f"{label}[{index}]", mode, problems) + if isinstance(value, dict): + properties = node.get("properties", {}) + for name in node.get("required", []): + if name not in value: + problems.append(f"{key + '.' if key else ''}{name}: is required") + for name, item in value.items(): + child = f"{key}.{name}" if key else name + if name in properties: + if mode not in annotation(properties[name], "modes", MODES): + problems.append(f'{child}: does not apply when mode is "{mode}"; remove it') + else: + _validate(properties[name], item, child, mode, problems) + elif isinstance(node.get("additionalProperties"), dict): + _validate(node["additionalProperties"], item, child, mode, problems) + else: + problems.append(f"{child}: unknown key") + + +def _complete(node, value, mode): + for name, child in node.get("properties", {}).items(): + if mode not in annotation(child, "modes", MODES) or not annotation(child, "setting", True): + continue + if name not in value: + if "default" in child: + value[name] = copy.deepcopy(child["default"]) + elif "properties" in child: + value[name] = {} + else: + continue + if isinstance(value[name], dict) and "properties" in child: + _complete(child, value[name], mode) + + +def validate(config): + """Return config with defaults filled in for every applicable key, or raise ConfigError.""" + if not isinstance(config, dict): + raise ConfigError(["config.json: must be a JSON object"]) + mode = config.get("mode") + if mode not in MODES: + raise ConfigError([f"mode: must be one of {', '.join(MODES)}"]) + problems = [] + _validate(SCHEMA, config, "", mode, problems) + if problems: + raise ConfigError(problems) + full = copy.deepcopy(config) + _complete(SCHEMA, full, mode) + native = full.get("native_core", False) + ports = full.get("ports", {}) + if mode != "web-only" and native != ("database" in ports): + problems.append("ports.database: required exactly when native_core is true" + if native else "ports.database: applies only with native_core; remove it") + if mode == "web-only" and "core_url" not in full.get("web", {}): + problems.append("web.core_url: required for a web-only installation") + if len(set(ports.values())) != len(ports): + problems.append("ports: " + ", ".join(sorted(ports)) + " need different ports") + core = full.get("core") + if core and core["default_harness"] not in core["harnesses"]: + problems.append("core.default_harness: must be listed in core.harnesses") + if problems: + raise ConfigError(problems) + return ordered(full) + + +def ordered(config, node=None): + """The same content with keys in schema order, so written files read like the reference.""" + node = SCHEMA if node is None else node + if not isinstance(config, dict) or "properties" not in node: + return config + result = {} + for name, child in node["properties"].items(): + if name in config: + result[name] = ordered(config[name], child) + for name, value in config.items(): + result.setdefault(name, value) + return result + + +def initial(mode, native_core=False, **values): + """Every applicable field for a new installation, seeded from installer flags.""" + config = {"$schema": "generated/config.schema.json", "format": 1, "mode": mode} + if mode != "web-only": + config["native_core"] = native_core + for key, value in values.items(): + if value is None: + continue + target = config + *parents, last = key.split(".") + for part in parents: + target = target.setdefault(part, {}) + target[last] = value + return validate(config) + + +def is_setting(key, node, modes, config): + return (config["mode"] in modes and annotation(node, "setting", True) + and (key != "ports.database" or config.get("native_core", False))) + + +def values(config): + """Every applicable setting of a validated config, by dotted key.""" + return {key: lookup(config, key) for key, node, modes in leaves() if is_setting(key, node, modes, config)} + + +def settings(config): + """The non-secret snapshot Core serves at GET /core/v1/installation.""" + mode = config["mode"] + items = [] + for key, node, modes in leaves(): + if not is_setting(key, node, modes, config): + continue + value = lookup(config, key) + sensitive = annotation(node, "sensitive", False) + item = {"key": key, "value": None if sensitive else value} + if sensitive: + item["configured"] = bool(value) + item.update({"default": node.get("default"), "changeable": annotation(node, "changeable", True), + "sensitive": sensitive, + "restarts": [name for name in annotation(node, "restarts", []) if name in SERVICES[mode]]}) + items.append(item) + return items + + +def sensitive_keys(): + return [key for key, node, _ in leaves() if annotation(node, "sensitive", False)] diff --git a/deploy/install/configuration.py b/deploy/install/configuration.py index aa09530ff..a7be2c78e 100644 --- a/deploy/install/configuration.py +++ b/deploy/install/configuration.py @@ -1,14 +1,29 @@ -"""Deployment files for the existing Core, Runtime and production console.""" -import os -import re -import stat +"""Derive every file under generated/ from config.json, state.json and secrets/. + +Generated files hold no secret except runtime-history.json, which carries the +operator's export headers. Secrets stay in secrets/, one copy each, and reach the +services as read-only single-file mounts or file paths. +""" +import hashlib +import json from pathlib import Path -from urllib.parse import urlsplit +import re +from urllib.parse import urlencode + +import config_model +import native_service + +# Where Core and Web containers see secrets and generated inputs. +RUN = "/run/parsar" +POOL = (("max_conns", "pool_max_conns"), ("min_conns", "pool_min_conns"), + ("max_conn_lifetime", "pool_max_conn_lifetime"), ("max_conn_idle_time", "pool_max_conn_idle_time"), + ("health_check_period", "pool_health_check_period")) -def environment_text(values): +def environment_text(values, header): """The shared Compose/systemd subset: quoted, single-line literal values.""" - lines = ["# Core process configuration. Escape backslash, double quote and dollar with backslash.\n"] + lines = ["# " + header + "\n", + "# Values are literal. Escape backslash, double quote and dollar with backslash.\n"] for key, value in values.items(): if (not isinstance(key, str) or not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", key) or not isinstance(value, str) or any(char in value for char in "\x00\r\n")): @@ -18,38 +33,16 @@ def environment_text(values): return "".join(lines) -def read_core_environment(root, state=None): - """Read the persisted authority without shell or ambient-variable expansion.""" - root = Path(root) - path = root / "config/core.env" - failure = "Core configuration is missing or unsafe; restore the private config/core.env file" - try: - directory = path.parent.lstat() - if (not stat.S_ISDIR(directory.st_mode) or stat.S_IMODE(directory.st_mode) & 0o077 - or directory.st_uid != os.geteuid()): - raise RuntimeError(failure) - descriptor = os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK) - with os.fdopen(descriptor, encoding="utf-8") as stream: - info = os.fstat(stream.fileno()) - if (not stat.S_ISREG(info.st_mode) or stat.S_IMODE(info.st_mode) & 0o077 - or info.st_uid != os.geteuid() or info.st_nlink != 1): - raise RuntimeError(failure) - content = stream.read() - except (OSError, UnicodeError): - raise RuntimeError(failure) from None +def read_environment(text): + """Parse the environment_text format without shell or variable expansion.""" result = {} - for line in content.split("\n"): + for line in text.split("\n"): if not line.strip() or line.lstrip().startswith("#"): continue match = re.fullmatch(r'([A-Za-z_][A-Za-z0-9_]*)="((?:[^\\"$\x00\r\n]|\\[\\"$])*)"', line) if not match or match[1] in result: - raise RuntimeError("Core configuration requires unique names and double-quoted literal values") + raise RuntimeError("An environment file requires unique names and double-quoted literal values") result[match[1]] = re.sub(r'\\([\\"$])', r'\1', match[2]) - if state is not None: - expected = str(path) if state["native_core"] else "/config/core.env" - if (result.get("AGENTS_API_SANDBOX_INSTALLATION_ID") != state["installation_id"] - or result.get("AGENTS_API_CONFIG_FILE") != expected): - raise RuntimeError("Core configuration installation identity or file path differs") return result @@ -57,81 +50,194 @@ def bind(source, target, readonly=True): return {"type": "bind", "source": str(source).replace("$", "$$"), "target": target, "read_only": readonly} -def core_environment(root, state, database_password): - native = state["native_core"] - 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'127.0.0.1:{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() +def sha256(data): + return hashlib.sha256(data.encode() if isinstance(data, str) else data).hexdigest() + + +def edit_hint(root): + return f"Generated from {Path(root) / 'config.json'}. Do not edit; change config.json and run {Path(root) / 'parsar'} apply." + + +def read_core_key(root): + """The Core key, checked the way Web checks it.""" + key = (Path(root) / "secrets/core.key").read_text().strip() + if len(key) < 32 or any(char.isspace() or char == "\x00" for char in key): + raise RuntimeError("secrets/core.key must hold one Core key of at least 32 characters without whitespace") + return key + + +def secret_digests(root, mode): + names = ["core.key"] if mode == "web-only" else ["core.key", "credential.key", "database.password"] + return {name: sha256((Path(root) / "secrets" / name).read_bytes()) for name in names} + + +def local_public_url(config): + """AGENTS_API_PUBLIC_URL: the public origin, or Core's loopback origin for local use.""" + return config["public_url"] or f'http://127.0.0.1:{config["ports"]["core"]}' + + +def log_environment(log): + result = {"PARSAR_LOG_LEVEL": log["level"]} + if log["format"] != "auto": + result["PARSAR_LOG_FORMAT"] = log["format"] + if log["add_source"]: + result["PARSAR_LOG_ADD_SOURCE"] = "1" + return result + + +def core_environment(root, config, state): + root = Path(root) + native = config["native_core"] + generated = str(root / "generated") if native else RUN + secrets = str(root / "secrets") if native else RUN + ports, core = config["ports"], config["core"] + database = f'127.0.0.1:{ports["database"]}' if native else "database:5432" + query = [("sslmode", "disable")] + [(name, str(core["database_pool"][key])) for key, name in POOL + if core["database_pool"][key] is not None] result = { - "AGENTS_API_DATABASE_URL": f"postgres://agents_api:{database_password}@{database}/agents_api?sslmode=disable", - "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": daemon_url, - "AGENTS_API_E2B_PROVIDER_BIN": (str(Path(root) / "native/e2b/agents-api-e2b-provider") - if native else "/opt/parsar/e2b/agents-api-e2b-provider"), - "AGENTS_API_E2B_STATE_DIR": str(Path(root) / "state/e2b") if native else "/state/e2b", + "AGENTS_API_ADDR": f'127.0.0.1:{ports["core"]}' if native else ":8091", + "AGENTS_API_PUBLIC_URL": local_public_url(config), + "AGENTS_API_DATABASE_URL": f"postgres://agents_api@{database}/agents_api?" + urlencode(query), + "AGENTS_API_DATABASE_PASSWORD_FILE": secrets + "/database.password", + "AGENTS_API_CREDENTIAL_KEY_FILE": secrets + "/credential.key", + "AGENTS_API_CORE_KEY_DIGESTS_FILE": generated + "/core-key-digests.json", + "AGENTS_API_SANDBOX_INSTALLATION_ID": state["installation_id"], + "AGENTS_API_SETTINGS_FILE": generated + "/settings.json", + "AGENTS_API_E2B_STATE_DIR": str(root / "state/e2b") if native else "/state/e2b", + "AGENTS_API_ENGINE": core["default_harness"], + "AGENTS_API_HARNESSES": ",".join(core["harnesses"]), + "AGENTS_API_EXECUTION_CONCURRENCY": str(core["execution_concurrency"]), + "AGENTS_API_WRITE_AUDIT_RETENTION": core["write_audit_retention"], } - result["AGENTS_API_CORE_KEY_DIGESTS_FILE"] = ( - str(Path(root) / "admin/core-key-digests.json") if native else "/admin/core-key-digests.json") - result["AGENTS_API_SANDBOX_INSTALLATION_ID"] = state["installation_id"] - result["AGENTS_API_CONFIG_FILE"] = config + "/core.env" + if native: + result["AGENTS_API_E2B_PROVIDER_BIN"] = str(root / "native/e2b/agents-api-e2b-provider") + if core["oauth_trusted_origins"]: + result["AGENTS_API_OAUTH_TRUSTED_ORIGINS"] = ",".join(core["oauth_trusted_origins"]) + if core["runtime_history"] is not None: + result["AGENTS_API_RUNTIME_HISTORY_FILE"] = generated + "/runtime-history.json" + retained = state.get("execution_options_file") + if retained: + # Kept by --convert until phase 3's one-time import moves it into Core. + result["AGENTS_API_EXECUTION_OPTIONS_FILE"] = retained["variable"] + result.update(log_environment(config["log"])) return result -def compose_config(root, state, manifest, database_password): +def settings_document(root, config, applied_at): + root = Path(root) + return {"path": str(root / "config.json"), "apply_command": f"{root / 'parsar'} apply", + "applied_at": applied_at, "settings": config_model.settings(config)} + + +def compose_config(root, config, state, inputs): root = Path(root) - config = root / "config" + mode, native = config["mode"], config.get("native_core", False) identity = f'{state["uid"]}:{state["gid"]}' - doc = {"name": state["project"], "services": {}} + images = state["images"] + doc = {"name": state["project"], + "x-parsar": {"generated_from": str(root / "config.json"), "edit": "config.json, then parsar apply"}, + "services": {}} services = doc["services"] - native = state["native_core"] and state["mode"] != "web-only" - if state["mode"] != "web-only": + if mode != "web-only": services["database"] = { - "image": manifest["images"]["database"], "restart": "unless-stopped", + "image": images["database"], "restart": "unless-stopped", "environment": {"POSTGRES_USER": "agents_api", "POSTGRES_DB": "agents_api", - "POSTGRES_PASSWORD": database_password}, - "volumes": ["database:/var/lib/postgresql/data"], + "POSTGRES_PASSWORD_FILE": "/run/secrets/database.password"}, + "volumes": ["database:/var/lib/postgresql/data", + bind(root / "secrets/database.password", "/run/secrets/database.password")], "healthcheck": {"test": ["CMD-SHELL", "pg_isready -U agents_api -d agents_api"], "interval": "2s", "timeout": "5s", "retries": 30}, } - shared = {"image": manifest["images"]["core"], "user": identity, - "env_file": [str(config / "core.env").replace("$", "$$")], "volumes": [bind(config, "/config")], - "read_only": True, "tmpfs": ["/tmp:mode=1777"], "init": True, - "security_opt": ["no-new-privileges:true"]} - services["migrate"] = dict(shared, command=["/usr/local/bin/agents-api-migrate"], - depends_on={"database": {"condition": "service_healthy"}}) - core = dict(shared, restart="unless-stopped", ports=[f'127.0.0.1:{state["core_port"]}:8091'], - depends_on={"migrate": {"condition": "service_completed_successfully"}}) - core["volumes"] = list(shared["volumes"]) + doc["volumes"] = {"database": {}} if native: - services["database"]["ports"] = [f'127.0.0.1:{state["database_port"]}:5432'] - services.pop("migrate") + services["database"]["ports"] = [f'127.0.0.1:{config["ports"]["database"]}:5432'] else: - core["volumes"].append(bind(root / "admin/core-key-digests.json", "/admin/core-key-digests.json")) - core["volumes"].append(bind(root / "state/e2b", "/state/e2b", False)) - services["core"] = core - doc["volumes"] = {"database": {}} - if state["mode"] != "core-only": - services["web"] = { - "image": manifest["images"]["web"], "user": identity, "restart": "unless-stopped", - "ports": [f'127.0.0.1:{state["web_port"]}:8080'], "read_only": True, - "security_opt": ["no-new-privileges:true"], - "volumes": [bind(root / "admin/core.key", "/admin/core.key")], - "environment": {"CORE_CONSOLE_ORIGIN": state.get("public_url") or f'http://127.0.0.1:{state["web_port"]}', - "CORE_CONSOLE_UPSTREAM": (f'http://127.0.0.1:{state["core_port"]}' if native - else state.get("core_url") or "http://core:8091"), - "CORE_CONSOLE_CORE_KEY_FILE": "/admin/core.key"}, + mounts = [bind(root / "secrets" / name, f"{RUN}/{name}") for name in ("credential.key", "database.password")] + mounts += [bind(root / "generated" / name, f"{RUN}/{name}") for name in ("core-key-digests.json", "settings.json")] + if config["core"]["runtime_history"] is not None: + mounts.append(bind(root / "generated/runtime-history.json", f"{RUN}/runtime-history.json")) + retained = state.get("execution_options_file") + if retained and not retained["variable"].startswith("/state/e2b/"): + mounts.append(bind(retained["path"], retained["variable"])) + shared = {"image": images["core"], "user": identity, + "env_file": [str(root / "generated/core.env").replace("$", "$$")], "volumes": mounts, + "read_only": True, "tmpfs": ["/tmp:mode=1777"], "init": True, + "security_opt": ["no-new-privileges:true"], "labels": {"io.parsar.inputs": inputs["core"]}} + services["migrate"] = dict(shared, command=["/usr/local/bin/agents-api-migrate"], + depends_on={"database": {"condition": "service_healthy"}}) + services["core"] = dict(shared, restart="unless-stopped", ports=[f'127.0.0.1:{config["ports"]["core"]}:8091'], + depends_on={"migrate": {"condition": "service_completed_successfully"}}, + volumes=mounts + [bind(root / "state/e2b", "/state/e2b", False)]) + if mode != "core-only": + if mode == "web-only": + upstream = config["web"]["core_url"] + else: + upstream = f'http://127.0.0.1:{config["ports"]["core"]}' if native else "http://core:8091" + environment = { + "CORE_CONSOLE_ORIGIN": config["public_url"] or f'http://127.0.0.1:{config["ports"]["web"]}', + "CORE_CONSOLE_UPSTREAM": upstream, + "CORE_CONSOLE_CORE_KEY_FILE": f"{RUN}/core.key", + "CORE_CONSOLE_NODE_PAYLOAD_DIR": "/node-payload", } - services["web"]["volumes"].append(bind(root / "node-payload", "/node-payload")) - services["web"]["environment"]["CORE_CONSOLE_NODE_PAYLOAD_DIR"] = "/node-payload" - if state["mode"] == "web-only" or native: - services["web"].pop("ports") - services["web"]["network_mode"] = "host" - services["web"]["environment"]["CORE_CONSOLE_ADDR"] = f'127.0.0.1:{state["web_port"]}' + environment.update(log_environment(config["log"])) + web = {"image": images["web"], "user": identity, "restart": "unless-stopped", + "ports": [f'127.0.0.1:{config["ports"]["web"]}:8080'], "read_only": True, + "security_opt": ["no-new-privileges:true"], + "volumes": [bind(root / "secrets/core.key", f"{RUN}/core.key"), bind(root / "node-payload", "/node-payload")], + "environment": environment, "labels": {"io.parsar.inputs": inputs["web"]}} + if mode == "web-only" or native: + web.pop("ports") + web["network_mode"] = "host" + environment["CORE_CONSOLE_ADDR"] = f'127.0.0.1:{config["ports"]["web"]}' + services["web"] = web return doc + + +class Rendered: + """Generated file contents plus one digest per service for restart planning.""" + + def __init__(self, files, services, unit): + self.files, self.services, self.unit = files, services, unit + + +def render(root, config, state, applied_at): + root = Path(root) + mode, native = config["mode"], config.get("native_core", False) + secrets = secret_digests(root, mode) + settings = settings_document(root, config, applied_at) + files = {"config.schema.json": config_model.SCHEMA_TEXT, + "settings.json": json.dumps(settings, indent=2) + "\n"} + inputs, core_env, unit = {}, "", None + if mode != "web-only": + files["core-key-digests.json"] = json.dumps([sha256(read_core_key(root))]) + "\n" + history = config["core"]["runtime_history"] + if history is not None: + files["runtime-history.json"] = json.dumps(history, indent=2) + "\n" + core_env = environment_text(core_environment(root, config, state), edit_hint(root)) + files["core.env"] = core_env + # Only settings Core itself restarts for enter its inputs, so a Web-only change + # leaves Core running. Its snapshot then refreshes on Core's next restart. + core_settings = [item for item in settings["settings"] if "core" in item["restarts"]] + retained = state.get("execution_options_file") + inputs["core"] = sha256(json.dumps({ + "core-key-digests.json": sha256(files["core-key-digests.json"]), + "settings": sha256(json.dumps([settings["path"], settings["apply_command"], core_settings], sort_keys=True)), + "runtime-history.json": sha256(files.get("runtime-history.json", "")), + "credential.key": secrets["credential.key"], "database.password": secrets["database.password"], + "execution-options": sha256(Path(retained["path"]).read_bytes()) if retained else None, + }, sort_keys=True)) + if native: + unit = native_service.unit_name(state) + files[unit] = native_service.unit_text(root, edit_hint(root)) + if mode != "core-only": + inputs["web"] = sha256(json.dumps({"core.key": secrets["core.key"]})) + compose = compose_config(root, config, state, inputs) + files["compose.json"] = json.dumps(compose, indent=2) + "\n" + services = {} + for name, service in compose["services"].items(): + # Compose resolves env_file into the service configuration, so its content counts. + text = json.dumps(service, sort_keys=True) + (core_env if "env_file" in service else "") + services[name] = sha256(text) + if native: + services["core"] = sha256(core_env + files[unit] + inputs["core"]) + return Rendered(files, services, unit) diff --git a/deploy/install/convert.py b/deploy/install/convert.py new file mode 100644 index 000000000..aebbede88 --- /dev/null +++ b/deploy/install/convert.py @@ -0,0 +1,479 @@ +"""install.sh --convert: move an installation made before config.json to the new layout. + +Conversion is also an upgrade to this bundle's release, because the earlier Core does +not read the new environment names. The preflight changes no file: it builds +config.json and state.json in memory, and anything it can't convert stops it with +"nothing was changed". An interrupted conversion resumes while installation.json and +config.json both exist. +""" +import datetime +import hashlib +import json +import os +from pathlib import Path +import re +import stat +import subprocess +import sys +from urllib.parse import parse_qsl, urlsplit + +import config_model +import configuration +import native_service +import parsar_cli + +POOL = {name: key for key, name in configuration.POOL} +MAPPED = {"AGENTS_API_EXECUTION_CONCURRENCY", "PARSAR_LOG_LEVEL", "PARSAR_LOG_FORMAT", "PARSAR_LOG_ADD_SOURCE", + "AGENTS_API_WRITE_AUDIT_RETENTION", "AGENTS_API_ENGINE", "AGENTS_API_HARNESSES", + "AGENTS_API_OAUTH_TRUSTED_ORIGINS", "AGENTS_API_RUNTIME_HISTORY_FILE", "AGENTS_API_EXECUTION_OPTIONS_FILE", + "AGENTS_API_DATABASE_URL"} +DERIVED = ("AGENTS_API_ADDR", "AGENTS_API_CREDENTIAL_KEY_FILE", "AGENTS_API_CORE_KEY_DIGESTS_FILE", + "AGENTS_API_SANDBOX_INSTALLATION_ID", "AGENTS_API_CONFIG_FILE", "AGENTS_API_E2B_PROVIDER_BIN", + "AGENTS_API_E2B_STATE_DIR", "AGENTS_API_DAEMON_WS_URL") + + +class ConvertError(Exception): + pass + + +# The installer's generator before config.json (deploy/install/configuration.py at +# 5c3dcc16, unchanged since #124). Conversion compares the retained files with it. + +def legacy_core_environment(root, state, database_password): + native = state["native_core"] + 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'127.0.0.1:{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() + return { + "AGENTS_API_DATABASE_URL": f"postgres://agents_api:{database_password}@{database}/agents_api?sslmode=disable", + "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": daemon_url, + "AGENTS_API_E2B_PROVIDER_BIN": (str(Path(root) / "native/e2b/agents-api-e2b-provider") + if native else "/opt/parsar/e2b/agents-api-e2b-provider"), + "AGENTS_API_E2B_STATE_DIR": str(Path(root) / "state/e2b") if native else "/state/e2b", + "AGENTS_API_CORE_KEY_DIGESTS_FILE": (str(Path(root) / "admin/core-key-digests.json") if native + else "/admin/core-key-digests.json"), + "AGENTS_API_SANDBOX_INSTALLATION_ID": state["installation_id"], + "AGENTS_API_CONFIG_FILE": config + "/core.env", + } + + +def legacy_compose_config(root, state, images, database_password): + root = Path(root) + config = root / "config" + identity = f'{state["uid"]}:{state["gid"]}' + doc = {"name": state["project"], "services": {}} + services = doc["services"] + native = state["native_core"] and state["mode"] != "web-only" + bind = configuration.bind + if state["mode"] != "web-only": + services["database"] = { + "image": images.get("database"), "restart": "unless-stopped", + "environment": {"POSTGRES_USER": "agents_api", "POSTGRES_DB": "agents_api", + "POSTGRES_PASSWORD": database_password}, + "volumes": ["database:/var/lib/postgresql/data"], + "healthcheck": {"test": ["CMD-SHELL", "pg_isready -U agents_api -d agents_api"], + "interval": "2s", "timeout": "5s", "retries": 30}, + } + shared = {"image": images.get("core"), "user": identity, + "env_file": [str(config / "core.env").replace("$", "$$")], "volumes": [bind(config, "/config")], + "read_only": True, "tmpfs": ["/tmp:mode=1777"], "init": True, + "security_opt": ["no-new-privileges:true"]} + services["migrate"] = dict(shared, command=["/usr/local/bin/agents-api-migrate"], + depends_on={"database": {"condition": "service_healthy"}}) + core = dict(shared, restart="unless-stopped", ports=[f'127.0.0.1:{state["core_port"]}:8091'], + depends_on={"migrate": {"condition": "service_completed_successfully"}}) + core["volumes"] = list(shared["volumes"]) + if native: + services["database"]["ports"] = [f'127.0.0.1:{state["database_port"]}:5432'] + services.pop("migrate") + else: + core["volumes"].append(bind(root / "admin/core-key-digests.json", "/admin/core-key-digests.json")) + core["volumes"].append(bind(root / "state/e2b", "/state/e2b", False)) + services["core"] = core + doc["volumes"] = {"database": {}} + if state["mode"] != "core-only": + services["web"] = { + "image": images.get("web"), "user": identity, "restart": "unless-stopped", + "ports": [f'127.0.0.1:{state["web_port"]}:8080'], "read_only": True, + "security_opt": ["no-new-privileges:true"], + "volumes": [bind(root / "admin/core.key", "/admin/core.key")], + "environment": {"CORE_CONSOLE_ORIGIN": state.get("public_url") or f'http://127.0.0.1:{state["web_port"]}', + "CORE_CONSOLE_UPSTREAM": (f'http://127.0.0.1:{state["core_port"]}' if native + else state.get("core_url") or "http://core:8091"), + "CORE_CONSOLE_CORE_KEY_FILE": "/admin/core.key"}, + } + services["web"]["volumes"].append(bind(root / "node-payload", "/node-payload")) + services["web"]["environment"]["CORE_CONSOLE_NODE_PAYLOAD_DIR"] = "/node-payload" + if state["mode"] == "web-only" or native: + services["web"].pop("ports") + services["web"]["network_mode"] = "host" + services["web"]["environment"]["CORE_CONSOLE_ADDR"] = f'127.0.0.1:{state["web_port"]}' + return doc + + +def legacy_read_environment(root): + """The earlier installer's strict reader: a private file of literal values.""" + path = Path(root) / "config/core.env" + directory = path.parent.lstat() + if (not stat.S_ISDIR(directory.st_mode) or stat.S_IMODE(directory.st_mode) & 0o077 + or directory.st_uid != os.geteuid()): + raise RuntimeError("config/ must be a private directory") + parsar_cli.check_private(path, "config/core.env") + return configuration.read_environment(path.read_text()) + + +# Preflight ------------------------------------------------------------------- + +def json_differences(expected, actual, path=""): + if isinstance(expected, dict) and isinstance(actual, dict): + result = [] + for key in sorted(set(expected) | set(actual), key=str): + result += json_differences(expected.get(key), actual.get(key), f"{path}.{key}" if path else str(key)) + return result + if isinstance(expected, list) and isinstance(actual, list) and len(expected) == len(actual): + result = [] + for index, (left, right) in enumerate(zip(expected, actual)): + result += json_differences(left, right, f"{path}[{index}]") + return result + return [] if expected == actual else [path] + + +def host_path(root, old, variable): + """Where a file named in the old core.env lives on this host.""" + if old["native_core"]: + return Path(variable) + for prefix, directory in (("/config/", root / "config"), ("/state/e2b/", root / "state/e2b")): + if variable.startswith(prefix) and "/" not in variable[len(prefix):]: + return directory / variable[len(prefix):] + return None + + +def detect(root): + """The old installation record, or ConvertError naming what is missing.""" + try: + old = json.loads((root / "installation.json").read_text()) + except (OSError, ValueError): + raise ConvertError("installation.json can't be read; this installation can't be converted") from None + if old.get("version") != 1: + raise ConvertError("installation.json has an unknown version; this installation can't be converted") + required = ["compose.json", "admin/core.key"] + if old["mode"] != "web-only": + required += ["admin/core-key-digests.json", "config/core.env", "config/credential.key", + "config/database.password"] + missing = [name for name in required if not (root / name).is_file()] + if missing: + raise ConvertError("This installation predates the supported layout and can't be converted; missing: " + + ", ".join(missing)) + return old + + +class Plan: + def __init__(self): + self.problems, self.notes, self.leftovers = [], [], [] + self.config = self.state = None + self.moves, self.deletions = [], [] + + def problem(self, where, key, reason): + self.problems.append(f"{where}{': ' + key if key else ''}: {reason}") + + +def old_core_deployment(root, old, plan, run): + """Read the sandbox deployment's core_url from the old Core, starting it with its old files.""" + base = f'http://127.0.0.1:{old["core_port"]}' + started = False + if parsar_cli.http(base + "/healthz")[0] != 200: + run(["docker", "compose", "-f", str(root / "compose.json"), "up", "--detach", "--wait"]) + if native_service.is_native(old): + run(["systemctl", "--user", "start", native_service.unit_name(old)]) + started = True + parsar_cli.wait_status(base + "/healthz") + key = (root / "admin/core.key").read_text().strip() + status, body = parsar_cli.http(base + "/core/v1/sandbox/deployment", parsar_cli.bearer(key)) + if status != 200: + plan.problem("sandbox deployment", "core_url", "the old Core did not return it; make sure it starts with its " + "own files and that admin/core.key is current") + return None, started + return json.loads(body).get("core_url") or "", started + + +def preflight(root, bundle_manifest, images, public_url_override, run): + """Build config.json and state.json in memory; change nothing.""" + old = detect(root) + plan = Plan() + mode = old["mode"] + native = bool(old["native_core"]) and mode != "web-only" + if (root / "config/managed-runtimes.json").exists(): + plan.problem("config/managed-runtimes.json", "", "retired file-managed provider configuration; preserve its " + "resources and follow the deployment replacement guide") + password = (root / "config/database.password").read_text() if mode != "web-only" else "" + compose = json.loads((root / "compose.json").read_text()) + own = {("core" if name == "migrate" else name): service.get("image") + for name, service in compose.get("services", {}).items()} + for key in json_differences(legacy_compose_config(root, old, own, password), compose): + plan.problem("compose.json", key, "differs from the file the installer generated; undo the edit") + + config = {"$schema": "generated/config.schema.json", "format": 1, "mode": mode} + if mode != "web-only": + config["native_core"] = native + ports = {} + if mode != "web-only": + ports["core"] = old["core_port"] + if mode != "core-only": + ports["web"] = old["web_port"] + if native: + ports["database"] = old["database_port"] + config["ports"] = ports + if mode == "web-only": + config["web"] = {"core_url": old.get("core_url")} + config["log"] = {"level": "info", "format": "auto", "add_source": False} + retained = None + known = {"config": set(), "admin": {"core.key"}} + if mode != "web-only": + known["config"] = {"core.env", "credential.key", "database.password"} + known["admin"].add("core-key-digests.json") + if native: + known["config"].add(native_service.unit_name(old)) + try: + env = legacy_read_environment(root) + except (RuntimeError, parsar_cli.ParsarError, OSError, UnicodeError): + plan.problem("config/core.env", "", "must be a private file of double-quoted literal values") + env = None + if env is not None: + retained = map_environment(root, old, env, password, config, plan, known) + try: + digests = json.loads((root / "admin/core-key-digests.json").read_text()) + except ValueError: + digests = None + current = hashlib.sha256((root / "admin/core.key").read_text().strip().encode()).hexdigest() + if not isinstance(digests, list) or current not in digests: + plan.problem("admin/core-key-digests.json", "", "does not authorize admin/core.key; restore it") + elif len(digests) > 1: + plan.notes.append(f"admin/core-key-digests.json lists {len(digests) - 1} more Core key digest(s). " + "They are dropped; only admin/core.key works after conversion.") + + public = old.get("public_url") + started = False + if mode != "web-only" and not plan.problems: + deployment, started = old_core_deployment(root, old, plan, run) + if deployment: + if public is None: + public = deployment + plan.notes.append(f"public_url {deployment} is adopted from the sandbox deployment, whose nodes use it.") + elif public != deployment: + if public_url_override in (public, deployment): + if public_url_override == public: + plan.notes.append(f"public_url {public} differs from the deployment's {deployment}: its nodes " + "get no new sandboxes; remove them in Web and add them again.") + public = public_url_override + public_url_override = None + else: + plan.problem("public_url", "", f"the installation's {public} differs from the sandbox deployment's " + f"{deployment}. Rerun with --public-url {public} or --public-url {deployment}; " + f"choosing {public} means the deployment's nodes must be removed and added again") + if public_url_override is not None and not plan.problems: + plan.problem("--public-url", "", "only settles a conflict between the installation and the sandbox deployment; " + "change public_url in config.json after conversion") + config["public_url"] = public + try: + plan.config = config_model.validate(config) + except config_model.ConfigError as error: + for item in error.problems: + plan.problem("config.json", "", item) + + for directory in ("config", "admin"): + path = root / directory + if path.is_dir(): + for entry in sorted(path.iterdir()): + if entry.name not in known[directory]: + plan.leftovers.append(f"{directory}/{entry.name}") + plan.moves = [("admin/core.key", "secrets/core.key")] + plan.deletions = ["compose.json"] + if mode != "web-only": + plan.moves += [("config/credential.key", "secrets/credential.key"), + ("config/database.password", "secrets/database.password")] + plan.deletions += ["config/core.env", "admin/core-key-digests.json"] + if native: + plan.deletions.append("config/" + native_service.unit_name(old)) + plan.state = { + "format": 1, "installation_id": old["installation_id"], "project": old["project"], + "uid": old["uid"], "gid": old["gid"], "mode": mode, "native_core": native, + "source_commit": bundle_manifest["source_commit"], "images": images, + "secrets_sha256": {Path(target).name: configuration.sha256((root / source).read_bytes()) + for source, target in plan.moves}, + "core_installation_id": None, "execution_options_file": retained, "applied": None, + "converted_from": {"source_commit": old["source_commit"], + "at": datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")}} + return old, plan, started + + +def map_environment(root, old, env, password, config, plan, known): + expected = legacy_core_environment(root, old, password) + for name in DERIVED: + if env.get(name) != expected[name]: + plan.problem("config/core.env", name, "differs from the value the installer generated; restore it") + for name in sorted(set(env) - set(DERIVED) - MAPPED): + plan.problem("config/core.env", name, "is not a setting config.json can hold; remove it") + base, _, query = env.get("AGENTS_API_DATABASE_URL", "").partition("?") + pool = {} + parameters = parse_qsl(query, keep_blank_values=True) + if base + "?sslmode=disable" != expected["AGENTS_API_DATABASE_URL"] or ("sslmode", "disable") not in parameters: + plan.problem("config/core.env", "AGENTS_API_DATABASE_URL", "names another database, user or password; " + "only the installation's own PostgreSQL can be converted") + for name, value in parameters: + if name == "sslmode": + continue + if name not in POOL or POOL[name] in pool: + plan.problem("config/core.env", "AGENTS_API_DATABASE_URL", f"query parameter {name} can't be converted") + continue + pool[POOL[name]] = int(value) if value.isdigit() and POOL[name].endswith("conns") else value + core = {"database_pool": {key: pool.get(key) for key, _ in configuration.POOL}} + concurrency = env.get("AGENTS_API_EXECUTION_CONCURRENCY", "4") + core["execution_concurrency"] = int(concurrency) if concurrency.isdigit() else concurrency + engine = env.get("AGENTS_API_ENGINE") or "codex" + enabled = {item.strip() for item in env.get("AGENTS_API_HARNESSES", "").split(",") if item.strip()} | {engine} + core["harnesses"], core["default_harness"] = sorted(enabled), engine + core["write_audit_retention"] = env.get("AGENTS_API_WRITE_AUDIT_RETENTION") or "2160h" + core["oauth_trusted_origins"] = [item.strip() for item in env.get("AGENTS_API_OAUTH_TRUSTED_ORIGINS", "").split(",") + if item.strip()] + core["runtime_history"] = None + if env.get("AGENTS_API_RUNTIME_HISTORY_FILE"): + path = host_path(root, old, env["AGENTS_API_RUNTIME_HISTORY_FILE"]) + try: + core["runtime_history"] = json.loads(path.read_text()) + plan.notes.append(f"Runtime history settings from {path} are now in config.json. The file is kept; " + "delete it if it holds export credentials.") + if path.parent == root / "config": + known["config"].add(path.name) + except (AttributeError, OSError, ValueError): + plan.problem("config/core.env", "AGENTS_API_RUNTIME_HISTORY_FILE", "names a file that can't be read as JSON") + retained = None + if env.get("AGENTS_API_EXECUTION_OPTIONS_FILE"): + variable = env["AGENTS_API_EXECUTION_OPTIONS_FILE"] + path = host_path(root, old, variable) + if path is None or path.is_symlink() or not path.is_file(): + plan.problem("config/core.env", "AGENTS_API_EXECUTION_OPTIONS_FILE", "names a file that can't be found") + else: + retained = {"variable": variable, "path": str(path)} + if path.parent == root / "config": + known["config"].add(path.name) + plan.notes.append(f"AGENTS_API_EXECUTION_OPTIONS_FILE and {path} stay as they are; config.json does not " + "hold model settings. A later release imports them into Core once, through parsar apply.") + log = config["log"] + level = env.get("PARSAR_LOG_LEVEL", "").strip().lower() + log["level"] = {"": "info", "warning": "warn"}.get(level, level) + log["format"] = env.get("PARSAR_LOG_FORMAT", "").strip().lower() or "auto" + add_source = env.get("PARSAR_LOG_ADD_SOURCE", "") + log["add_source"] = add_source == "1" if add_source in ("", "0", "1") else add_source + if any(name.startswith("PARSAR_LOG_") for name in env): + plan.notes.append("PARSAR_LOG_* had no effect before this release; the log settings now apply.") + config["core"] = core + return retained + + +# Steps ----------------------------------------------------------------------- + +def confirm(root, old, plan, yes, interactive, out): + out("Conversion writes config.json and state.json, moves the secrets into secrets/, removes the old generated") + out("files and upgrades Core to this release. Database migrations can't be undone; back up first:") + if old["mode"] != "web-only": + out(f' docker compose -f {root / "compose.json"} exec -T database pg_dump -U agents_api agents_api > parsar-backup.sql') + for note in plan.notes: + out("Note: " + note) + for name in plan.leftovers: + out(f"Left in place: {name} (not part of the installer's layout)") + defaults = {key: node.get("default") for key, node, _ in config_model.leaves()} + shown = {key: value for key, value in config_model.values(plan.config).items() + if key not in config_model.sensitive_keys() and (key in ("mode", "public_url") or value != defaults[key])} + out("config.json, apart from defaults: " + json.dumps(shown)) + if not (yes or (interactive and input("Type yes to convert this installation: ").strip() == "yes")): + raise ConvertError("Conversion was not confirmed; nothing was changed. Rerun with --yes to skip the prompt") + + +def stop_old(root, old, run): + if native_service.is_native(old): + native_service.disable(old) + run(["docker", "compose", "-f", str(root / "compose.json"), "stop"]) + + +def move(root, source, target): + source, target = root / source, root / target + if source.exists() and target.exists(): + if source.read_bytes() != target.read_bytes(): + raise ConvertError(f"Both {source} and {target} exist with different content; keep the right one and rerun") + source.unlink() + elif source.exists(): + os.rename(source, target) + elif not target.exists(): + raise ConvertError(f"{source} is missing; restore it and rerun --convert") + + +def write_layout(root, old, plan): + """Step 5. Every operation checks its source and target, so a rerun resumes.""" + for name in ("secrets", "generated"): + (root / name).mkdir(mode=0o700, exist_ok=True) + os.chmod(root / name, 0o700) + parsar_cli.write_private(root / "state.json", json.dumps(plan.state, indent=2) + "\n") + parsar_cli.create_private(root / "config.json", json.dumps(plan.config, indent=2) + "\n") + finish_layout(root, old, plan.moves, plan.deletions) + + +def finish_layout(root, old, moves, deletions): + for source, target in moves: + move(root, source, target) + for name in deletions: + (root / name).unlink(missing_ok=True) + for directory in ("config", "admin"): + path = root / directory + if path.is_dir() and not any(path.iterdir()): + path.rmdir() + (root / "installation.json").unlink() + + +def resume(root, out): + old = json.loads((root / "installation.json").read_text()) + moves = [("admin/core.key", "secrets/core.key")] + deletions = ["compose.json"] + if old["mode"] != "web-only": + moves += [("config/credential.key", "secrets/credential.key"), + ("config/database.password", "secrets/database.password")] + deletions += ["config/core.env", "admin/core-key-digests.json"] + if native_service.is_native(old): + deletions.append("config/" + native_service.unit_name(old)) + out("Resuming an interrupted conversion.") + finish_layout(root, old, moves, deletions) + + +def convert(root, manifest, load_images, public_url_override, yes, run, interactive=None, out=print): + """Steps 1-5. The caller then loads the bundle's files and applies (step 6).""" + interactive = sys.stdin.isatty() if interactive is None else interactive + if (root / "config.json").exists(): + load_images(json.loads((root / "state.json").read_text())["images"]) + resume(root, out) + return + old = detect(root) + wanted = ["web"] if old["mode"] == "web-only" else (["database"] if old["native_core"] else ["core", "database"]) + if old["mode"] == "all": + wanted.append("web") + old, plan, started = preflight(root, manifest, {name: None for name in wanted}, public_url_override, run) + try: + if plan.problems: + raise ConvertError("This installation can't be converted yet:\n" + + "\n".join(" - " + item for item in plan.problems) + "\nNothing was changed.") + confirm(root, old, plan, yes, interactive, out) + except ConvertError: + if started: + # Leave the services as they were found. + run(["docker", "compose", "-f", str(root / "compose.json"), "stop"]) + native_service.stop(root, old) + raise + plan.state["images"] = load_images(wanted) + stop_old(root, old, run) + write_layout(root, old, plan) diff --git a/deploy/install/install.py b/deploy/install/install.py index f8a19ca27..71df339d7 100644 --- a/deploy/install/install.py +++ b/deploy/install/install.py @@ -1,10 +1,14 @@ #!/usr/bin/env python3 -"""Install one matched Core distribution without changing execution ownership.""" +"""Install one matched Core distribution, repair it, or convert an earlier installation. + +A new installation's flags seed /config.json. Afterwards, edit that file +and run /parsar apply; rerunning this installer only repairs. +""" import argparse import base64 import hashlib -import json import ipaddress +import json import os from pathlib import Path import platform @@ -16,17 +20,19 @@ import subprocess import sys import tempfile -import time -import urllib.error -import urllib.request -from urllib.parse import urlsplit import uuid +from urllib.parse import urlsplit -from configuration import compose_config, core_environment, environment_text, read_core_environment +import config_model +import configuration +import convert import local_node import native_service +import parsar_cli from distribution import DistributionError, artifact, image_identities, ensure_docker_image +SETTING_FLAGS = ("core_only", "web_only", "native_core", "core_port", "web_port", "core_url", "public_url") + class InstallError(Exception): pass @@ -37,16 +43,6 @@ def run(args, **kwargs): return subprocess.run(args, check=True, **kwargs) -def private_write(path, value): - descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600) - with os.fdopen(descriptor, "w") as stream: - stream.write(value) - - -def write_json(path, value): - private_write(path, json.dumps(value, indent=2) + "\n") - - def digest(path): result = hashlib.sha256() with path.open("rb") as stream: @@ -67,7 +63,10 @@ def verify_bundle(bundle): raise InstallError("Invalid distribution path") if digest(path) != expected: raise InstallError("Distribution checksum mismatch: " + name) - required = {"manifest.json", "install.sh", "install.py", "configuration.py", "native_service.py", "local_node.py", "node_spec.py", "node-install.pyz", "self-hosted-install.pyz", "distribution.py", "runtime/seccomp.json"} + required = {"manifest.json", "install.sh", "install.py", "configuration.py", "config_model.py", + "config.schema.json", "parsar_cli.py", "convert.py", "parsar.pyz", "native_service.py", + "local_node.py", "node_spec.py", "node-install.pyz", "self-hosted-install.pyz", "distribution.py", + "runtime/seccomp.json"} required.update(f"images/{name}.tar" for name in ("core", "web", "database")) required.update("native/bin/" + name for name in ("agents-api", "agents-api-migrate")) required.add("native/e2b/agents-api-e2b-provider") @@ -97,72 +96,124 @@ def database_port(): return sock.getsockname()[1] -def core_target(value): - from urllib.parse import urlsplit - parsed = urlsplit(value) - if (parsed.scheme not in ("http", "https") or not parsed.hostname or parsed.username or - parsed.password or parsed.query or parsed.fragment or parsed.path not in ("", "/")): - raise argparse.ArgumentTypeError("Core URL must be an HTTP(S) origin without credentials") - if parsed.scheme != "https" and parsed.hostname not in ("127.0.0.1", "localhost"): - raise argparse.ArgumentTypeError("Remote Core requires HTTPS") - return value.rstrip("/") +def origin(value): + check, message = config_model.CHECKS["origin"] + if not check(value): + raise argparse.ArgumentTypeError(message) + return value def arguments(argv=None): parser = argparse.ArgumentParser(description=__doc__) modes = parser.add_mutually_exclusive_group() - modes.add_argument("--core-only", action="store_true") - modes.add_argument("--web-only", action="store_true") - parser.add_argument("--native-core", action="store_true", help="Run Core as a systemd user service") - parser.add_argument("--sandbox-provider", choices=("true", "false"), nargs="?", const="true", default="false", + modes.add_argument("--core-only", action="store_true", default=None) + modes.add_argument("--web-only", action="store_true", default=None) + parser.add_argument("--native-core", action="store_true", default=None, help="Run Core as a systemd user service") + parser.add_argument("--sandbox-provider", choices=("true", "false"), nargs="?", const="true", help="Prepare a local sandbox provider (default: false)") parser.add_argument("--provider", choices=("microsandbox", "docker"), help="Local sandbox provider when enabled (default: microsandbox)") parser.add_argument("--install-dir", type=Path, default=Path.home() / ".parsar/core") - parser.add_argument("--core-port", type=int, default=8091) - parser.add_argument("--web-port", type=int, default=8080) - parser.add_argument("--core-url", type=core_target) - parser.add_argument("--public-url", type=core_target, help="Public HTTPS Core/Web origin behind your TLS reverse proxy") + parser.add_argument("--core-port", type=int) + parser.add_argument("--web-port", type=int) + parser.add_argument("--core-url", type=origin, help="Web-only: origin of the existing Core") + parser.add_argument("--public-url", type=origin, help="Public HTTPS origin behind your TLS reverse proxy") parser.add_argument("--core-key-file", type=Path, help="Web-only: private file containing the existing Core's Core key") + parser.add_argument("--config", type=Path, help="Seed a new installation's config.json from this file") + parser.add_argument("--convert", action="store_true", + help="Convert an installation made before config.json; also upgrades it to this release") + parser.add_argument("--yes", action="store_true", help="With --convert: do not ask for confirmation") parser.add_argument("--admin-token-file", type=Path, help=argparse.SUPPRESS) - parser.add_argument("--status", action="store_true", help="Read installation health; never invoke a model") - parser.add_argument("--stop", action="store_true", help="Stop installed services; retain all data") + parser.add_argument("--status", action="store_true", help=argparse.SUPPRESS) + parser.add_argument("--stop", action="store_true", help=argparse.SUPPRESS) args = parser.parse_args(argv) if args.admin_token_file: parser.error("--admin-token-file was renamed; use --core-key-file") - if args.web_only and args.native_core: - parser.error("--web-only cannot install native Core") - args.sandbox_provider = args.sandbox_provider == "true" - if args.provider and not args.sandbox_provider: - parser.error("--provider requires --sandbox-provider true") - if args.web_only and args.sandbox_provider: - parser.error("--web-only cannot install a sandbox provider") - args.provider = (args.provider or "microsandbox") if args.sandbox_provider else None - if args.provider and not (args.status or args.stop): - endpoint = urlsplit(args.public_url or "") + for retired in ("status", "stop"): + if getattr(args, retired): + parser.error(f"--{retired} is retired; run {args.install_dir / 'parsar'} {retired}") + if not args.install_dir.is_absolute(): + parser.error("--install-dir must be absolute") + args.given = [name for name, value in vars(args).items() + if name not in ("install_dir", "given") and value not in (None, False)] + return args + + +def seed_document(args): + """The --config file, which replaces the setting flags.""" + if args.config is None: + return None + if any(getattr(args, name) is not None for name in SETTING_FLAGS): + raise InstallError("--config replaces the setting flags; put those settings in the file") + try: + document = json.loads(args.config.read_text()) + except (OSError, ValueError): + raise InstallError("--config must name a readable JSON file") from None + if not isinstance(document, dict): + raise InstallError("--config must hold a JSON object") + return document + + +def check_flags(args, document): + """Flag combinations for a new installation, before config.json is seeded.""" + if document is None: + mode, native = ("core-only" if args.core_only else "web-only" if args.web_only else "all"), args.native_core + else: + mode, native = document.get("mode"), document.get("native_core") + sandbox = args.sandbox_provider == "true" + if args.provider and not sandbox: + raise InstallError("--provider requires --sandbox-provider true") + if mode == "web-only" and sandbox: + raise InstallError("--web-only cannot install a sandbox provider") + if mode == "web-only" and native: + raise InstallError("--web-only cannot install native Core") + if mode == "web-only" and not args.core_key_file: + raise InstallError("--web-only requires --core-key-file (and --core-url, or web.core_url in --config)") + if mode != "web-only" and args.core_key_file: + raise InstallError("--core-key-file requires --web-only") + return (args.provider or "microsandbox") if sandbox else None + + +def seed_config(args, document): + """config.json for a new installation, from flags or from --config.""" + if document is not None: + document.setdefault("$schema", "generated/config.schema.json") + if document.get("native_core"): + document.setdefault("ports", {}).setdefault("database", database_port()) + return config_model.validate(document) + mode = "core-only" if args.core_only else "web-only" if args.web_only else "all" + native = bool(args.native_core) + return config_model.initial(mode, native, public_url=args.public_url, **{ + "ports.core": args.core_port, "ports.web": args.web_port, "web.core_url": args.core_url, + "ports.database": database_port() if native else None}) + + +def check_public_url(config, provider): + if provider: + endpoint = urlsplit(config["public_url"] or "") try: loopback = ipaddress.ip_address(endpoint.hostname).is_loopback except ValueError: loopback = endpoint.hostname == "localhost" if endpoint.scheme != "https" or loopback: - parser.error("Local sandbox installation requires --public-url with HTTPS reachable from sandbox guests; loopback origins cannot be used") - if args.status and args.stop: - parser.error("Choose status or stop") - if not args.install_dir.is_absolute(): - parser.error("--install-dir must be absolute") - if any(not 1024 <= p <= 65535 for p in (args.core_port, args.web_port)): - parser.error("Ports must be between 1024 and 65535") - if not args.core_only and not args.web_only and args.core_port == args.web_port: - parser.error("Core and Web need different ports") - if args.web_only and not (args.core_url and args.core_key_file): - parser.error("--web-only requires --core-url and --core-key-file") - if not args.web_only and (args.core_url or args.core_key_file): - parser.error("Existing Core connection flags require --web-only") - return args + raise InstallError("Local sandbox installation requires public_url with HTTPS reachable from sandbox " + "guests; loopback origins cannot be used") -def compose(root, *args, **kwargs): - return run(["docker", "compose", "-f", str(root / "compose.json"), *args], **kwargs) +def read_core_key_file(source): + try: + info = source.lstat() + except OSError: + raise InstallError("Core key file must be an absolute, private regular file") from None + if (not source.is_absolute() or not stat.S_ISREG(info.st_mode) + or stat.S_IMODE(info.st_mode) & 0o077 or info.st_size > 4096): + raise InstallError("Core key file must be an absolute, private regular file") + token = source.read_text().strip() + if not token or any(c.isspace() for c in token) or "\x00" in token: + raise InstallError("Invalid Core key file") + if len(token) < 32: + raise InstallError("The Core key must have at least 32 characters") + return token def check_compose(): @@ -172,125 +223,28 @@ def check_compose(): raise InstallError("Docker Compose 2.26.0 or newer is required for literal Core environment values") -def wait_http(url, headers=None, attempts=60): - class NoRedirect(urllib.request.HTTPRedirectHandler): - def redirect_request(self, req, fp, code, msg, response_headers, newurl): - return None +def check_host(): + if platform.system() != "Linux" or platform.machine() not in ("x86_64", "amd64") or os.getuid() == 0: + raise InstallError("Run as a non-root user on Linux amd64 with Docker access") + check_compose() + run(["docker", "info", "--format", "{{.ServerVersion}}"], stdout=subprocess.DEVNULL) - # Probe credentials belong only to this endpoint, never a redirect or an - # ambient HTTP proxy. This also applies to the remote web-only Core probe. - opener = urllib.request.build_opener(urllib.request.ProxyHandler({}), NoRedirect()) - for attempt in range(attempts): - try: - request = urllib.request.Request(url, headers=headers or {}) - with opener.open(request, timeout=2) as response: - if response.status == 200: - return True - except (urllib.error.URLError, TimeoutError): - pass - if attempt + 1 < attempts: - time.sleep(1) - return False - - -def status(root, state): - output = compose(root, "ps", "--all", "--format", "json", capture_output=True, text=True).stdout - # Compose versions may return one array or one object per line. - rows = json.loads(output) if output.lstrip().startswith("[") else [json.loads(line) for line in output.splitlines() if line] - required = {"web"} if state["mode"] == "web-only" else {"database", "core"} - if state["mode"] == "all": - required.add("web") - observed = {row["Service"]: row for row in rows} - if native_service.is_native(state): - observed["core"] = {"State": "running" if native_service.active(state) else "stopped"} - print("Core service: " + observed["core"]["State"]) - healthy = all(name in observed and observed[name]["State"] == "running" - and observed[name].get("Health", "") in ("", "healthy") for name in required) - for row in rows: - print(f'{row["Service"]}: {row["State"]} {row.get("Health", "")}') - if state["mode"] != "web-only": - core_ok = wait_http(f'http://127.0.0.1:{state["core_port"]}/healthz', attempts=1) - healthy = healthy and core_ok - print("Core API: " + ("healthy" if core_ok else "unavailable")) - if state["mode"] != "core-only": - web_ok = wait_http(f'http://127.0.0.1:{state["web_port"]}/healthz', attempts=1) - healthy = healthy and web_ok - print("Web: " + ("healthy" if web_ok else "unavailable")) - print("Service health does not prove model execution. This check makes no model requests.") - if not healthy: - raise InstallError("One or more installed services are unavailable") - - -def initialize(root, args, manifest): - mode = "core-only" if args.core_only else "web-only" if args.web_only else "all" - if (root / "installation.json").exists(): - state = json.loads((root / "installation.json").read_text()) - wanted = (mode, args.native_core, args.core_port, args.web_port, args.core_url, args.public_url) - actual = (state["mode"], state["native_core"], 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 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 (root / "config/managed-runtimes.json").exists(): - raise InstallError("Retired file-managed provider configuration exists; preserve its resources and follow the deployment replacement guide") - if mode != "web-only": - read_core_environment(root, state) - directory = root / "state/e2b" - if (not directory.is_dir() or directory.is_symlink() or - stat.S_IMODE(directory.stat().st_mode) & 0o077): - raise InstallError("Private provider receipts are missing or unsafe; restore the retained installation") - return state - if root.exists() and any(root.iterdir()): - raise InstallError("Installation directory is not empty; refusing to overwrite existing state") - if mode == "web-only": - source = args.core_key_file - info = source.stat() - if (not source.is_absolute() or source.is_symlink() or not stat.S_ISREG(info.st_mode) - or stat.S_IMODE(info.st_mode) & 0o077 or info.st_size > 4096): - raise InstallError("Core key file must be an absolute, private regular file") - token = source.read_text().strip() - if not token or any(c.isspace() for c in token) or "\x00" in token: - raise InstallError("Invalid Core key file") - if len(token) < 32: - raise InstallError("The Core key must have at least 32 characters") - if mode != "web-only": - free_port(args.core_port) - if mode != "core-only": - free_port(args.web_port) - root.mkdir(mode=0o700, parents=True, exist_ok=True) - os.chmod(root, 0o700) - directories = ["config", "admin"] - if mode != "web-only": - directories += ["state", "state/e2b"] - for name in directories: - (root / name).mkdir(mode=0o700) - state = {"version": 1, "source_commit": manifest["source_commit"], "mode": mode, - "native_core": args.native_core, "installation_id": str(uuid.uuid4()), - "project": "parsar-" + secrets.token_hex(5), "uid": os.getuid(), "gid": os.getgid(), - "core_port": args.core_port, "web_port": args.web_port, "core_url": args.core_url, "public_url": args.public_url} - if native_service.is_native(state): - state["database_port"] = database_port() - config = root / "config" - if mode != "web-only": - private_write(config / "credential.key", base64.b64encode(secrets.token_bytes(32)).decode()) - private_write(config / "database.password", secrets.token_hex(32)) - core_key = secrets.token_hex(32) - private_write(root / "admin/core.key", core_key) - write_json(root / "admin/core-key-digests.json", [hashlib.sha256(core_key.encode()).hexdigest()]) + +def image_names(mode, native): if mode == "web-only": - private_write(root / "admin/core.key", token) - password = (config / "database.password").read_text() if mode != "web-only" else "" - if mode != "web-only": - private_write(config / "core.env", environment_text(core_environment(root, state, password))) - write_json(root / "compose.json", compose_config(root, state, manifest, password)) - write_json(root / "installation.json", state) - return state + return ["web"] + names = ["database"] if native else ["core", "database"] + return names + (["web"] if mode == "all" else []) + + +def image_loader(manifest, bundle): + def load(names): + return {name: ensure_docker_image(manifest, name, lambda name=name: bundle / f"images/{name}.tar") + for name in names} + return load -def prepare_node_payload(root, state, bundle): +def prepare_node_payload(root, state, bundle, replace=False): if state["mode"] == "core-only": return destination = root / "node-payload" @@ -307,6 +261,8 @@ def prepare_node_payload(root, state, bundle): or source.stat().st_size != entry["size"] or digest(source) != entry["sha256"]): raise InstallError("Offline artifact verification failed: " + logical) names.append(name) + if replace and destination.is_dir() and not destination.is_symlink(): + shutil.rmtree(destination) for name in names: source, target = bundle / name, destination / name target.parent.mkdir(parents=True, exist_ok=True, mode=0o700) @@ -324,84 +280,154 @@ def prepare_node_payload(root, state, bundle): os.unlink(temporary) +def install_parsar(root, bundle): + """Copy the parsar command into the installation; replace a missing or different copy.""" + target = root / "parsar" + source = bundle / "parsar.pyz" + if target.is_file() and not target.is_symlink() and digest(target) == digest(source): + os.chmod(target, 0o700) + return + descriptor, temporary = tempfile.mkstemp(prefix=".parsar-", dir=root) + os.close(descriptor) + try: + shutil.copyfile(source, temporary) + os.chmod(temporary, 0o700) + os.replace(temporary, target) + finally: + if os.path.exists(temporary): + os.unlink(temporary) + + +def layout(root): + if not root.exists() or not any(root.iterdir()): + return "empty" + config, legacy = (root / "config.json").exists(), (root / "installation.json").exists() + if config and legacy: + return "interrupted" + return "config" if config else "legacy" if legacy else "other" + + +def create(root, args, config, manifest, images): + """Write the new installation's secrets, config.json and state.json.""" + mode = config["mode"] + token = read_core_key_file(args.core_key_file) if mode == "web-only" else secrets.token_hex(32) + root.mkdir(mode=0o700, parents=True, exist_ok=True) + os.chmod(root, 0o700) + for name in ["secrets", "generated"] + ([] if mode == "web-only" else ["state", "state/e2b"]): + (root / name).mkdir(mode=0o700) + write = parsar_cli.create_private + write(root / "secrets/core.key", token) + if mode != "web-only": + write(root / "secrets/credential.key", base64.b64encode(secrets.token_bytes(32)).decode()) + write(root / "secrets/database.password", secrets.token_hex(32)) + state = {"format": 1, "installation_id": str(uuid.uuid4()), "project": "parsar-" + secrets.token_hex(5), + "uid": os.getuid(), "gid": os.getgid(), "mode": mode, "native_core": config.get("native_core", False), + "source_commit": manifest["source_commit"], "images": images, + "secrets_sha256": configuration.secret_digests(root, mode), "core_installation_id": None, + "execution_options_file": None, "applied": None, "converted_from": None} + state["secrets_sha256"].pop("core.key") + write(root / "config.json", json.dumps(config, indent=2) + "\n") + write(root / "state.json", json.dumps(state, indent=2) + "\n") + + +def finish(root, bundle, manifest, provider=None): + """Put the bundle's files in place, apply config.json and start the services.""" + state = parsar_cli.load_state(root) + # A converted installation replaces the earlier release's binaries and payload once. + replace = bool(state.get("converted_from")) and not state.get("applied") + prepare_node_payload(root, state, bundle, replace) + native_service.prepare(root, state, bundle, replace) + install_parsar(root, bundle) + parsar_cli.apply(root, start=True) + config, _ = parsar_cli.load_config(root) + mode = config["mode"] + if mode == "web-only" and parsar_cli.paired_core(root, config)[0] != 200: + raise InstallError("Core key authentication failed. Inspect secrets/core.key and web.core_url; no model was called") + if provider: + local_node.install(root, dict(state, provider=provider, core_port=config["ports"]["core"], + public_url=config["public_url"]), manifest, bundle, run) + public_url = config["public_url"] + if mode != "core-only": + url = f'http://127.0.0.1:{config["ports"]["web"]}' + print("Console: " + (public_url or url + " (local only)")) + if mode != "web-only": + print("API base URL: " + configuration.local_public_url(config) + "/v1" + ("" if public_url else " (local only)")) + print(f"Core key: {root / 'secrets/core.key'}. Keep it private; it also authorizes the Core management API.") + print(f"Settings: {root / 'config.json'}. Edit it, then run {root / 'parsar'} apply.") + print(f"Manage the services with {root / 'parsar'} status, start and stop.") + if mode != "web-only" and not provider: + print("No execution node was installed by this run. Choose a sandbox backend and add nodes in Web.") + elif provider: + print("Provider: " + provider + ". Local node enrolled; Core provisions Sessions on demand.") + print("Services installed. No model request was made. See docs/getting-started/quickstart.md.") + + def main(argv=None): args = arguments(argv) root = args.install_dir if root.is_symlink() or root.resolve() != root: raise InstallError("Installation directory must be canonical and not a symlink") - if args.status or args.stop: - state = json.loads((root / "installation.json").read_text()) - if args.stop: - if native_service.is_native(state): - native_service.stop(root, state) - compose(root, "stop") - print("Control-plane services stopped. Sandbox resources and data retained; running sandbox work may continue.") - else: - status(root, state) - return - if platform.system() != "Linux" or platform.machine() not in ("x86_64", "amd64") or os.getuid() == 0: - raise InstallError("Run as a non-root user on Linux amd64 with Docker access") - check_compose() - run(["docker", "info", "--format", "{{.ServerVersion}}"], stdout=subprocess.DEVNULL) bundle = Path(__file__).resolve().parent + kind = layout(root) + if kind in ("legacy", "interrupted") and not args.convert: + raise InstallError("This installation predates config.json; run ./install.sh --convert --install-dir " + f"{root}" if kind == "legacy" else + f"A conversion was interrupted; run ./install.sh --convert --install-dir {root} to finish it") + if kind in ("legacy", "interrupted"): + if set(args.given) - {"convert", "yes", "public_url"}: + raise InstallError("--convert accepts only --install-dir, --yes and --public-url") + check_host() + manifest = verify_bundle(bundle) + convert.convert(root, manifest, image_loader(manifest, bundle), args.public_url, args.yes, run) + finish(root, bundle, manifest) + return + if args.convert: + raise InstallError("--convert needs an installation made by an earlier release in --install-dir") + if kind == "config": + if args.given: + raise InstallError(f"This installation is configured by {root / 'config.json'}. Edit it and run " + f"{root / 'parsar'} apply; install.sh accepts only --install-dir to repair it") + state = parsar_cli.load_state(root) + check_host() + manifest = verify_bundle(bundle) + if state["source_commit"] != manifest["source_commit"]: + raise InstallError("This installation runs another release; upgrading an installation arrives with " + "parsar upgrade") + if native_service.is_native(state): + native_service.preflight(bundle, root) + images = image_loader(manifest, bundle)(list(state["images"])) + if images != state["images"]: + parsar_cli.save_state(root, dict(state, images=images)) + finish(root, bundle, manifest) + return + if kind == "other": + raise InstallError("Installation directory is not empty; refusing to overwrite existing state") + document = seed_document(args) + provider = check_flags(args, document) + config = seed_config(args, document) + check_public_url(config, provider) + if config["mode"] == "web-only": + read_core_key_file(args.core_key_file) + check_host() manifest = verify_bundle(bundle) - if args.native_core: + if config.get("native_core"): native_service.preflight(bundle, root) - if args.web_only: - images = ["web"] - elif args.native_core: - images = ["database"] - else: - images = ["core", "database"] - if not args.core_only and not args.web_only: - images.append("web") - local_images = dict(manifest["images"]) - for name in images: - 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 native_service.is_native(state): - native_service.prepare(root, state, bundle) - compose(root, "up", "--detach", "--wait") - if native_service.is_native(state): - migration_environment = read_core_environment(root, state) - run([str(root / "native/bin/agents-api-migrate")], env=dict(os.environ, **migration_environment), - stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) - native_service.start(root, state) - if state["mode"] != "web-only" and not wait_http(f'http://127.0.0.1:{state["core_port"]}/healthz'): - raise InstallError("Core did not become healthy. Use --status; retained state has not been removed") - if args.provider: - local_node.install(root, dict(state, provider=args.provider), manifest, bundle, run) - if state["mode"] != "core-only": - url = f'http://127.0.0.1:{state["web_port"]}' - host = urlsplit(state.get("public_url") or url).netloc - if not wait_http(url + "/console/auth", {"Host": host}): - raise InstallError("Web sign-in is unavailable. Use --status and inspect the Web service") - core_url = state.get("core_url") or f'http://127.0.0.1:{state["core_port"]}' - token = (root / "admin/core.key").read_text().strip() - if not wait_http(core_url + "/core/v1/projects", {"Authorization": "Bearer " + token}): - raise InstallError("Core key authentication failed. Inspect private configuration; no model was called") - print("Console: " + (state.get("public_url") or url)) - print("Sign in to Web with the Core key. Keep it private; it also authorizes the Core management API.") - if state["mode"] != "web-only": - print(f'API: http://127.0.0.1:{state["core_port"]}/v1') - print("Create a Project and issue its API key through the administrator API before calling the direct Core API.") - print("Core configuration file: " + str(root / "config/core.env")) - if args.provider: - print("Provider: " + args.provider + ". Local node enrolled; Core provisions Sessions on demand.") - else: - print("No execution node was installed by this run. Open Hosted Sandbox Manager to manage providers and nodes.") - print("Core key file: " + str(root / "admin/core.key")) - print("Services installed. No model request was made. See docs/getting-started/quickstart.md.") + for key in ("core", "web", "database"): + if key in config["ports"]: + free_port(config["ports"][key]) + images = image_loader(manifest, bundle)(image_names(config["mode"], config.get("native_core", False))) + create(root, args, config, manifest, images) + finish(root, bundle, manifest, provider) if __name__ == "__main__": try: main() - except (InstallError, local_node.LocalNodeError, DistributionError, RuntimeError, OSError, ValueError, KeyError, subprocess.CalledProcessError) as error: + except (InstallError, convert.ConvertError, parsar_cli.ParsarError, config_model.ConfigError, + local_node.LocalNodeError, DistributionError, RuntimeError) as error: + print(str(error), file=sys.stderr) + sys.exit(1) + except (OSError, ValueError, KeyError, subprocess.CalledProcessError): # Errors never include generated configuration or external process output. - print(str(error) if isinstance(error, (InstallError, local_node.LocalNodeError, DistributionError, RuntimeError)) else "Installation failed; inspect prerequisites and private deployment files", file=sys.stderr) + print("Installation failed; inspect prerequisites and private deployment files", file=sys.stderr) sys.exit(1) diff --git a/deploy/install/installer_fakes.py b/deploy/install/installer_fakes.py new file mode 100644 index 000000000..080c6f6ba --- /dev/null +++ b/deploy/install/installer_fakes.py @@ -0,0 +1,236 @@ +"""A fake Docker, systemd and HTTP host for installer tests. Nothing real is started. + +Compose is modelled by its observable contract: `up` recreates a service exactly when +its resolved configuration (including env_file content) changed, and Core loads its +Core key digests when it starts. +""" +import hashlib +import json +from pathlib import Path +import subprocess +from types import SimpleNamespace +from unittest import mock + +import parsar_cli + + +def sha256(data): + return hashlib.sha256(data.encode() if isinstance(data, str) else data).hexdigest() + + +class FakeHost: + def __init__(self, test): + self.commands = [] + self.containers = {} # service -> configuration hash of its container + self.running = set() + self.recreated = [] + self.native = {"active": False, "starts": 0, "restarts": 0, "reloads": 0, "addr": None, "digests": []} + self.core = {"port": None, "digests": [], "fails": False, "log": ""} + self.web_port = None + self.native_root = None # the installation whose native unit systemctl manages + self.missing_images = set() + self.core_installation_id = "11111111-2222-4333-8444-555555555555" + self.bindings = {"nodes": 0, "nodes_on_other_address": 0, "hosted_sandboxes": 0, "self_hosted_executors": 0} + self.nodes = [] + self.remote_core = {} # web-only: origin -> (status, installation_id) + self.deployment_core_url = "" # what an old Core reports for its sandbox deployment + for target, replacement in ((subprocess, "run"),): + patcher = mock.patch.object(target, replacement, side_effect=self.run) + patcher.start() + test.addCleanup(patcher.stop) + for name, value in (("http", self.http), ("time", SimpleNamespace(sleep=lambda seconds: None))): + patcher = mock.patch.object(parsar_cli, name, value) + patcher.start() + test.addCleanup(patcher.stop) + + # Commands --------------------------------------------------------------- + def run(self, args, check=False, **kwargs): + args = [str(item) for item in args] + self.commands.append(args) + code, stdout = 0, "" + if args[:2] == ["docker", "compose"] and args[2:3] == ["-f"]: + code, stdout = self.compose(Path(args[3]), args[4:]) + elif args[:2] == ["docker", "compose"]: + stdout = "2.30.0" + elif args[:3] == ["docker", "image", "inspect"]: + code = 1 if args[3] in self.missing_images else 0 + stdout = "" if code else args[3] + " linux/amd64" + elif args[0] == "systemctl": + code, stdout = self.systemctl(args[2:]) + elif args[0] == "loginctl": + stdout = "yes" + elif args[0] == "journalctl": + stdout = self.core["log"] + if check and code: + raise subprocess.CalledProcessError(code, args) + text = kwargs.get("text") or kwargs.get("universal_newlines") + return subprocess.CompletedProcess(args, code, stdout if text else stdout.encode(), "" if text else b"") + + def service_hash(self, document, name): + service = document["services"][name] + text = json.dumps(service, sort_keys=True) + for path in service.get("env_file", []): + text += Path(path.replace("$$", "$")).read_text() + return sha256(text) + + def compose(self, path, args): + document = json.loads(path.read_text()) if path.exists() else {"services": {}} + services = document["services"] + if args[:1] == ["ps"]: + rows = [{"Service": name, "State": "running" if name in self.running else "exited", "Health": ""} + for name in self.containers if name in services] + return 0, json.dumps(rows) + if args[:1] == ["stop"]: + self.running -= set(args[1:] or services) + return 0, "" + if args[:1] == ["logs"]: + return 0, self.core["log"] + if args[:1] == ["exec"]: + return 0, "" + if args[:1] == ["up"]: + targets = [item for item in args[1:] if not item.startswith("-") and not item.isdigit()] + names = list(services) if not targets else [name for name in ("database", "migrate", "core") + if name in services] + for name in names: + digest = self.service_hash(document, name) + if self.containers.get(name) != digest: + self.containers[name] = digest + self.recreated.append(name) + if name == "core": + self.load_core(path.parent.parent, services[name]) + if name != "migrate": + self.running.add(name) + if "web" in names: + web = services["web"] + self.web_port = int(web["ports"][0].split(":")[1]) if "ports" in web else int( + web["environment"]["CORE_CONSOLE_ADDR"].rsplit(":", 1)[1]) + if "core" in names and self.core["fails"]: + self.running.discard("core") + return 1, "" + return 0, "" + + def load_core(self, root, service): + if "ports" in service: + self.core["port"] = int(service["ports"][0].split(":")[1]) + digests = root / "generated/core-key-digests.json" + self.core["digests"] = json.loads(digests.read_text()) if digests.exists() else json.loads( + (root / "admin/core-key-digests.json").read_text()) + + def systemctl(self, args): + native = self.native + if args[0] == "daemon-reload": + native["reloads"] += 1 + elif args[0] in ("enable", "restart", "start"): + native["starts" if args[0] != "restart" else "restarts"] += 1 + native["active"] = not self.core["fails"] + if (self.native_root / "generated/core.env").exists(): + environment = (self.native_root / "generated/core.env").read_text() + address = next(line for line in environment.splitlines() if line.startswith("AGENTS_API_ADDR=")) + native["addr"] = int(address.rsplit(":", 1)[1].rstrip('"')) + native["digests"] = json.loads((self.native_root / "generated/core-key-digests.json").read_text()) + elif args[0] in ("stop", "disable"): + native["active"] = False + elif args[0] == "is-active": + return (0 if native["active"] else 3), "" + elif args[0] == "show": + return 0, "252" + return 0, "" + + # HTTP ------------------------------------------------------------------- + def core_listening(self, port): + if "core" in self.running and self.core["port"] == port: + return self.core["digests"] + if self.native["active"] and self.native["addr"] == port: + return self.native["digests"] + return None + + def http(self, url, headers=None, timeout=5): + headers = headers or {} + origin, _, path = url.partition("://")[2].partition("/") + path = "/" + path + key = headers.get("Authorization", "").removeprefix("Bearer ") + for remote, (status, installation) in self.remote_core.items(): + if url.startswith(remote + "/"): + return status, json.dumps({"installation_id": installation}).encode() + port = int(origin.rsplit(":", 1)[1]) + digests = self.core_listening(port) + if path == "/healthz": + if digests is not None or ("web" in self.running and self.web_port == port): + return 200, b"{}" + return 0, b"" + if path == "/console/auth": + return (200, b"{}") if "web" in self.running and self.web_port == port else (0, b"") + if digests is None: + return 0, b"" + if sha256(key) not in digests: + return 401, b"" + if path == "/core/v1/installation": + return 200, json.dumps({"installation_id": self.core_installation_id, + "address_bindings": self.bindings}).encode() + if path == "/core/v1/sandbox/nodes": + return 200, json.dumps({"data": self.nodes}).encode() + if path == "/core/v1/sandbox/deployment": + return 200, json.dumps({"core_url": self.deployment_core_url}).encode() + return 404, b"" + + +MANIFEST = { + "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": "parsar-core-runtime@sha256:" + "b" * 64, + "microsandbox": {"runtime_sha256": "5" * 64, "firmware_sha256": "6" * 64}, +} +MODULES = ("install.py", "configuration.py", "config_model.py", "config.schema.json", "parsar_cli.py", "convert.py", + "native_service.py", "local_node.py", "node_spec.py", "distribution.py", "install.sh") + + +def write_checksums(bundle): + files = sorted(path for path in bundle.rglob("*") if path.is_file() and path.name != "SHA256SUMS") + (bundle / "SHA256SUMS").write_text("".join( + sha256(path.read_bytes()) + " " + str(path.relative_to(bundle)) + "\n" for path in files)) + + +def make_bundle(directory, manifest, commit=None): + """A synthetic distribution with this checkout's installer modules.""" + manifest = json.loads(json.dumps(manifest)) + if commit: + manifest["source_commit"] = commit + bundle = Path(directory) + (bundle / "images").mkdir(parents=True) + (bundle / "runtime").mkdir() + (bundle / "runtime/seccomp.json").write_text('{"defaultAction":"SCMP_ACT_ERRNO"}') + for name in MODULES: + (bundle / name).write_bytes(Path(__file__).with_name(name).read_bytes()) + for name in ("node-install.pyz", "self-hosted-install.pyz"): + (bundle / name).write_bytes(b"synthetic verified Python bootstrap") + (bundle / "parsar.pyz").write_bytes(b"synthetic parsar command " + manifest["source_commit"].encode()) + manifest["artifacts"] = {} + 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"): + manifest["artifacts"][name] = {"filename": "parsar-" + manifest["source_commit"] + "-" + name.replace("/", "-"), + "sha256": "a" * 64, "size": 1} + (bundle / "manifest.json").write_text(json.dumps(manifest)) + for name in manifest["images"]: + (bundle / "images" / (name + ".tar")).write_bytes(("synthetic " + name).encode()) + for name in ("bin/agents-api", "bin/agents-api-migrate", "bin/agents-api-microsandbox-provider", + "bin/parsar-sandbox-node", "microsandbox/msb", "microsandbox/libkrunfw.so.5.6.1", + "e2b/agents-api-e2b-provider"): + path = bundle / "native" / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_bytes(b"\x7fELFsynthetic native file " + manifest["source_commit"].encode()) + write_checksums(bundle) + return bundle, manifest + + +def run_installer(install, bundle, argv): + """install.main from the bundle on a Linux amd64 host.""" + 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.native_service.platform, "system", return_value="Linux"), \ + mock.patch.object(install, "free_port"): + return install.main([str(item) for item in argv]) diff --git a/deploy/install/local_node.py b/deploy/install/local_node.py index 6809c2631..a45757da8 100644 --- a/deploy/install/local_node.py +++ b/deploy/install/local_node.py @@ -38,7 +38,7 @@ def request(core, token, method, path, value=None): def install(root, state, manifest, bundle, run): core = f'http://127.0.0.1:{state["core_port"]}' - admin = (root / "admin/core.key").read_text().strip() + admin = (root / "secrets/core.key").read_text().strip() current = request(core, admin, "GET", "deployment") if current.get("installation_id") != state["installation_id"]: raise LocalNodeError("Core deployment identity differs; preserve its installation state") diff --git a/deploy/install/native_service.py b/deploy/install/native_service.py index 9a8eb92bd..2d0649b90 100644 --- a/deploy/install/native_service.py +++ b/deploy/install/native_service.py @@ -9,8 +9,6 @@ import subprocess import tempfile -from configuration import read_core_environment - REQUIRED = ("bin/agents-api", "bin/agents-api-migrate", "e2b/agents-api-e2b-provider") @@ -19,7 +17,7 @@ def is_native(state): return state["mode"] != "web-only" and state["native_core"] -def _unit_name(state): +def unit_name(state): project = state.get("project", "") if not isinstance(project, str) or not re.fullmatch(r"parsar-[0-9a-f]{10}", project): raise RuntimeError("Native Core requires its installation's generated project name") @@ -103,36 +101,22 @@ def _digest(path): return digest.digest() -def _private_write(path, content): - if path.is_symlink() or (path.exists() and not path.is_file()): - raise RuntimeError("Native Core configuration must be a regular private file") - if path.exists() and path.read_text() == content: - os.chmod(path, 0o600) - return - descriptor, temporary = tempfile.mkstemp(prefix=".core-", dir=path.parent) - try: - with os.fdopen(descriptor, "w", encoding="utf-8") as stream: - stream.write(content) - os.replace(temporary, path) - finally: - if os.path.exists(temporary): - os.unlink(temporary) - - -def prepare(root, state, bundle): +def prepare(root, state, bundle, replace=False): + """Install the bundle's native Core binaries. Only a conversion replaces different ones.""" if not is_native(state): return try: root, bundle = _path(root), _path(bundle) - unit_name = _unit_name(state) - read_core_environment(root, state) + unit_name(state) source, target = bundle / "native", root / "native" incoming = _files(source) if target.exists() or target.is_symlink(): installed = _files(target) - if incoming.keys() != installed.keys() or any(_digest(path) != _digest(installed[name]) for name, path in incoming.items()): + if incoming.keys() == installed.keys() and all(_digest(path) == _digest(installed[name]) for name, path in incoming.items()): + replace = False + elif not replace: raise RuntimeError("Installed native Core files differ; preserve the installation and follow the upgrade guide") - else: + if not target.exists() or replace: root.mkdir(parents=True, mode=0o700, exist_ok=True) with tempfile.TemporaryDirectory(prefix=".native-", dir=root) as temporary: staged = Path(temporary) / "native" @@ -140,46 +124,64 @@ def prepare(root, state, bundle): destination = staged / name destination.parent.mkdir(parents=True, exist_ok=True) shutil.copyfile(path, destination) + if target.exists(): + os.replace(target, Path(temporary) / "previous") os.replace(staged, target) for path in [target, target / "bin", target / "e2b", *(target / name for name in REQUIRED)]: os.chmod(path, 0o700) - config = root / "config" - if config.is_symlink(): - raise RuntimeError("Native Core configuration directory must not be a symlink") - config.mkdir(mode=0o700, exist_ok=True) - os.chmod(config, 0o700) - # ':' disables command-line environment substitution. The executable is - # still Core itself; no shell, wrapper or provider shutdown hook is used. - executable = str(target / "bin/agents-api").replace("%", "%%").replace('"', '\\"') - unit = ("[Unit]\nDescription=Parsar Core\n\n[Service]\nType=exec\n" - + 'ExecStart=:"' + executable + '"\n' - + "WorkingDirectory=" + str(root).replace("%", "%%") + "\n" - + "EnvironmentFile=" + str(config / "core.env").replace("%", "%%") + "\n" - + "Restart=on-failure\nKillMode=process\nUMask=0077\n\n[Install]\nWantedBy=default.target\n") - _private_write(config / unit_name, unit) except (OSError, UnicodeError): raise RuntimeError("Cannot prepare private native Core files; existing state was not removed") from None +def unit_text(root, header): + root = _path(root) + # ':' disables command-line environment substitution. The executable is + # still Core itself; no shell, wrapper or provider shutdown hook is used. + executable = str(root / "native/bin/agents-api").replace("%", "%%").replace('"', '\\"') + return ("# " + header + "\n" + + "[Unit]\nDescription=Parsar Core\n\n[Service]\nType=exec\n" + + 'ExecStart=:"' + executable + '"\n' + + "WorkingDirectory=" + str(root).replace("%", "%%") + "\n" + + "EnvironmentFile=" + str(root / "generated/core.env").replace("%", "%%") + "\n" + + "Restart=on-failure\nKillMode=process\nUMask=0077\n\n[Install]\nWantedBy=default.target\n") + + +def daemon_reload(): + _checked(["systemctl", "--user", "daemon-reload"], "Cannot reload the systemd user manager") + + def start(root, state): + """Enable the generated unit by path and start it.""" if not is_native(state): return - unit = _path(root) / "config" / _unit_name(state) + unit = _path(root) / "generated" / unit_name(state) if unit.is_symlink() or not unit.is_file(): - raise RuntimeError("Native Core service must be prepared before starting it") - _checked(["systemctl", "--user", "daemon-reload"], "Cannot reload the systemd user manager") + raise RuntimeError("Native Core service must be generated before starting it; run parsar apply") + daemon_reload() _checked(["systemctl", "--user", "enable", "--now", str(unit)], "Cannot enable or start this installation's native Core service") +def restart(state): + if is_native(state): + _checked(["systemctl", "--user", "restart", unit_name(state)], "Cannot restart this installation's native Core service") + + def stop(root, state): if is_native(state): _path(root) - _checked(["systemctl", "--user", "stop", _unit_name(state)], "Cannot stop this installation's native Core service") + _checked(["systemctl", "--user", "stop", unit_name(state)], "Cannot stop this installation's native Core service") + + +def disable(state): + """Stop the unit and remove its enablement link, as a conversion does before moving it.""" + if is_native(state): + _checked(["systemctl", "--user", "disable", "--now", unit_name(state)], + "Cannot disable this installation's native Core service") def active(state): if not is_native(state): return False - result = _run(["systemctl", "--user", "is-active", "--quiet", _unit_name(state)], + result = _run(["systemctl", "--user", "is-active", "--quiet", unit_name(state)], "Cannot query this installation's native Core service") return result.returncode == 0 diff --git a/deploy/install/parsar_cli.py b/deploy/install/parsar_cli.py new file mode 100644 index 000000000..b5e6392ec --- /dev/null +++ b/deploy/install/parsar_cli.py @@ -0,0 +1,676 @@ +"""parsar: status, start, stop, apply and rotate-core-key for one installation. + +The installation directory is the directory that holds the command. The bundle it +was installed from is never needed. config.json is the only file an operator edits; +apply derives generated/ from it and restarts exactly the services whose inputs +changed. +""" +import argparse +import contextlib +import datetime +import fcntl +import json +import os +from pathlib import Path +import secrets +import stat +import subprocess +import sys +import tempfile +import time +import urllib.error +import urllib.request +from urllib.parse import urlsplit + +import config_model +import configuration +import native_service + + +class ParsarError(Exception): + pass + + +def run(args, **kwargs): + # Never print a generated Compose file, process environment or secret value. + return subprocess.run(args, **dict({"check": True}, **kwargs)) + + +def now(): + return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + +# Private files --------------------------------------------------------------- + +def check_private(path, what): + try: + info = os.lstat(path) + except FileNotFoundError: + raise ParsarError(f"{what} is missing") from None + if (not stat.S_ISREG(info.st_mode) or stat.S_IMODE(info.st_mode) & 0o077 + or info.st_uid != os.geteuid() or info.st_nlink != 1): + raise ParsarError(f"{what} must be a regular file with mode 0600, owned by you, and not a link") + + +def read_private(path, what): + check_private(path, what) + descriptor = os.open(path, os.O_RDONLY | os.O_NOFOLLOW) + with os.fdopen(descriptor, "rb") as stream: + return stream.read() + + +def create_private(path, data): + descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600) + with os.fdopen(descriptor, "wb") as stream: + stream.write(data.encode() if isinstance(data, str) else data) + + +def write_private(path, data): + """Replace a file atomically with a 0600 copy in the same directory.""" + path = Path(path) + descriptor, temporary = tempfile.mkstemp(prefix="." + path.name + ".", dir=path.parent) + try: + with os.fdopen(descriptor, "wb") as stream: + stream.write(data.encode() if isinstance(data, str) else data) + os.replace(temporary, path) + finally: + if os.path.exists(temporary): + os.unlink(temporary) + + +def load_config(root): + raw = read_private(root / "config.json", "config.json") + try: + document = json.loads(raw) + except ValueError as error: + line = getattr(error, "lineno", "?") + raise ParsarError(f"config.json is not valid JSON (line {line})") from None + return config_model.validate(document), configuration.sha256(raw) + + +def load_state(root): + state = json.loads(read_private(root / "state.json", "state.json")) + if state.get("format") != 1: + raise ParsarError("state.json has an unknown format; use the parsar command of this installation's release") + return state + + +def save_state(root, state): + write_private(root / "state.json", json.dumps(state, indent=2) + "\n") + + +@contextlib.contextmanager +def locked(root): + descriptor = os.open(root / ".parsar.lock", os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW, 0o600) + try: + try: + fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + raise ParsarError("Another parsar command is running for this installation") from None + yield + finally: + os.close(descriptor) + + +# Services -------------------------------------------------------------------- + +def compose(root, *args, **kwargs): + return run(["docker", "compose", "-f", str(root / "generated/compose.json"), *args], **kwargs) + + +def service_states(root): + if not (root / "generated/compose.json").exists(): + return {} + output = compose(root, "ps", "--all", "--format", "json", capture_output=True, text=True).stdout + # Compose versions may return one array or one object per line. + rows = json.loads(output) if output.lstrip().startswith("[") else [ + json.loads(line) for line in output.splitlines() if line.strip()] + return {row["Service"]: row for row in rows} + + +def running(root, state): + names = {name for name, row in service_states(root).items() if row.get("State") == "running"} + if native_service.is_native(state) and native_service.active(state): + names.add("core") + return names + + +def migrate_native(root): + environment = configuration.read_environment((root / "generated/core.env").read_text()) + run([str(root / "native/bin/agents-api-migrate")], env=dict(os.environ, **environment), + stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) + + +def plan_running(state, services, was_running, start): + """Services running after apply. Compose starts every container once any runs.""" + native = native_service.is_native(state) + containers = set(services) - {"migrate"} - ({"core"} if native else set()) + result = containers if start or was_running & containers else set() + if native and (start or "core" in was_running): + result = result | {"core"} + return result + + +def converge(root, state, services, was_running, core_port, start=False, reload_unit=False, core_first=False): + """Bring the services that run to the written files. Stopped services stay stopped.""" + native = native_service.is_native(state) + will_run = plan_running(state, services, was_running, start) + if native and reload_unit: + native_service.daemon_reload() + if core_first and "core" in will_run: + if native: + native_service.restart(state) + else: + compose(root, "up", "--detach", "--wait", "--wait-timeout", "300", "core") + if wait_status(f"http://127.0.0.1:{core_port}/healthz") is None: + raise ParsarError("Core did not become healthy") + if will_run - ({"core"} if native else set()): + compose(root, "up", "--detach", "--wait", "--wait-timeout", "300") + if native and "core" in will_run and not core_first: + if start: + migrate_native(root) + native_service.start(root, state) + else: + native_service.restart(state) + + +def describe(error): + """A failure message without command paths, output or environment.""" + if isinstance(error, subprocess.CalledProcessError): + words = [word for word in error.cmd if not word.startswith(("/", "-"))][:3] + return f"`{' '.join(words)}` failed" + return str(error) + + +def core_error_line(root, state): + """Core's single startup failure line. Core logs no environment values.""" + try: + if native_service.is_native(state): + output = run(["journalctl", "--user", "--unit", native_service.unit_name(state), "--lines", "200", + "--no-pager", "--output", "cat"], capture_output=True, text=True).stdout + else: + output = compose(root, "logs", "--no-log-prefix", "--tail", "200", "core", + capture_output=True, text=True).stdout + except (subprocess.CalledProcessError, OSError): + return None + lines = [line.strip() for line in output.splitlines() if "agents-api startup failed" in line] + return lines[-1] if lines else None + + +# HTTP ------------------------------------------------------------------------ + +class _NoRedirect(urllib.request.HTTPRedirectHandler): + def redirect_request(self, *args, **kwargs): + return None + + +def http(url, headers=None, timeout=5): + """(status, body). Credentials go only to this URL: no redirects, no ambient proxy.""" + opener = urllib.request.build_opener(urllib.request.ProxyHandler({}), _NoRedirect()) + try: + with opener.open(urllib.request.Request(url, headers=headers or {}), timeout=timeout) as response: + return response.status, response.read(1024 * 1024) + except urllib.error.HTTPError as error: + return error.code, b"" + except (urllib.error.URLError, OSError, ValueError): + return 0, b"" + + +def wait_status(url, headers=None, expect=200, attempts=60): + for attempt in range(attempts): + status, body = http(url, headers) + if status == expect: + return body + if attempt + 1 < attempts: + time.sleep(1) + return None + + +def bearer(key): + return {"Authorization": "Bearer " + key} + + +def core_base(config): + return f'http://127.0.0.1:{config["ports"]["core"]}' + + +def health(root, config, expected): + """config needs public_url and ports; it may be the applied view of a changed config.json.""" + if "core" in expected: + base = core_base(config) + if wait_status(base + "/healthz") is None: + raise ParsarError("Core did not become healthy") + if wait_status(base + "/core/v1/installation", bearer(configuration.read_core_key(root)), attempts=10) is None: + raise ParsarError("Core did not accept the Core key at /core/v1/installation") + if "web" in expected: + url = f'http://127.0.0.1:{config["ports"]["web"]}' + if wait_status(url + "/console/auth", {"Host": urlsplit(config["public_url"] or url).netloc}) is None: + raise ParsarError("Web sign-in is unavailable") + + +# Apply ----------------------------------------------------------------------- + +def check_fixed(config, state): + for key in ("mode", "native_core"): + if config.get(key, False) != state[key]: + raise ParsarError(f"{key} is fixed after installation ({json.dumps(state[key])}). " + "Install into a new directory to change it; nothing was applied.") + + +def check_secrets(root, config, state): + reasons = {"credential.key": "Stored credentials can only be read with the original key.", + "database.password": "PostgreSQL keeps the password it was initialized with."} + if config["mode"] != "web-only": + for name, reason in reasons.items(): + data = read_private(root / "secrets" / name, "secrets/" + name) + if configuration.sha256(data) != state["secrets_sha256"][name]: + raise ParsarError(f"secrets/{name} changed since installation. {reason} " + "Restore the original file; nothing was applied.") + check_private(root / "secrets/core.key", "secrets/core.key") + configuration.read_core_key(root) + + +def execution_options_import(root, state, out): + """Hook for phase 3's one-time import of a retained execution options file. + + install.sh --convert keeps an existing AGENTS_API_EXECUTION_OPTIONS_FILE and its + file unchanged (state.json execution_options_file). Phase 3 replaces this body with + the import into Core's deployment model providers and then clears the record. + """ + retained = state.get("execution_options_file") + if retained: + out(f"Note: Core still reads {retained['path']} through AGENTS_API_EXECUTION_OPTIONS_FILE. " + "A later release imports it into Core once and removes it.") + return state + + +def edited_files(root, state): + """Generated files whose content differs from what the last apply wrote.""" + edited = [] + for name, digest in ((state.get("applied") or {}).get("files") or {}).items(): + path = root / "generated" / name + current = None if path.is_symlink() or not path.is_file() else configuration.sha256(path.read_bytes()) + if current != digest: + edited.append(name) + return edited + + +def previous_values(root, state, edited): + """The applied settings, from the snapshot and the one generated file with sensitive values.""" + if not state.get("applied") or "settings.json" in edited: + return None + values = {item["key"]: item["value"] for item in + json.loads((root / "generated/settings.json").read_text())["settings"]} + history = root / "generated/runtime-history.json" + if "core.runtime_history.headers" in values and "runtime-history.json" not in edited: + values["core.runtime_history.headers"] = json.loads(history.read_text()).get("headers") if history.exists() else None + return values + + +def confirm_public_url(root, config, previous, was_running, args, interactive, out): + """Changing the public URL strands what is bound to the old one; list it and confirm.""" + old = configuration.local_public_url({"public_url": previous["public_url"], + "ports": {"core": previous["ports.core"]}}) + new = configuration.local_public_url(config) + if old == new: + return + counts, nodes = None, [] + if "core" in was_running: + base, key = f'http://127.0.0.1:{previous["ports.core"]}', configuration.read_core_key(root) + status, body = http(base + "/core/v1/installation", bearer(key)) + if status == 200: + counts = json.loads(body).get("address_bindings") or {} + status, body = http(base + "/core/v1/sandbox/nodes", bearer(key)) + nodes = json.loads(body).get("data", []) if status == 200 else [] + out(f"The public URL changes from {old} to {new}.") + if counts is not None: + out(f'Bound to the current address: {counts.get("nodes", 0)} node(s), ' + f'{counts.get("hosted_sandboxes", 0)} hosted sandbox(es), ' + f'{counts.get("self_hosted_executors", 0)} self-hosted executor credential(s).') + for node in nodes: + out(f' node {node.get("name")}: {"online" if node.get("online") else "offline"}') + if not any(counts.get(name) for name in ("nodes", "hosted_sandboxes", "self_hosted_executors")): + return + else: + out("Core is not running, so the nodes, sandboxes and executors bound to the address can't be counted.") + out("After the change, nodes on the old address get no new sandboxes; remove them in Web and add them again.\n" + "Existing sandboxes and executors keep working only while the old address still reaches this Core, so\n" + "keep the old route until they are replaced. Self-hosted executors must restart with the new remote_url.") + if args.dry_run: + out("Applying this change needs confirmation.") + elif args.confirm_public_url_change is not None: + if args.confirm_public_url_change != new: + raise ParsarError(f"--confirm-public-url-change must equal the new public URL {new}; nothing was applied") + elif interactive: + if input("Type the new public URL to continue: ").strip() != new: + raise ParsarError("The public URL change was not confirmed; nothing was applied") + else: + raise ParsarError(f"Confirm with --confirm-public-url-change {new}; nothing was applied") + + +def paired_core(root, config): + """(HTTP status, installation ID) of the Core that web.core_url reaches.""" + status, body = http(config["web"]["core_url"] + "/core/v1/installation", bearer(configuration.read_core_key(root))) + return status, (json.loads(body).get("installation_id") if status == 200 else None) + + +def check_paired_core(root, config, state, args, interactive, out): + """Web-only: detect a web.core_url that now reaches a different Core.""" + status, installation = paired_core(root, config) + if status == 401: + out("Warning: Core rejects this Web host's Core key; the key is out of date. " + "Copy secrets/core.key from the Core host, then run parsar apply.") + if status != 200: + return state.get("core_installation_id") + recorded = state.get("core_installation_id") + if recorded and installation != recorded: + out(f"web.core_url reaches Core installation {installation}, not the paired installation {recorded}.") + if args.dry_run: + out("Applying this change needs confirmation.") + elif not (args.yes or (interactive and input("Type yes to pair Web with this Core: ").strip() == "yes")): + raise ParsarError("Pairing Web with a different Core was not confirmed; nothing was applied") + return installation + + +def read_generated(root, names): + result = {} + for name in names: + path = root / "generated" / name + result[name] = path.read_bytes() if path.is_file() and not path.is_symlink() else None + return result + + +def apply(root, dry_run=False, yes=False, discard_edits=False, confirm_public_url_change=None, + start=False, interactive=None, out=print): + root = Path(root) + args = argparse.Namespace(dry_run=dry_run, yes=yes, confirm_public_url_change=confirm_public_url_change) + interactive = sys.stdin.isatty() if interactive is None else interactive + with locked(root): + return _apply(root, args, discard_edits, start, interactive, out) + + +def _apply(root, args, discard_edits, start, interactive, out, core_first=False, was_running=None, restore=None): + config, config_digest = load_config(root) + state = load_state(root) + check_fixed(config, state) + check_secrets(root, config, state) + state = execution_options_import(root, state, out) + applied = state.get("applied") or {} + edited = edited_files(root, state) + if edited and not discard_edits: + raise ParsarError("\n".join(f"generated/{name} was edited by hand." for name in edited) + + "\nPut the change in config.json and run parsar apply --discard-edits, which keeps the" + " edited copy as generated/.edited-