diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 804abe4fd..6b2875a12 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -1714,15 +1714,28 @@ across upgrades. This is the same managed Runtime, not user-managed enrollment. #### Matched Core and console distribution [Configuration](docs/configuration.md) is the canonical operator parameter reference. -Compose and native launchers load the same private `config/core.env`; the installer -creates it once, validates retained literal values and never replaces user edits. -Core logs the loaded path without values. The installation receipt records packaging -and identity, not provider overrides. Explicit local-node flags call the ordinary -administrator API once; PostgreSQL owns the resulting selection. Administrator-issued -enrollment approves capacity (default two active/eight retained); a node cannot -supply or overwrite those limits. Downloaded specification copies remain validated -against the existing database-owned resources/Runtime contract. Do not add a new -configuration format, loader precedence, hot reload or embedded Core node. +Every process setting has one home: the installation's private `config.json`, +described by `deploy/install/config.schema.json`. The operator edits only that +file; `parsar apply` validates it, derives `generated/` (Compose file, `core.env`, +native unit, Core key digest file, settings snapshot) and converges on what actually +runs: each service carries the digest of its inputs (Compose label +`io.parsar.inputs`, native `PARSAR_INPUTS`), and exactly the services whose running +inputs differ are recreated or restarted. Decide restarts from what runs, never +from recorded bookkeeping, so the next apply finishes any interrupted one. Installation flags only seed it, and +rerunning the installer rejects them. Runtime settings stay in PostgreSQL and +change through Web or `/core/v1`. Secrets live once each in `secrets/`; identity and +install facts live in tool-written `state.json`. Core still reads only its +environment and has no config loader; it serves the non-secret snapshot at +`GET /core/v1/installation`. Keep the schema, the subset validator +(`config_model.py`), the generator and the generated reference table in +`docs/configuration.md` (`scripts/config-reference.py`) in step. Do not add a +second operator configuration file, loader precedence, hot reload, compatibility +reading of retired names, or an embedded Core node. Explicit local-node flags call +the ordinary administrator API once; PostgreSQL owns the resulting selection. +Administrator-issued enrollment approves capacity (default two active/eight +retained); a node cannot supply or overwrite those limits. Downloaded specification +copies remain validated against the existing database-owned resources/Runtime +contract. The installer packages Core and the Web console together, with independent `--core-only` and `--web-only` modes. `site/` is the public static @@ -1803,7 +1816,8 @@ allocates, wakes a sandbox or grants project resource access. It is an `/api/v1` machine route that reaches Core directly, never through the console. Bounded polling and reruns retain the original container and history; timeout is a diagnostic failure, not permission to relaunch. The installation public URL -(`AGENTS_API_PUBLIC_URL`, from the installer's `--public-url`) is the one origin for +(`public_url` in the installation's `config.json`, seeded by `--public-url`, and +`AGENTS_API_PUBLIC_URL` for Core) is the one origin for applications, nodes, sandbox guests and self-hosted executors, and also the console origin. Core derives the daemon `wss` URL, the self-hosted `remote_url`, hosted Runtime bootstrap and the deployment's read-only `core_url` from it; the deployment 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/README.md b/README.md index a1dcf3099..de8d84ab8 100644 --- a/README.md +++ b/README.md @@ -38,7 +38,7 @@ retain their node across disconnects and resume. for obtaining/building a matching bundle and the host/network prerequisites. 2. **Sign in to Web.** Open the console address printed by the installer and sign in with the [Core key](docs/getting-started/operations.md#core-key) from - `~/.parsar/core/admin/core.key`. Keep it private. The console connects to Core + `~/.parsar/core/secrets/core.key`. Keep it private. The console connects to Core automatically. Create a Project on the **Projects and keys** page, then issue a key within it for your application. Save the one-time plaintext response privately; Core stores its digest. Rotate by issuing another key in the @@ -68,7 +68,9 @@ run an optional API example with your own model credentials. Run these from an extracted distribution. The plain command uses loopback for local console/API access. For node enrollment, use the reachable origin described -above; the installer does not change an existing installation's public URL. +above. Flags only seed the installation's `config.json`; later changes go there +and take effect with `~/.parsar/core/parsar apply`, and `parsar status`, `start` +and `stop` replace `install.sh --status` and `--stop`. Installing a local provider is optional, needs that HTTPS `--public-url`, and is not required for adding nodes in Web. diff --git a/deploy/install/config.schema.json b/deploy/install/config.schema.json new file mode 100644 index 000000000..1c9d07f0a --- /dev/null +++ b/deploy/install/config.schema.json @@ -0,0 +1,255 @@ +{ + "$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"], + "native_restarts": ["core", "web"], + "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..e123ac073 --- /dev/null +++ b/deploy/install/config_model.py @@ -0,0 +1,265 @@ +"""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 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", "native_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): + """Core's ValidateSandboxCoreURL rule, through the installer's one implementation of it.""" + from configuration import valid_core_origin # configuration imports this module at load time + return valid_core_origin(value) and (not https_only or value.startswith("https://")) + + +_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"] or type(value) is not type(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) + # With native Core some settings restart more; native_restarts names them all. + restarts = annotation(node, "native_restarts" if config.get("native_core") and + annotation(node, "native_restarts") else "restarts", []) + item.update({"default": node.get("default"), "changeable": annotation(node, "changeable", True), + "sensitive": sensitive, + "restarts": [name for name in 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 b79f653ce..6e76a8ca9 100644 --- a/deploy/install/configuration.py +++ b/deploy/install/configuration.py @@ -1,10 +1,24 @@ -"""Deployment files for the existing Core, Runtime and production console.""" +"""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 ipaddress -import os -import re -import stat +import json from pathlib import Path -from urllib.parse import urlsplit +import re +from urllib.parse import urlencode, urlsplit + +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")) _HOST_LABEL = re.compile(r"[a-z0-9]([a-z0-9-]{0,61}[a-z0-9])?") @@ -34,7 +48,10 @@ def valid_core_origin(value): if port and not (port.isdigit() and str(int(port)) == port and 1 <= int(port) <= 65535): return False try: - loopback = ipaddress.ip_address(host).is_loopback + address = ipaddress.ip_address(host) + # Go's IsLoopback also counts an IPv4-mapped loopback address. + mapped = getattr(address, "ipv4_mapped", None) + loopback = address.is_loopback or bool(mapped and mapped.is_loopback) except ValueError: if netloc.startswith("[") or len(host) > 253 or not all(_HOST_LABEL.fullmatch(label) for label in host.split(".")): return False @@ -42,9 +59,10 @@ def valid_core_origin(value): return parsed.scheme == "https" or loopback -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")): @@ -54,35 +72,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 and result.get("AGENTS_API_SANDBOX_INSTALLATION_ID") != state["installation_id"]: - raise RuntimeError("Core configuration installation identity differs") return result @@ -90,78 +89,212 @@ def bind(source, target, readonly=True): return {"type": "bind", "source": str(source).replace("$", "$$"), "target": target, "read_only": readonly} -def core_environment(root, state): - 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" +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 = { - # The password stays in its own file; the URL never carries it. - "AGENTS_API_DATABASE_URL": f"postgres://agents_api@{database}/agents_api?sslmode=disable", - "AGENTS_API_DATABASE_PASSWORD_FILE": config + "/database.password", - "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", - # The one origin applications, nodes, sandboxes and self-hosted executors - # use. Without a public URL, only this host reaches Core's loopback port. - "AGENTS_API_PUBLIC_URL": state.get("public_url") or f'http://127.0.0.1:{state["core_port"]}', - "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"] + 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" + 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) - config = root / "config" + 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): + root = Path(root) + 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"], + # Compose interpolates every string, extension fields included. + "x-parsar": {"generated_from": str(root / "config.json").replace("$", "$$"), + "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")) + 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"]} + 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} + 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 + + +LABEL = "io.parsar.inputs" + + +class Rendered: + """Generated file contents, plus the inputs digest each service must run with. + + The digest covers everything a service reads: its Compose definition, the + env_file content and the files and secrets it mounts. It is the service's + `io.parsar.inputs` label, or PARSAR_INPUTS in the native unit, so the running + services can be compared with a render at any time. + """ + + 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"} + external, 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"]] + external["core"] = 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"], + }, sort_keys=True) + if mode != "core-only": + external["web"] = json.dumps({"core.key": secrets["core.key"]}) + compose = compose_config(root, config, state) + 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 + external.get("core" if name == "migrate" else name, "")) + service["labels"] = {LABEL: services[name]} + files["compose.json"] = json.dumps(compose, indent=2) + "\n" + if native: + unit = native_service.unit_name(state) + services["core"] = sha256(core_env + native_service.unit_text(root, edit_hint(root)) + external["core"]) + files[unit] = native_service.unit_text(root, edit_hint(root), services["core"]) + return Rendered(files, services, unit) + + +def rendered_inputs(files, unit): + """The inputs digest per service that a set of generated files asks for.""" + result = {} + if files.get("compose.json"): + try: + for name, service in json.loads(files["compose.json"])["services"].items(): + result[name] = service.get("labels", {}).get(LABEL) + except (ValueError, KeyError, AttributeError): + return {} + if unit and files.get(unit): + text = files[unit].decode() if isinstance(files[unit], bytes) else files[unit] + match = re.search(r"^Environment=PARSAR_INPUTS=([0-9a-f]{64})$", text, re.M) + result["core"] = match[1] if match else None + return result diff --git a/deploy/install/convert.py b/deploy/install/convert.py new file mode 100644 index 000000000..7510135f5 --- /dev/null +++ b/deploy/install/convert.py @@ -0,0 +1,605 @@ +"""install.sh --convert: move an installation made before config.json to the new layout. + +It accepts the layout of the installers since #124 (deploy/install/configuration.py +at 5c3dcc16) and of #138 (b37e43b9), which changed only config/core.env: the public +URL, the database password file and no retired variables. + +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 ipaddress +import json +import os +from pathlib import Path +import stat +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_E2B_PROVIDER_BIN", "AGENTS_API_E2B_STATE_DIR") +# Written before #138 and retired by it; either layout may lack them. +RETIRED = ("AGENTS_API_DAEMON_WS_URL", "AGENTS_API_CONFIG_FILE") +# Written since #138. +PUBLIC_URL, PASSWORD_FILE = "AGENTS_API_PUBLIC_URL", "AGENTS_API_DATABASE_PASSWORD_FILE" + + +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 legacy138_core_environment(root, state): + """core.env as the #138 installer wrote it (b37e43b9).""" + 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" + return { + "AGENTS_API_DATABASE_URL": f"postgres://agents_api@{database}/agents_api?sslmode=disable", + PASSWORD_FILE: config + "/database.password", + PUBLIC_URL: state.get("public_url") or f'http://127.0.0.1:{state["core_port"]}', + } + + +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 origin(value): + return config_model.CHECKS["origin"][0](value) + + +def canonical(value, key, plan, override=None): + """An earlier installer kept the case it was given; lowercase it when that is the only problem.""" + if value is None or origin(value): + return value + if origin(value.lower()): + plan.notes.append(f"{key} {value} is written as {value.lower()}, the canonical form Core requires.") + return value.lower() + if override is None: + remedy = {"public_url": "rerun with --public-url naming the canonical origin", + "AGENTS_API_PUBLIC_URL": "write the canonical origin in config/core.env and rerun"}.get( + key, "install Web again with a canonical --core-url") + plan.problem("installation.json", key, f"{value} is not a canonical origin ({config_model.CHECKS['origin'][1]}); {remedy}") + return value + + +def loopback(value): + host = urlsplit(value).hostname or "" + try: + return ipaddress.ip_address(host).is_loopback + except ValueError: + return host == "localhost" + + +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): + # By path: an interrupted conversion may already have disabled the unit. + run(["systemctl", "--user", "enable", "--now", str(root / "config" / 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" + try: + info = os.lstat(root) + if not stat.S_ISDIR(info.st_mode) or stat.S_IMODE(info.st_mode) & 0o077 or info.st_uid != os.geteuid(): + raise OSError + except OSError: + plan.problem(str(root), "", "must be a directory with mode 0700, owned by you, and not a link") + private = ["admin/core.key"] + ([] if mode == "web-only" else ["config/credential.key", "config/database.password"]) + for name in private: + try: + parsar_cli.check_private(root / name, name) + except parsar_cli.ParsarError: + plan.problem(name, "", "must be a private regular file (mode 0600, owned by you, not a link); fix it and rerun") + if plan.problems: + return old, plan, False + 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": canonical(old.get("core_url"), "web.core_url", plan)} + config["log"] = {"level": "info", "format": "auto", "add_source": False} + env_public, from_env = None, False + 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: + env_public, from_env = 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.") + + if from_env: + public = env_public + if mode == "all" and public != old.get("public_url"): + # Web's origin follows public_url, so a hand-set Core address would move it. + if public_url_override is not None and public_url_override in (public, old.get("public_url")): + public, public_url_override = public_url_override, None + else: + plan.problem("config/core.env", PUBLIC_URL, f"is {public or 'the loopback fallback'}, but " + f"installation.json's public_url, which Web uses, is {old.get('public_url') or 'none'}. " + f"Rerun with --public-url naming the one to keep" + + ("" if public and old.get("public_url") else ", or make them agree in config/core.env")) + else: + public = canonical(old.get("public_url"), "public_url", plan, public_url_override) + if public_url_override is not None and old.get("public_url") and not origin(old["public_url"].lower()): + public, public_url_override = public_url_override, None + started = False + # Before #138, nodes enrolled with the sandbox deployment's own core_url. Since then + # Core derives that from AGENTS_API_PUBLIC_URL, so core.env already names the address. + if mode != "web-only" and not from_env and not plan.problems: + deployment, started = old_core_deployment(root, old, plan, run) + if deployment == f'http://127.0.0.1:{old["core_port"]}': + deployment = None # Core's loopback address, saved by a local-only setup: no public URL. + 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"] + plan.deletions + 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 if target != "secrets/core.key"}, + "core_installation_id": None, "generated": {}, "local_node": None, + "converted_from": {"source_commit": old["source_commit"], + "at": datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), + "remove": [name for name in plan.deletions if name.startswith(("config/", "state/"))]}} + return old, plan, started + + +def map_environment(root, old, env, password, config, plan, known): + """Map config/core.env into config.json. + + Returns the public URL Core uses when core.env names it (AGENTS_API_PUBLIC_URL, + since #138), None for Core's loopback fallback, and whether it did. + """ + before, since = legacy_core_environment(root, old, password), legacy138_core_environment(root, old) + where = "config/core.env" + for name in DERIVED: + if name in env and env[name] != before[name]: + plan.problem(where, name, "differs from the value the installer generated, and config.json can't hold " + "it; undo the edit") + for name in RETIRED: + if name in env and env[name] != before[name]: + plan.problem(where, name, "is retired and was edited; remove the line") + if PASSWORD_FILE in env and env[PASSWORD_FILE] != since[PASSWORD_FILE]: + plan.problem(where, PASSWORD_FILE, "names another file; only the installation's config/database.password " + "can be converted") + public, from_env = None, PUBLIC_URL in env + fallback = f'http://127.0.0.1:{old["core_port"]}' + if from_env and env[PUBLIC_URL] == since[PUBLIC_URL]: + # The #138 installer's own value: the public URL, or the loopback fallback for none. + public = old.get("public_url") + elif from_env: + # Set by hand when Core was upgraded in place: it is the address Core uses. + public = canonical(env[PUBLIC_URL], "AGENTS_API_PUBLIC_URL", plan) + if public == fallback: + public = None + elif public and configuration.valid_core_origin(public) and loopback(public): + plan.problem(where, PUBLIC_URL, f"is a loopback address other than Core's own {fallback}; set the public " + "HTTPS origin there, or remove the line") + plan.notes.append(f"public_url is taken from AGENTS_API_PUBLIC_URL in config/core.env, the address Core uses.") + for name in sorted(set(env) - set(DERIVED) - set(RETIRED) - {PUBLIC_URL, PASSWORD_FILE} - MAPPED): + plan.problem(where, name, "is not a setting config.json can hold; remove it") + base, _, query = env.get("AGENTS_API_DATABASE_URL", "").partition("?") + parameters = parse_qsl(query, keep_blank_values=True) + with_password = before["AGENTS_API_DATABASE_URL"].partition("?")[0] + without_password = since["AGENTS_API_DATABASE_URL"].partition("?")[0] + if base not in (with_password, without_password) or ("sslmode", "disable") not in parameters: + plan.problem(where, "AGENTS_API_DATABASE_URL", "names another database, user or password; " + "only the installation's own PostgreSQL can be converted") + elif base == without_password and PASSWORD_FILE not in env: + plan.problem(where, "AGENTS_API_DATABASE_URL", f"has no password and {PASSWORD_FILE} is not set") + pool = {} + for name, value in parameters: + if name == "sslmode": + continue + if name not in POOL or POOL[name] in pool: + plan.problem(where, "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()) + if path.parent in (root / "config", root / "state/e2b"): + # config.json holds the settings now; the old copy may hold export credentials. + plan.deletions.append(str(path.relative_to(root))) + known["config"].add(path.name) + plan.notes.append(f"Runtime history settings move from {path} into config.json; the old file is removed.") + else: + plan.notes.append(f"Runtime history settings from {path} are now in config.json. Delete {path}, " + "which may hold export credentials; it is a second copy.") + except (AttributeError, OSError, ValueError): + plan.problem("config/core.env", "AGENTS_API_RUNTIME_HISTORY_FILE", "names a file that can't be read as JSON") + if env.get("AGENTS_API_EXECUTION_OPTIONS_FILE"): + # This release retires the operator options file; Core refuses to start while it is set. + path = host_path(root, old, env["AGENTS_API_EXECUTION_OPTIONS_FILE"]) + if path is not None and path.parent == root / "config": + known["config"].add(path.name) + plan.notes.append(f"AGENTS_API_EXECUTION_OPTIONS_FILE is retired by this release and is not carried over. " + f"Hosted Sessions without a model provider of their own now need a deployment default: " + f"set it per harness in Web (System) or with PUT /core/v1/harnesses/{{harness}}/model-provider. " + f"{path or env['AGENTS_API_EXECUTION_OPTIONS_FILE']} is left in place; it may hold model keys, " + f"so delete it once the defaults are set.") + 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 public, from_env + + +# 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 check_resumed_public_url(root, value): + """A resumed conversion keeps the public URL it chose; --public-url may only repeat it.""" + if value is not None: + chosen = json.loads((root / "config.json").read_text()).get("public_url") + if value != chosen: + raise ConvertError(f"This conversion already set public_url to {chosen}; rerun without --public-url, " + "and change it in config.json after the conversion") + + +def resume(root, state, 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)) + deletions += [name for name in state["converted_from"].get("remove", []) if name not in deletions] + 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(): + state = json.loads((root / "state.json").read_text()) + if state["source_commit"] != manifest["source_commit"]: + raise ConvertError("Finish the conversion with the bundle it started with, release " + state["source_commit"]) + check_resumed_public_url(root, public_url_override) + load_images(state["images"]) + resume(root, state, 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 bbccf60d6..a9d02c210 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,20 @@ 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, valid_core_origin +import config_model +import configuration +from configuration import valid_core_origin +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 +44,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 +64,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,30 +97,6 @@ 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 loopback_origin(value): - hostname = urlsplit(value or "").hostname - try: - return ipaddress.ip_address(hostname).is_loopback - except ValueError: - return hostname == "localhost" - - -def origin_port(value): - parsed = urlsplit(value) - return parsed.port or (443 if parsed.scheme == "https" else 80) - - def public_origin(value): # Normalize case and a trailing slash, then apply Core's exact origin rule. try: @@ -139,53 +115,122 @@ def public_origin(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=public_origin, 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=public_origin, help="Web-only: origin of the existing Core") + parser.add_argument("--public-url", type=public_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): - if urlsplit(args.public_url or "").scheme != "https" or loopback_origin(args.public_url): - 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") + 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") - 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") + 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 compose(root, *args, **kwargs): - return run(["docker", "compose", "-f", str(root / "compose.json"), *args], **kwargs) +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 loopback_origin(value): + hostname = urlsplit(value or "").hostname + try: + return ipaddress.ip_address(hostname).is_loopback + except ValueError: + return hostname == "localhost" + + +def origin_port(value): + parsed = urlsplit(value) + return parsed.port or (443 if parsed.scheme == "https" else 80) + + +def check_public_url(config, provider): + if provider: + if urlsplit(config["public_url"] or "").scheme != "https" or loopback_origin(config["public_url"]): + raise InstallError("Local sandbox installation requires public_url with HTTPS reachable from sandbox " + "guests; loopback origins cannot be used") + + +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(): @@ -195,125 +240,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))) - 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" @@ -330,6 +278,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) @@ -347,102 +297,233 @@ def prepare_node_payload(root, state, bundle): os.unlink(temporary) -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) +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 - 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 - manifest = verify_bundle(bundle) - if args.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) - public_url = state.get("public_url") - if state["mode"] != "core-only": - url = f'http://127.0.0.1:{state["web_port"]}' - host = urlsplit(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") + 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" + if config or legacy: + return "config" if config else "legacy" + generated = root / "generated" + never_applied = not (generated.is_dir() and any(generated.iterdir())) + if (root / "state.json").exists(): + try: + recorded = json.loads((root / "state.json").read_text()).get("generated") + except (OSError, ValueError, AttributeError): + recorded = True + # Never applied: nothing started, so nothing depends on these secrets yet. + return "incomplete" if never_applied and not recorded else "missing-config" + if never_applied and {path.name for path in root.iterdir()} <= {"secrets", "generated", "state", ".parsar.lock"}: + return "incomplete" + return "other" + + +def create(root, args, config, manifest, images, provider=None): + """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, + "converted_from": None, "generated": {}, + # A requested local node is enrolled once the services first start; a repair retries it. + "local_node": provider} + state["secrets_sha256"].pop("core.key") + # state.json first: whenever config.json exists, the installation can be repaired. + write(root / "state.json", json.dumps(state, indent=2) + "\n") + write(root / "config.json", json.dumps(config, indent=2) + "\n") + + +def unfinished_conversion(state): + return bool(state.get("converted_from")) and not state["converted_from"].get("finished") + + +def finish(root, bundle, manifest, fresh=False): + """Put the bundle's files in place, apply config.json and start the services.""" + state = parsar_cli.load_state(root) + converting = unfinished_conversion(state) + # A converted installation replaces the earlier release's binaries and payload once. + prepare_node_payload(root, state, bundle, converting) + native_service.prepare(root, state, bundle, converting) + install_parsar(root, bundle) + retry = f"rerun ./install.sh {'--convert ' if converting else ''}--install-dir {root}" + parsar_cli.apply(root, start=True, retry=retry) + config = parsar_cli.load_config(root) + mode = config["mode"] + # An earlier-release Core has no /core/v1/installation (404); apply noted it and Web still works. + if mode == "web-only" and parsar_cli.paired_core(root, config)[0] not in (200, 404): + raise InstallError("Core key authentication failed. Inspect secrets/core.key and web.core_url; no model was called") + provider = state.get("local_node") + if provider: + local_node.install(root, dict(state, provider=provider, core_port=config["ports"]["core"], + public_url=config["public_url"]), manifest, bundle, run) + state = parsar_cli.load_state(root) + if provider or converting: + state["local_node"] = None + if converting: + state["converted_from"] = dict(state["converted_from"], finished=True) + parsar_cli.save_state(root, state) + summary(root, config, provider, fresh) + + +def summary(root, config, provider, fresh): + mode, public_url, ports = config["mode"], config["public_url"], config["ports"] + if mode != "core-only": # The console accepts only its configured origin, so a public URL has no loopback console. - console = public_url or url + console = public_url or f'http://127.0.0.1:{ports["web"]}' print("Console: " + console + (" (local only)" if loopback_origin(console) else "")) - if state["mode"] != "web-only": - api = f'http://127.0.0.1:{state["core_port"]}/v1' + if mode != "web-only": + api = f'http://127.0.0.1:{ports["core"]}/v1' if public_url and not loopback_origin(public_url): print("API base URL: " + public_url + "/v1") print("Local-only API on this host: " + api) - elif public_url and origin_port(public_url) != state["web_port"]: + elif public_url and origin_port(public_url) != ports.get("web"): print("API base URL: " + public_url + "/v1 (local only)") else: # Web answers 404 on /v1, so only Core's own port serves the API locally. print("API base URL: " + api + " (local only)") - core_key = root / "admin/core.key" - if state["mode"] == "core-only": + core_key = root / "secrets/core.key" + if mode == "core-only": print(f'Next: create a Project and its API key through the Core management API at ' - f'http://127.0.0.1:{state["core_port"]}/core/v1 (local only) with the Core key in {core_key}.') + f'http://127.0.0.1:{ports["core"]}/core/v1 (local only) with the Core key in {core_key}.') else: print(f"Next: sign in to Web with the Core key in {core_key}, then create a Project and its API key on the Projects and keys page.") print("Keep the Core key private; it also authorizes the Core management API.") - if state["mode"] != "web-only": - print("Core configuration file: " + str(root / "config/core.env")) - if args.provider: - print("Provider: " + args.provider + ". Local node enrolled; Core provisions Sessions on demand.") - elif state["mode"] == "all": - print("No execution node was installed by this run. Choose a sandbox backend and add nodes on the Nodes page in Web.") - else: - print("No execution node was installed by this run.") + 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 provider: + print("Provider: " + provider + ". Local node enrolled; Core provisions Sessions on demand.") + # Only a new installation says it has no nodes; a repaired or converted one keeps its own. + elif fresh and mode == "all": + print("No execution node was installed by this run. Choose a sandbox backend and add nodes on the Nodes page in Web.") + elif fresh and mode == "core-only": + print("No execution node was installed by this run.") 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") + 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) + described = convert.detect(root) if kind == "legacy" else json.loads((root / "state.json").read_text()) + if described.get("native_core") and described.get("mode") != "web-only": + # A host that can't run this release's native Core is refused before anything changes. + native_service.preflight(bundle, root) + convert.convert(root, manifest, image_loader(manifest, bundle), args.public_url, args.yes, run) + finish(root, bundle, manifest) + return + if kind == "config" and args.convert: + state = parsar_cli.load_state(root) + # The layout is converted, but the new release has not started successfully yet. + if not unfinished_conversion(state): + raise InstallError(f"{root} already uses config.json; edit it and run {root / 'parsar'} apply") + if set(args.given) - {"convert", "yes", "public_url"}: + raise InstallError("Finishing a conversion accepts only --install-dir, --yes and --public-url") + convert.check_resumed_public_url(root, args.public_url) + check_host() + manifest = verify_bundle(bundle) + if state["source_commit"] != manifest["source_commit"]: + raise InstallError("Finish the conversion with the bundle it started with, release " + state["source_commit"]) + if native_service.is_native(state): + native_service.preflight(bundle, root) + image_loader(manifest, bundle)(list(state["images"])) + 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 == "missing-config": + raise InstallError(f"{root / 'config.json'} is missing. Restore it from a backup; " + f"{root / 'generated/settings.json'} lists the last applied values. The secrets and " + "database belong to this installation, so keep the directory. Nothing was changed") + if kind == "incomplete": + raise InstallError(f"An earlier installation into {root} stopped before writing config.json and started no " + "service. Remove the directory and install again") + 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 config.get("native_core"): + native_service.preflight(bundle, root) + 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, provider) + finish(root, bundle, manifest, fresh=True) + + 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..78dfbade2 --- /dev/null +++ b/deploy/install/installer_fakes.py @@ -0,0 +1,297 @@ +"""A fake Docker, systemd and HTTP host for installer tests. Nothing real is started. + +Compose is modelled by its observable contract: `up` recreates a container exactly +when its resolved configuration (including env_file content) changed, each container +keeps the labels it was created with, and Core loads its Core key digests when it +starts. systemd runs the unit it last loaded; its process keeps the environment it +started with. +""" +import hashlib +import json +from pathlib import Path +import re +import subprocess +from types import SimpleNamespace +from unittest import mock + +import native_service +import parsar_cli + +LABEL = "io.parsar.inputs" + + +def sha256(data): + return hashlib.sha256(data.encode() if isinstance(data, str) else data).hexdigest() + + +class FakeHost: + def __init__(self, test): + self.commands, self.requests = [], [] + self.containers = {} # service -> {hash, inputs, running} + self.project = None + self.recreated = [] + self.native = {"active": False, "starts": 0, "restarts": 0, "reloads": 0, "addr": None, "digests": [], + "inputs": None, "loaded": None, "environment": ""} + # fails: Core never starts; rejects(core.env text): Core refuses that configuration. + self.core = {"port": None, "digests": [], "fails": False, "log": "", "rejects": lambda environment: False} + 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, name, value in ((subprocess, "run", mock.Mock(side_effect=self.run)), + (parsar_cli, "http", self.http), + (parsar_cli, "time", SimpleNamespace(sleep=lambda seconds: None)), + (native_service, "_process_environment", self.process_environment)): + patcher = mock.patch.object(target, name, value) + patcher.start() + test.addCleanup(patcher.stop) + + def running(self): + return {name for name, item in self.containers.items() if item["running"]} + + def run_container(self, name, port=None, digests=None): + """A container of an installation made outside the test, such as an earlier release.""" + self.containers[name] = {"hash": "external", "inputs": None, "running": True} + if name == "core": + self.core.update(port=port, digests=digests or []) + + # 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[:2] == ["docker", "ps"]: + stdout = "\n".join("id-" + name for name in self.containers) + elif args[:2] == ["docker", "inspect"]: + stdout = "\n".join(f'{name}\t{self.containers[name]["inputs"] or ""}\t' + f'{"running" if self.containers[name]["running"] else "exited"}\t' + for name in (item[3:] for item in args[4:]) if name in self.containers) + 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] == ["stop"]: + for name in args[1:] or services: + if name in self.containers: + self.containers[name]["running"] = False + return 0, "" + if args[:1] == ["restart"]: + for name in args[1:]: + self.containers[name]["running"] = True + if name == "core": + self.load_core(path.parent.parent, services[name]) + return 0, "" + if args[:1] == ["logs"]: + return 0, self.core["log"] + 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", *targets) if name in services and + (name in targets or name in ("database", "migrate") and "core" in targets)] + for name in names: + digest = self.service_hash(document, name) + current = self.containers.get(name) + if current is None or current["hash"] != digest: + self.containers[name] = {"hash": digest, "inputs": services[name].get("labels", {}).get(LABEL), + "running": False} + self.recreated.append(name) + if name == "core": + self.load_core(path.parent.parent, services[name]) + if name != "migrate": + self.containers[name]["running"] = True + 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]) + environment = path.parent / "core.env" + if "core" in names and (self.core["fails"] or + self.core["rejects"](environment.read_text() if environment.exists() else "")): + self.containers["core"]["running"] = False + self.core["failed"] = True + 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()) + environment = root / "generated/core.env" + self.core["environment"] = environment.read_text() if environment.exists() else "" + + def unit_file(self): + found = sorted(self.native_root.glob("generated/parsar-*-core.service")) + \ + sorted(self.native_root.glob("config/parsar-*-core.service")) + return found[0].read_text() if found else "" + + def systemctl(self, args): + native = self.native + if args[0] == "daemon-reload": + native["reloads"] += 1 + native["loaded"] = self.unit_file() + elif args[0] in ("enable", "restart", "start"): + if args[0] != "restart" and native["active"]: + return 0, "" # systemd leaves an active unit alone on start and enable --now + native["starts" if args[0] != "restart" else "restarts"] += 1 + if native["loaded"] is None: + native["loaded"] = self.unit_file() + environment = self.native_root / "generated/core.env" + refused = self.core["fails"] or self.core["rejects"](environment.read_text() if environment.exists() else "") + native["active"] = not refused + self.core["failed"] = self.core.get("failed") or refused + match = re.search(r"^Environment=PARSAR_INPUTS=(\w+)$", native["loaded"], re.M) + native["inputs"] = match[1] if match and native["active"] else None + root = self.native_root + for environment, digests in ((root / "generated/core.env", root / "generated/core-key-digests.json"), + (root / "config/core.env", root / "admin/core-key-digests.json")): + if environment.exists(): + address = next(line for line in environment.read_text().splitlines() + if line.startswith("AGENTS_API_ADDR=")) + native["addr"] = int(address.rsplit(":", 1)[1].rstrip('"')) + native["digests"] = json.loads(digests.read_text()) + native["environment"] = environment.read_text() + break + elif args[0] in ("stop", "disable"): + native["active"], native["inputs"] = False, None + elif args[0] == "is-active": + return (0 if native["active"] else 3), "" + elif args[0] == "show" and "--property=MainPID" in args: + return 0, "4242" if native["active"] else "0" + elif args[0] == "show": + return 0, "252" + return 0, "" + + def process_environment(self, pid): + return {"PARSAR_INPUTS": self.native["inputs"]} if self.native["inputs"] else {} + + # HTTP ------------------------------------------------------------------- + def core_listening(self, port): + if self.containers.get("core", {}).get("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): + self.requests.append(url) + 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) + web_up = self.containers.get("web", {}).get("running") and self.web_port == port + if path == "/healthz": + return (200, b"{}") if digests is not None or web_up else (0, b"") + if path == "/console/auth": + return (200, b"{}") if web_up else (0, b"") + if digests is None: + return 0, b"" + if sha256(key) not in digests: + return 401, b"" + if path == "/core/v1/installation": + native = self.native["active"] and self.native["addr"] == port + environment = self.native["environment"] if native else self.core.get("environment", "") + match = re.search(r'^AGENTS_API_PUBLIC_URL="([^"]+)"$', environment, re.M) + public = match[1] if match else None + return 200, json.dumps({"installation_id": self.core_installation_id, "public_url": public, + "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 8fd64b31b..2c4350574 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/model_provider_sessions.py b/deploy/install/model_provider_sessions.py index 015c3b042..ec9c5a969 100644 --- a/deploy/install/model_provider_sessions.py +++ b/deploy/install/model_provider_sessions.py @@ -30,7 +30,9 @@ class CheckError(Exception): def count(root): """Return {environment: count} for Sessions without a frozen provider.""" - compose = Path(root) / "compose.json" + # An installation made with config.json keeps its Compose file under generated/. + compose = next((path for path in (Path(root) / "generated/compose.json", Path(root) / "compose.json") + if path.exists()), Path(root) / "compose.json") try: services = json.loads(compose.read_text())["services"] except (OSError, ValueError, KeyError, TypeError): @@ -46,7 +48,7 @@ def count(root): except (OSError, subprocess.SubprocessError): raise CheckError("Cannot run the database query; is Docker available?") from None if result.returncode: - raise CheckError("The database query failed; start the installation (install.sh --status shows its state) and retry") + raise CheckError("The database query failed; start the installation and retry") counts = dict.fromkeys(ENVIRONMENTS, 0) for line in result.stdout.splitlines(): if not line.strip(): diff --git a/deploy/install/native_service.py b/deploy/install/native_service.py index 9a8eb92bd..ab7a63ab6 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,86 @@ 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, inputs=None): + """The unit. PARSAR_INPUTS carries the inputs digest the running Core started with.""" + 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" + + ("Environment=PARSAR_INPUTS=" + inputs + "\n" if inputs else "") + + "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 _process_environment(pid): + try: + raw = Path(f"/proc/{pid}/environ").read_bytes() + except OSError: + return {} + return dict(item.split("=", 1) for item in raw.decode(errors="replace").split("\0") if "=" in item) + + +def running_inputs(state): + """PARSAR_INPUTS of the running Core process, or None when it is not running.""" + if not is_native(state): + return None + result = _run(["systemctl", "--user", "show", "--property=MainPID", "--value", unit_name(state)], + "Cannot query this installation's native Core service") + pid = result.stdout.strip() + if result.returncode or not pid.isdigit() or pid == "0": + return None + return _process_environment(int(pid)).get("PARSAR_INPUTS") 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..8cb73e946 --- /dev/null +++ b/deploy/install/parsar_cli.py @@ -0,0 +1,809 @@ +"""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 renders generated/ from it and converges the running services on that render. +What runs is the truth: each Compose container carries the inputs digest it was +created with (label io.parsar.inputs) and native Core carries it in PARSAR_INPUTS, +so an interrupted apply, rotation or rollback is finished by the next apply. +""" +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 check_directories(root): + """The installation, secrets/ and generated/ are real private directories of this user.""" + for path, what in ((root, str(root)), (root / "secrets", "secrets/"), (root / "generated", "generated/")): + try: + info = os.lstat(path) + except FileNotFoundError: + raise ParsarError(f"{what} is missing") from None + if (not stat.S_ISDIR(info.st_mode) or stat.S_IMODE(info.st_mode) & 0o077 + or info.st_uid != os.geteuid()): + raise ParsarError(f"{what} must be a directory with mode 0700, 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) + + +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) + + +# What runs --------------------------------------------------------------------- + +def compose(root, *args, **kwargs): + return run(["docker", "compose", "-f", str(root / "generated/compose.json"), *args], **kwargs) + + +INSPECT = ('{{index .Config.Labels "com.docker.compose.service"}}\t{{index .Config.Labels "' + configuration.LABEL + + '"}}\t{{.State.Status}}\t{{if .State.Health}}{{.State.Health.Status}}{{end}}') + + +def observe(state): + """{service: {running, inputs, health}} of this installation's containers and native Core.""" + result = {} + ids = run(["docker", "ps", "-aq", "--filter", f'label=com.docker.compose.project={state["project"]}', + "--filter", "label=com.docker.compose.oneoff=False"], capture_output=True, text=True).stdout.split() + if ids: + for line in run(["docker", "inspect", "--format", INSPECT, *ids], capture_output=True, text=True).stdout.splitlines(): + service, inputs, status, health = (line.split("\t") + ["", "", "", ""])[:4] + if service: + result[service] = {"running": status == "running", "inputs": inputs or None, "health": health} + if native_service.is_native(state): + inputs = native_service.running_inputs(state) + result["core"] = {"running": inputs is not None or native_service.active(state), "inputs": inputs, "health": ""} + return result + + +def stale(actual, desired, will_run): + """Services that should run but don't, or run with other inputs.""" + return {name for name in will_run if name != "migrate" and ( + not actual.get(name, {}).get("running") or actual[name]["inputs"] != desired.get(name))} + + +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 converge(root, state, desired, will_run, force=()): + """Bring every service in will_run to the desired inputs; Core first, then the rest. + + Compose recreates exactly the containers whose configuration (and so label) + differs; native Core restarts when its running PARSAR_INPUTS differ. + """ + native = native_service.is_native(state) + actual = observe(state) + todo = stale(actual, desired, will_run) | set(force) + if not todo: + return set() + up = ["up", "--detach", "--wait", "--wait-timeout", "300"] + if "core" in todo: + if native: + if "database" in will_run and "database" in stale(actual, desired, {"database"}): + compose(root, *up, "database") + native_service.daemon_reload() + if actual.get("core", {}).get("running"): + native_service.restart(state) + else: + migrate_native(root) + native_service.start(root, state) + elif "core" in force and not stale(actual, desired, {"core"}): + compose(root, "restart", "core") + else: + compose(root, *up, "core") + containers = will_run - ({"core"} if native else set()) + if todo - {"core"} or ("core" in todo and not native): + if containers: + compose(root, *up) + for name in sorted(set(force) & containers - {"core"}): + if not stale(actual, desired, {name}): + compose(root, "restart", name) + return todo + + +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 + + +def describe(error): + """A failure message without command paths, output or environment.""" + if isinstance(error, subprocess.CalledProcessError): + # A command's own path shows as its name; other paths and options are left out. + words = [Path(word).name if index == 0 else word for index, word in enumerate(error.cmd) + if index == 0 or not word.startswith(("/", "-"))][:3] + return f"`{' '.join(words)}` failed" + if isinstance(error, KeyboardInterrupt): + return "it was interrupted" + return str(error) + + +# 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 view of what was last written.""" + 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") + + +# Files ------------------------------------------------------------------------- + +def check_fixed(config, state): + """state.json records mode and native_core at installation; it wins over config.json.""" + 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 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 comparable(name, data): + """File content as compared for edits; the snapshot's applied_at is not an edit.""" + if data is not None and name == "settings.json": + try: + document = json.loads(data) + document.pop("applied_at", None) + return json.dumps(document, sort_keys=True).encode() + except (ValueError, AttributeError): + return data + return data.encode() if isinstance(data, str) else data + + +def edited_files(state, disk, rendered): + """Generated files that match neither a digest parsar wrote nor the current render. + + state.json keeps the last few digests written for each file, so a run + interrupted between writing files and state.json is not taken for an edit. A + missing file is simply written again. + """ + edited = [] + for name, digests in sorted((state.get("generated") or {}).items()): + data = disk.get(name) + if data is None or configuration.sha256(data) in digests: + continue + if name in rendered.files and comparable(name, data) == comparable(name, rendered.files[name]): + continue + edited.append(name) + return edited + + +def record_digests(state, files): + generated = dict(state.get("generated") or {}) + for name, text in files.items(): + digest = configuration.sha256(text) + generated[name] = ([item for item in generated.get(name, []) if item != digest] + [digest])[-3:] + return dict(state, generated=generated) + + +def disk_view(disk): + """The settings last written, by key, and the snapshot's stamp. + + Sensitive values are not in the snapshot; the one sensitive setting is read back + from the runtime-history.json that carries it. + """ + try: + document = json.loads(disk.get("settings.json") or b"") + values = {item["key"]: item["value"] for item in document["settings"]} + except (ValueError, KeyError, TypeError): + return None, None + if "core.runtime_history.headers" in values: + history = disk.get("runtime-history.json") + try: + values["core.runtime_history.headers"] = json.loads(history).get("headers") if history else None + except (ValueError, AttributeError): + values["core.runtime_history.headers"] = None + return values, document.get("applied_at") + + +def render_now(root, config, state): + """The render of config.json, keeping the written stamp unless something changed.""" + names = set((state.get("generated") or {})) + disk = read_generated(root, names | {"settings.json", "runtime-history.json"}) + previous, stamp = disk_view(disk) + rendered = configuration.render(root, config, state, stamp or now()) + disk = read_generated(root, names | set(rendered.files) | {"runtime-history.json"}) + if any(comparable(name, disk.get(name)) != comparable(name, text) for name, text in rendered.files.items()) \ + or any(disk.get(name) is not None for name in names - set(rendered.files)): + rendered = configuration.render(root, config, state, now()) + return rendered, disk, previous + + +# Apply ----------------------------------------------------------------------- + +def old_public_url(root, config, previous, disk, actual): + """The public URL things are bound to: Core's own answer, else the written core.env.""" + port = (previous or {}).get("ports.core") or config["ports"]["core"] + if actual.get("core", {}).get("running"): + status, body = http(f"http://127.0.0.1:{port}/core/v1/installation", bearer(configuration.read_core_key(root))) + if status == 200: + return json.loads(body).get("public_url"), port, True + try: + written = configuration.read_environment((disk.get("core.env") or b"").decode()) + except RuntimeError: + written = {} + return written.get("AGENTS_API_PUBLIC_URL"), port, False + + +def confirm_public_url(root, config, old, port, core_answered, args, interactive, out): + """Changing the public URL strands what is bound to the old one; list it and confirm. + + old is None when it can't be read; then the change always needs confirmation. + """ + new = configuration.local_public_url(config) + if old == new: + return + counts, nodes = None, [] + if core_answered: + base, key = f"http://127.0.0.1:{port}", 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 {} + # The node list only names them; the count comes from Core's own bindings. + status, body = http(base + "/core/v1/sandbox/nodes", bearer(key)) + nodes = [node for node in (json.loads(body).get("data", []) if status == 200 else []) + if node.get("core_url") == old] + if old is None: + out(f"The public URL in use can't be read, so the change to {new} needs confirmation.") + else: + out(f"The public URL changes from {old} to {new}.") + if counts is not None: + bound = max(0, counts.get("nodes", 0) - counts.get("nodes_on_other_address", 0)) + out(f'Bound to the current address: {bound} node(s), {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 (bound or counts.get("hosted_sandboxes") or counts.get("self_hosted_executors")): + return + elif old is not None: + 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, previous, args, interactive, out): + """Web-only: detect a web.core_url that now reaches a different Core.""" + if args.dry_run and previous is not None and previous.get("web.core_url") != config["web"]["core_url"]: + # The Core key goes only to a Core that apply is asked to use. + out("web.core_url changes; apply checks which Core it reaches.") + return state.get("core_installation_id") + 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.") + elif status == 404: + out("Note: the paired Core runs an earlier release without /core/v1/installation. Convert or upgrade the " + "Core host, then run parsar apply here to record which Core Web is paired with.") + 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 apply(root, dry_run=False, yes=False, discard_edits=False, confirm_public_url_change=None, + start=False, interactive=None, out=print, retry=None): + 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, retry=retry) + + +def _apply(root, args, discard_edits, start, interactive, out, rollback=True, retry=None): + retry = retry or f"run {root / 'parsar'} apply again" + check_directories(root) + config = load_config(root) + state = load_state(root) + check_fixed(config, state) + check_secrets(root, config, state) + if not args.dry_run: + # A rotation that stopped before using its new key leaves only this file. + (root / "secrets/core.key.new").unlink(missing_ok=True) + rendered, disk, previous = render_now(root, config, state) + edited = edited_files(state, disk, rendered) + 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-