diff --git a/ouroboros/contracts/plugin_api.py b/ouroboros/contracts/plugin_api.py index ccb815a9a..a9cd3c972 100644 --- a/ouroboros/contracts/plugin_api.py +++ b/ouroboros/contracts/plugin_api.py @@ -438,11 +438,36 @@ class ExtensionRegistrationError(Exception): """Raised when a registration violates namespace, permission, or schema.""" +def normalize_extension_route_methods(methods: Any, *, subject: str) -> tuple[str, ...]: + """One route declaration's HTTP methods, normalized against the vocabulary. + + Used by the in-process ``register_route`` AND by the host-side re-check of + a child catalog's route descriptors, so the vocabulary cannot admit + out-of-process what it refuses in-process. A bare string is a single + method; case and surrounding space are normalized, order is preserved and + duplicates collapse. ``subject`` names the route in the refusal. + """ + declared = (methods,) if isinstance(methods, str) else (methods or ()) + normalized = tuple(dict.fromkeys( + str(method).strip().upper() for method in declared if str(method).strip() + )) + if not normalized: + raise ExtensionRegistrationError(f"{subject} methods must be non-empty") + unsupported = [m for m in normalized if m not in VALID_EXTENSION_ROUTE_METHODS] + if unsupported: + raise ExtensionRegistrationError( + f"{subject} methods {unsupported!r} are unsupported; " + f"expected subset of {sorted(VALID_EXTENSION_ROUTE_METHODS)}" + ) + return normalized + + __all__ = [ "PluginAPI", "ExtensionRegistrationError", "FORBIDDEN_SKILL_SETTINGS", "PLUGIN_API_VERSION", "LEGACY_PLUGIN_API_GENERATION", "PLUGIN_API_SURFACE_FINGERPRINTS", "PluginAPINegotiation", "RuntimeInfo", "VALID_EXTENSION_PERMISSIONS", "VALID_EXTENSION_ROUTE_METHODS", + "normalize_extension_route_methods", "ExecutionMode", "MATRIX_CAPABILITIES", "ALWAYS_AVAILABLE_CAPABILITIES", "OUT_OF_PROCESS_UNAVAILABLE_CAPABILITIES", "api_generation", "capability_available", "available_capabilities", diff --git a/ouroboros/extension_child_catalog.py b/ouroboros/extension_child_catalog.py index a468b6fb4..6ac77695c 100644 --- a/ouroboros/extension_child_catalog.py +++ b/ouroboros/extension_child_catalog.py @@ -9,19 +9,10 @@ publication snapshot before anything is installed. from __future__ import annotations -import os -import pathlib from typing import TYPE_CHECKING, Any, Dict -from ouroboros.contracts.plugin_api import ExtensionRegistrationError, VALID_EXTENSION_ROUTE_METHODS -from ouroboros.extension_registry_state import ( - _lock, - _routes, - _settings_sections, - _tools, - _ui_tabs, - _ws_handlers, -) +from ouroboros.contracts.plugin_api import ExtensionRegistrationError, normalize_extension_route_methods +from ouroboros.extension_registry_state import SURFACE_KINDS, _lock from ouroboros.extension_surface_names import ( _EXTENSION_NAME_RE, _widget_geometry_from_render, @@ -63,15 +54,19 @@ def _stage_out_of_process_surfaces( item["skills_repo_path"] = str(skill.skill_dir.parent) return item - kinds = ( - ("tools", _validate_child_tool_descriptor, "name", True, _tools, "tool"), - ("routes", _validate_child_route_descriptor, "path", True, _routes, "route"), - ("ws_handlers", _validate_child_ws_descriptor, "type", True, _ws_handlers, "ws handler"), - ("ui_tabs", _validate_child_ui_descriptor, None, False, _ui_tabs, "ui tab"), - ("settings_sections", _validate_child_settings_descriptor, None, False, _settings_sections, "settings section"), - ) + # Per-kind: the descriptor validator, and the descriptor field carrying the + # registry key (None = the child's own "key" field, for the UI kinds, which + # dispatch nothing and so need no handler proxy). + validators = { + "tools": (_validate_child_tool_descriptor, "name"), + "routes": (_validate_child_route_descriptor, "path"), + "ws_handlers": (_validate_child_ws_descriptor, "type"), + "ui_tabs": (_validate_child_ui_descriptor, None), + "settings_sections": (_validate_child_settings_descriptor, None), + } with _lock: - for kind, validate, key_field, proxied, live, label in kinds: + for kind, live, label in SURFACE_KINDS: + validate, key_field = validators[kind] staged = getattr(api._staged, kind) for raw in catalog.get(kind) or []: item = validate(skill.name, dict(raw or {})) @@ -82,7 +77,7 @@ def _stage_out_of_process_surfaces( if not key: continue api._stage_surface_locked( - live, staged, key, _proxy(item) if proxied else item, label, + live, staged, key, _proxy(item) if key_field else item, label, ) @@ -122,19 +117,9 @@ def _validate_child_tool_descriptor(skill_name: str, item: Dict[str, Any]) -> Di def _validate_child_route_descriptor(skill_name: str, item: Dict[str, Any]) -> Dict[str, Any]: path = str(item.get("path") or "") _validate_child_catalog_namespace(skill_name, "route", path) - methods_iter = item.get("methods") or ("GET",) - if isinstance(methods_iter, str): - methods_iter = (methods_iter,) - methods = tuple(dict.fromkeys(str(method).strip().upper() for method in methods_iter if str(method).strip())) - if not methods: - raise ExtensionRegistrationError(f"out-of-process route {path!r} methods must be non-empty") - invalid = [method for method in methods if method not in VALID_EXTENSION_ROUTE_METHODS] - if invalid: - raise ExtensionRegistrationError( - f"out-of-process route {path!r} methods {invalid!r} are unsupported; " - f"expected subset of {sorted(VALID_EXTENSION_ROUTE_METHODS)}" - ) - item["methods"] = methods + item["methods"] = normalize_extension_route_methods( + item.get("methods") or ("GET",), subject=f"out-of-process route {path!r}", + ) return item @@ -167,105 +152,3 @@ def _validate_child_settings_descriptor(skill_name: str, item: Dict[str, Any]) - raise ExtensionRegistrationError(f"out-of-process settings section {key!r} render must be an object") item["render"] = _validate_settings_schema(dict(item.get("render") or {})) return item - - -def skill_state_path(drive_root: pathlib.Path, name: str) -> pathlib.Path: - """The per-skill state-dir PATH without creating the directory. - - The generation-fenced companion recovery resolves this path BEFORE its - fence (fix-round-6); creation belongs to the post-fence attach in - ``_publish_registrations``, so a stale refusal creates no directories. - ``skill_state_dir`` remains the creating variant for load paths. - """ - from ouroboros.skill_loader import _sanitize_skill_name, _skills_state_root - - return _skills_state_root(pathlib.Path(drive_root)) / _sanitize_skill_name(name) - - -def companion_manifest_path_override(spec: Dict[str, Any]) -> bool: - """An explicit manifest PATH means the author owns the runtime lookup — - neither the bundled argv rewrite nor the emergency prepend may shadow it - (T14).""" - return any(str(key).upper() == "PATH" for key in (spec.get("env") or {})) - - -def companion_node_argv(spec: Dict[str, Any], expected_runtime: str, cmd: list) -> list: - """T14 node symmetry with the python->sys.executable rewrite: - platform_layer.select_skill_node_runtime owns the bundled-first precedence - (health rollback included). npm is never rewritten — the emergency PATH - prepend in ``materialize_companion_env`` covers its `#!/usr/bin/env node` - shebang. On a healthy PATH argv and env stay byte-identical; an npm - launcher rewritten to an ABSOLUTE node shebang ignores PATH and keeps - failing honestly (disclosed residual).""" - if expected_runtime not in {"node", "npm"} or companion_manifest_path_override(spec): - return cmd - from ouroboros.platform_layer import select_skill_node_runtime - - selected_node, _node_provenance = select_skill_node_runtime() - if selected_node and cmd and cmd[0] == "node": - return [selected_node, *cmd[1:]] - return cmd - - -def materialize_companion_env(api: "PluginAPIImpl", descriptor: Any, spec: Dict[str, Any], token: str) -> None: - """Fill one staged companion descriptor's env at publication (fix-round-6). - - Runs only inside ``_publish_registrations``'s post-swap attach, after the - generation fence admitted the publication: the settings-derived values - (``_scrub_env`` -> ``load_settings`` takes the settings lock and may - persist a settings migration), the manifest env overlay, the Host Service - bridge URL/token and the isolated-dep PYTHONPATH are all filled in HERE, - keeping the pre-fence descriptor build purely computational. - """ - from ouroboros.contracts.plugin_api import FORBIDDEN_SKILL_SETTINGS - from ouroboros.extension_isolated_deps import _isolated_python_site_dirs - from ouroboros.gateway.host_service import DEFAULT_HOST_SERVICE_HOST, host_service_port - from ouroboros.tools.skill_exec import _scrub_env - - env = _scrub_env( - list(api._env_allow), - api._state_dir, - api._skill, - granted_keys=list(api._granted_upper), - ) - reserved_env = {"HOST_SERVICE_TOKEN", "HOST_SERVICE_URL"} - manifest_env = { - str(key): str(value) - for key, value in (spec.get("env") or {}).items() - if str(key).upper() not in FORBIDDEN_SKILL_SETTINGS - and str(key).upper() not in reserved_env - } - # Case-aware merge (delta finding D2-8): a manifest "Path" must REPLACE - # the allowlisted "PATH" on Windows, never sit next to it — duplicate - # case-variant env keys make CreateProcess-era spawns fail or pick an - # undefined winner. Same contract as the executor-local service lane. - from ouroboros.workspace_executor import overlay_env - - env = overlay_env(env, manifest_env) - env["HOST_SERVICE_URL"] = f"http://{DEFAULT_HOST_SERVICE_HOST}:{host_service_port()}" - env["HOST_SERVICE_TOKEN"] = token - if api._skill_dir is not None: - site_dirs = [str(path) for path in _isolated_python_site_dirs(api._skill_dir)] - if site_dirs: - existing_pythonpath = env.get("PYTHONPATH") - env["PYTHONPATH"] = os.pathsep.join( - [*site_dirs, existing_pythonpath] if existing_pythonpath else site_dirs - ) - expected_runtime = str(spec.get("runtime") or "").strip() - if expected_runtime in {"node", "npm"} and not companion_manifest_path_override(spec): - # T14 emergency-only PATH prepend (see register_companion_process): - # descriptor env keys win over the supervisor's `_companion_base_env` - # merge, so the prepended PATH reaches the child and survives - # supervisor restarts. - from ouroboros.platform_layer import skill_node_emergency_path_dir - - node_prepend_dir = skill_node_emergency_path_dir() - if node_prepend_dir: - existing_path = env.get("PATH") or os.environ.get("PATH", "") - env["PATH"] = ( - os.pathsep.join([node_prepend_dir, existing_path]) - if existing_path - else node_prepend_dir - ) - descriptor.env.clear() - descriptor.env.update(env) diff --git a/ouroboros/extension_loader.py b/ouroboros/extension_loader.py index 7fdf35c1d..087f2f6a5 100644 --- a/ouroboros/extension_loader.py +++ b/ouroboros/extension_loader.py @@ -54,7 +54,6 @@ from ouroboros.extension_isolated_deps import _isolated_python_site_dirs, async_ from ouroboros.extension_child_catalog import ( _out_of_process_handler_proxy, # noqa: F401 _stage_out_of_process_surfaces, - skill_state_path, _validate_child_catalog_namespace, # noqa: F401 _validate_child_route_descriptor, # noqa: F401 _validate_child_settings_descriptor, # noqa: F401 @@ -128,7 +127,7 @@ from ouroboros.extension_surface_names import ( extension_surface_name, # noqa: F401 parse_extension_surface_name, # noqa: F401 ) -from ouroboros.skill_loader import _SKILL_DIR_CACHE_NAMES, _sanitize_skill_name, LoadedSkill, SkillPayloadUnreadable, compute_content_hash, discover_skills, find_skill, grant_status_for_skill, requested_core_setting_keys, skill_conflict_status, skill_review_gate, skill_state_dir # noqa: F401 +from ouroboros.skill_loader import _SKILL_DIR_CACHE_NAMES, _sanitize_skill_name, LoadedSkill, SkillPayloadUnreadable, compute_content_hash, discover_skills, find_skill, grant_status_for_skill, requested_core_setting_keys, skill_conflict_status, skill_review_gate, skill_state_dir, skill_state_dir_path # noqa: F401 from ouroboros.skill_token import SkillToken # noqa: F401 from ouroboros.tools.skill_exec import _scrub_env # noqa: F401 from ouroboros.utils import atomic_write_json, read_json_dict, utc_now_iso # noqa: F401 @@ -475,7 +474,7 @@ def ensure_companions_running( } # Fix-round-6: resolved WITHOUT mkdir — the post-fence attach creates it. - state_dir = skill_state_path(drive_root, skill.name) + state_dir = skill_state_dir_path(drive_root, skill.name) # ABI-9 generation-bound recovery: publish under the lifecycle lock; the # seam re-validates that the publication observed above is still live — # an unload/reload completing since the snapshot is a typed refusal with diff --git a/ouroboros/extension_plugin_api.py b/ouroboros/extension_plugin_api.py index 0cc5612ea..c9f95946c 100644 --- a/ouroboros/extension_plugin_api.py +++ b/ouroboros/extension_plugin_api.py @@ -31,18 +31,19 @@ from ouroboros.contracts.plugin_api import ( LEGACY_PLUGIN_API_GENERATION, RuntimeInfo, VALID_EXTENSION_PERMISSIONS, - VALID_EXTENSION_ROUTE_METHODS, available_capabilities, capability_available, + normalize_extension_route_methods, ) from ouroboros.event_bus import VALID_TOPICS as VALID_EVENT_TOPICS, get_global_event_bus from ouroboros.extension_companion import CompanionDescriptor, get_global_supervisor, is_server_process -from ouroboros.extension_child_catalog import materialize_companion_env from ouroboros.extension_isolated_deps import ( async_isolated_site_dirs_scope, isolated_site_dirs_scope, ) from ouroboros.extension_registry_state import ( + DISPATCH_SURFACE_KINDS, + SURFACE_KINDS, _PluginAPIConfig, _ExtensionRegistrations, _StagedCompanionSpawn, @@ -74,6 +75,11 @@ from ouroboros.extension_ui_validation import ( validate_ui_render as _validate_ui_render, ) from ouroboros.gateway.host_service import AUTH_TOKEN_FILENAME +from ouroboros.node_runtime import ( + prepend_skill_node_emergency_path, + skill_manifest_owns_path, + skill_node_argv, +) from ouroboros.provider_models import MODEL_PROVIDER_CREDENTIAL_KEYS from ouroboros.skill_loader import compute_content_hash, requested_core_setting_keys from ouroboros.skill_token import SkillToken @@ -108,11 +114,9 @@ def _reject_extension_child_side_effect(capability: str) -> None: Every side-effect registration method calls this; the matrix in ``contracts.plugin_api`` is the single source of truth for what an - out-of-process (isolated-dep) child may use. on_unload, send_ws_message, and - register_companion_process are supported out-of-process; subscribe_event and - register_supervised_task are not (use a companion_process instead). + out-of-process (isolated-dep) child may use, and the refusal reads the + available set from there rather than restating it. """ - mode = current_execution_mode() if not capability_available(capability, mode): available = ", ".join(sorted(available_capabilities(mode))) @@ -195,11 +199,7 @@ class PluginAPIImpl: self._env_allow = frozenset(str(k).strip() for k in (config.env_allowlist or [])) self._env_allow_upper = frozenset(k.upper() for k in self._env_allow) self._state_dir = pathlib.Path(config.state_dir) - self._drive_root = ( - pathlib.Path(config.drive_root) - if config.drive_root is not None - else self._state_dir - ) + self._drive_root = pathlib.Path(config.drive_root) if config.drive_root is not None else self._state_dir self._subscribe_events = frozenset(str(t).strip() for t in (config.subscribe_events or []) if str(t).strip()) self._companion_specs = { str(item.get("name") or "").strip(): dict(item) @@ -255,11 +255,7 @@ class PluginAPIImpl: """Whether this live in-process extension can read a funded model key.""" if "read_settings" not in self._permissions: return False - candidates = ( - self._env_allow_upper - & self._granted_upper - & MODEL_PROVIDER_CREDENTIAL_KEYS - ) + candidates = self._env_allow_upper & self._granted_upper & MODEL_PROVIDER_CREDENTIAL_KEYS if not candidates: return False settings = self._settings_reader() or {} @@ -313,8 +309,7 @@ class PluginAPIImpl: if self._skill_dir is None: return handler(*args, **kwargs) with isolated_site_dirs_scope(self._skill_dir, enabled=self._dependency_site_dirs_enabled): - result = handler(*args, **kwargs) - return result + return handler(*args, **kwargs) return _wrapped @@ -374,22 +369,7 @@ class PluginAPIImpl: ) -> None: self._require("route") rel = _assert_namespace_path(path) - methods_iter = (methods,) if isinstance(methods, str) else (methods or ()) - norm_methods = tuple( - dict.fromkeys( - str(m).strip().upper() - for m in methods_iter - if str(m).strip() - ) - ) - if not norm_methods: - raise ExtensionRegistrationError("route methods must be non-empty") - invalid_methods = [m for m in norm_methods if m not in VALID_EXTENSION_ROUTE_METHODS] - if invalid_methods: - raise ExtensionRegistrationError( - f"route methods {invalid_methods!r} are unsupported; " - f"expected subset of {sorted(VALID_EXTENSION_ROUTE_METHODS)}" - ) + norm_methods = normalize_extension_route_methods(methods, subject="route") mount = f"/api/extensions/{self._skill}/{rel}" with _lock: self._stage_surface_locked(_routes, self._staged.routes, mount, { @@ -552,10 +532,9 @@ class PluginAPIImpl: raise ExtensionRegistrationError("companion command must be declared in manifest") if expected_runtime in {"python", "python3"} and cmd[0] in {"python", "python3"}: cmd = [sys.executable, *cmd[1:]] - # T14 node symmetry with the python rewrite; policy in child_catalog. - from ouroboros.extension_child_catalog import companion_node_argv - - cmd = companion_node_argv(spec, expected_runtime, cmd) + # T14 node symmetry with the python rewrite above; the bundled-first + # precedence and its PATH-override carve-out live in node_runtime. + cmd = skill_node_argv(spec, expected_runtime, cmd) if not is_server_process(): with _lock: self._require_open_locked() @@ -588,6 +567,55 @@ class PluginAPIImpl: _StagedCompanionSpawn(name=clean_name, descriptor=descriptor, spec=dict(spec)) ) + def _companion_env(self, spec: Dict[str, Any], token: str) -> Dict[str, str]: + """The env one staged companion is spawned with (fix-round-6). + + Built only inside ``_publish_registrations``'s post-swap attach, after + the generation fence admitted the publication: the settings-derived + values (``_scrub_env`` -> ``load_settings`` takes the settings lock and + may persist a settings migration), the manifest env overlay, the Host + Service bridge URL/token and the isolated-dep PYTHONPATH are all + resolved HERE, so the pre-fence descriptor build stays purely + computational. + """ + from ouroboros.extension_isolated_deps import _isolated_python_site_dirs + from ouroboros.gateway.host_service import DEFAULT_HOST_SERVICE_HOST, host_service_port + from ouroboros.tools.skill_exec import _scrub_env + # Case-aware merge (delta finding D2-8): a manifest "Path" must REPLACE + # the allowlisted "PATH" on Windows, never sit next to it — duplicate + # case-variant env keys make CreateProcess-era spawns fail or pick an + # undefined winner. Same contract as the executor-local service lane. + from ouroboros.workspace_executor import overlay_env + + reserved = set(FORBIDDEN_SKILL_SETTINGS) | {"HOST_SERVICE_TOKEN", "HOST_SERVICE_URL"} + env = overlay_env( + _scrub_env( + list(self._env_allow), self._state_dir, self._skill, + granted_keys=list(self._granted_upper), + ), + { + str(key): str(value) + for key, value in (spec.get("env") or {}).items() + if str(key).upper() not in reserved + }, + ) + env["HOST_SERVICE_URL"] = f"http://{DEFAULT_HOST_SERVICE_HOST}:{host_service_port()}" + env["HOST_SERVICE_TOKEN"] = token + site_dirs = [] if self._skill_dir is None else [ + str(path) for path in _isolated_python_site_dirs(self._skill_dir) + ] + if site_dirs: + inherited = env.get("PYTHONPATH") + env["PYTHONPATH"] = os.pathsep.join([*site_dirs, inherited] if inherited else site_dirs) + if str(spec.get("runtime") or "").strip() in {"node", "npm"} and not skill_manifest_owns_path(spec): + # T14 emergency PATH prepend, the other half of the argv rewrite in + # register_companion_process: descriptor env keys win over the + # supervisor's `_companion_base_env` merge, so the prepend reaches + # the child (and the PATH it would otherwise inherit) and survives + # supervisor restarts. + prepend_skill_node_emergency_path(env, fallback_path=os.environ.get("PATH", "")) + return env + def _stage_companion_name_locked(self, name: str) -> None: if name not in self._staged.companion_names: self._staged.companion_names.append(name) @@ -690,30 +718,21 @@ class PluginAPIImpl: # --- ABI-9 atomic publication (stage -> validate -> swap) --- - def _run_staged_disposers(self, extra: Sequence[Callable[[], Any]] = ()) -> None: - for dispose in [*reversed(list(extra)), *reversed(self._staged.disposers)]: - try: - dispose() - except Exception: - log.warning("extension %s staged-registration disposer failed", self._skill, exc_info=True) - self._staged = _StagedRegistrations() - def _abort_registration(self) -> None: - """Discard the staged snapshot and undo its live side effects.""" + """Discard the staged snapshot. + + Nothing needs undoing: ABI-9 staging is purely computational, so an + abandoned snapshot has no live side effect anywhere — that is the + whole point of deferring the attach to publication. + """ self._close_registration() - self._run_staged_disposers() + self._staged = _StagedRegistrations() def _staged_surface_conflicts_locked(self) -> list[str]: return [ key - for live, staged in ( - (_tools, self._staged.tools), - (_routes, self._staged.routes), - (_ws_handlers, self._staged.ws_handlers), - (_ui_tabs, self._staged.ui_tabs), - (_settings_sections, self._staged.settings_sections), - ) - for key in staged + for kind, live, _label in SURFACE_KINDS + for key in getattr(self._staged, kind) if key in live ] @@ -728,48 +747,47 @@ class PluginAPIImpl: ) -> None: """Atomically publish the staged registration snapshot (ABI-9). - validate -> SWAP -> attach under ONE registry-lock hold: the - definitive unload/conflict validation runs FIRST, so a refused - publication has produced no externally visible effect at all — no - supervised runner, no companion process, no bus subscription. The - validated snapshot then swaps into the process-wide registries as - the authoritative bundle, stamped with a fresh generation digest, - and only AFTER the swap do the deferred side effects attach — the - event-bus subscriptions, the supervised runners, the companion - spawns — still inside the same critical section. A handler is - therefore visible to the bus only for an already-published - extension: a concurrent ``EventBus.publish()`` (which takes only - the bus's own lock) landing between the validation and the attach - cannot invoke a handler of a not-yet-published extension. A bundle - may publish MORE than once (the OOP companion recovery path): a - later publication mints a fresh digest and RE-STAMPS every - already-published descriptor the bundle owns, so per-surface - provenance never diverges from the bundle digest. Such a recovery - publication is GENERATION-BOUND: it passes - ``require_live_generation`` and is admitted only while the bundle it - observed is STILL the live publication — a vanished bundle or a - different generation raises a typed ``ExtensionStaleRecoveryError`` - before any mutation (the whole companion env included — state dir, - auth token and settings-derived values materialize into the staged - descriptors only in the post-swap attach below, never during the - purely computational descriptor build), and recovery never creates a - bundle, so a completed unload/reload cannot be resurrected. Every - attachable effect is recorded on the published bundle at the swap + validate -> SWAP -> attach under ONE registry-lock hold. The definitive + unload/conflict validation runs FIRST, so a refused publication has + produced no externally visible effect at all — no supervised runner, no + companion process, no bus subscription. The validated snapshot then + swaps into the process-wide registries as the authoritative bundle, + stamped with a fresh generation digest, and only AFTER the swap do the + deferred side effects attach — event-bus subscriptions, supervised + runners, companion spawns — still inside the same critical section. A + handler is therefore visible to the bus only for an already-published + extension: a concurrent ``EventBus.publish()`` (which takes only the + bus's own lock) landing between the validation and the attach cannot + invoke a handler of a not-yet-published extension. + + A bundle may publish MORE than once (the OOP companion recovery path): + a later publication mints a fresh digest and RE-STAMPS every + already-published descriptor the bundle owns, so per-surface provenance + never diverges from the bundle digest. Such a recovery publication is + GENERATION-BOUND: it passes ``require_live_generation`` and is admitted + only while the bundle it observed is STILL the live publication — a + vanished bundle or a different generation raises a typed + ``ExtensionStaleRecoveryError`` before any mutation (the whole companion + env included — state dir, auth token and settings-derived values + materialize into the staged descriptors only in the post-swap attach + below, never during the purely computational descriptor build), and + recovery never creates a bundle, so a completed unload/reload cannot be + resurrected. + + Every attachable effect is recorded on the published bundle at the swap (futures as they are created), so a failure while attaching leaves - nothing orphaned: it is disclosed and raised into the caller's - STANDARD dispose+unload path (``unload_extension`` reaps the - bundle's surfaces, subscriptions, futures and companions). No - concurrent unload or conflicting publication can interleave - anywhere inside (both mutate only under this lock). + nothing orphaned: it is disclosed and raised into the caller's STANDARD + dispose+unload path (``unload_extension`` reaps the bundle's surfaces, + subscriptions, futures and companions). No concurrent unload or + conflicting publication can interleave anywhere inside (both mutate only + under this lock). """ with _lock: try: self._require_open_locked() if require_live_generation is not None: live_bundle = _extensions.get(self._skill) - live_generation = ( - str(live_bundle.generation_digest or "") if live_bundle is not None else "" - ) + live_generation = str(getattr(live_bundle, "generation_digest", "") or "") if live_bundle is None or live_generation != str(require_live_generation): raise ExtensionStaleRecoveryError( f"skill {self._skill!r} recovery publication refused: " @@ -784,7 +802,7 @@ class PluginAPIImpl: "publication refused" ) except Exception: - self._run_staged_disposers() + self._staged = _StagedRegistrations() raise staged = self._staged bundle = _extensions.get(self._skill) @@ -801,13 +819,9 @@ class PluginAPIImpl: digest = uuid.uuid4().hex bundle.generation_digest = digest self._published_generation = digest - for live, staged_map, bundle_keys, stamp in ( - (_tools, staged.tools, bundle.tools, True), - (_routes, staged.routes, bundle.routes, True), - (_ws_handlers, staged.ws_handlers, bundle.ws_handlers, True), - (_ui_tabs, staged.ui_tabs, bundle.ui_tabs, False), - (_settings_sections, staged.settings_sections, bundle.settings_sections, False), - ): + for kind, live, _label in SURFACE_KINDS: + staged_map, bundle_keys = getattr(staged, kind), getattr(bundle, kind) + stamp = kind in DISPATCH_SURFACE_KINDS if stamp: # Staged-protocol restamp (ABI-9): a bundle may publish # MORE than once (the server-side companion recovery path @@ -855,7 +869,8 @@ class PluginAPIImpl: self._state_dir.mkdir(parents=True, exist_ok=True) token = mint_skill_token(self._state_dir, self._skill, self._skill_dir) for spawn in staged.companion_spawns: - materialize_companion_env(self, spawn.descriptor, spawn.spec, token) + spawn.descriptor.env.clear() + spawn.descriptor.env.update(self._companion_env(spawn.spec, token)) bus = get_global_event_bus() for sub in staged.event_subscriptions: bus.subscribe(self._skill, sub.topic, sub.handler, sub_id=sub.sub_id) @@ -883,22 +898,15 @@ class PluginAPIImpl: with _lock: self._registration_closed = True self._runtime_closing = True - with self._api_lock: - with _lock: - self._runtime_closed = True + with self._api_lock, _lock: + self._runtime_closed = True # --- runtime access --- def log(self, level: str, message: str, **fields: Any) -> None: - lvl = str(level or "info").lower() levels = {"debug": 10, "info": 20, "warning": 30, "error": 40} - log.log( - levels.get(lvl, 20), - "[ext %s] %s %s", - self._skill, - message, - fields if fields else "", - ) + level_no = levels.get(str(level or "info").lower(), logging.INFO) + log.log(level_no, "[ext %s] %s %s", self._skill, message, fields or "") def get_settings(self, keys: Sequence[str]) -> Dict[str, Any]: with self._api_lock: @@ -953,40 +961,32 @@ class PluginAPIImpl: def get_runtime_info(self) -> RuntimeInfo: """Return the PluginAPI runtime-info snapshot without manifest I/O.""" try: - from ouroboros.config import ( - get_runtime_mode as _get_runtime_mode, - DATA_DIR as _DATA_DIR, - ) - runtime_mode = _get_runtime_mode() - data_dir = str(_DATA_DIR) + from ouroboros.config import DATA_DIR, get_runtime_mode + runtime_mode, data_dir = get_runtime_mode(), str(DATA_DIR) except Exception: - runtime_mode = "advanced" - data_dir = "" + runtime_mode, data_dir = "advanced", "" try: - from ouroboros import get_version as _get_version - app_version = str(_get_version()) + from ouroboros import get_version + app_version = str(get_version()) except Exception: app_version = "" try: - from ouroboros.config import AGENT_SERVER_PORT as _agent_port, PORT_FILE as _PORT_FILE - server_port = 0 + from ouroboros.config import AGENT_SERVER_PORT, PORT_FILE try: - port_text = pathlib.Path(_PORT_FILE).read_text(encoding="utf-8").strip() - if port_text: - server_port = int(port_text) + # The live port file wins over the configured default: a server + # that started on a fallback port must still be reachable here. + live_port = int(pathlib.Path(PORT_FILE).read_text(encoding="utf-8").strip()) except Exception: - server_port = 0 - if server_port <= 0: - server_port = int(_agent_port) + live_port = 0 + server_port = live_port if live_port > 0 else int(AGENT_SERVER_PORT) except Exception: server_port = 0 - skill_dir = str(getattr(self, "_skill_dir", "") or "") mode = current_execution_mode() return { "runtime_mode": runtime_mode, "app_version": app_version, "data_dir": data_dir, - "skill_dir": skill_dir, + "skill_dir": str(self._skill_dir or ""), "state_dir": str(self._state_dir), "server_port": server_port, # Capability negotiation: an extension can branch on its execution mode diff --git a/ouroboros/extension_registry_state.py b/ouroboros/extension_registry_state.py index fb3e09843..e4958201d 100644 --- a/ouroboros/extension_registry_state.py +++ b/ouroboros/extension_registry_state.py @@ -100,9 +100,9 @@ class _StagedRegistrations: subscriptions — attaches only at publication, after the definitive unload/conflict validation AND after the snapshot swap (so a handler is visible to the bus only for an already-published extension); an aborted - registration leaves zero residue, and a post-swap attach failure is - disposed through the standard unload path. The disposers list is - loader-internal and never exposed through the PluginAPI ABI. + registration leaves zero residue — staging is purely computational, so + there is nothing to dispose — and a post-swap attach failure is disposed + through the standard unload path. """ tools: Dict[str, Any] = field(default_factory=dict) @@ -115,7 +115,6 @@ class _StagedRegistrations: companion_names: List[str] = field(default_factory=list) supervised_tasks: List[_StagedSupervisedTask] = field(default_factory=list) companion_spawns: List[_StagedCompanionSpawn] = field(default_factory=list) - disposers: List[Callable[[], Any]] = field(default_factory=list) @dataclass @@ -158,6 +157,24 @@ _ui_tabs: Dict[str, Any] = {} # {":": tab_spec} # Declarative settings sections keyed like UI tabs. _settings_sections: Dict[str, Any] = {} +# The five surface kinds, each naming the field that carries it on a staged +# snapshot AND on a published bundle, its live registry, and the word used in +# refusals. Every walk over "all surfaces" — conflict detection, the atomic +# swap, the child-catalog staging — reads this, so a sixth kind cannot be +# added to some walks and forgotten by others. +SURFACE_KINDS: Sequence[tuple[str, Dict[str, Any], str]] = ( + ("tools", _tools, "tool"), + ("routes", _routes, "route"), + ("ws_handlers", _ws_handlers, "ws handler"), + ("ui_tabs", _ui_tabs, "ui tab"), + ("settings_sections", _settings_sections, "settings section"), +) + +# Dispatch surfaces: a physical call arrives on these, so each descriptor is +# stamped with the publication that owns it. UI kinds are read as a snapshot +# and carry no per-surface provenance. +DISPATCH_SURFACE_KINDS = frozenset({"tools", "routes", "ws_handlers"}) + def _lifecycle_lock_for(skill_name: str) -> threading.RLock: with _lock: diff --git a/ouroboros/marketplace/isolated_deps.py b/ouroboros/marketplace/isolated_deps.py index 7c5d02899..d5713dd9e 100644 --- a/ouroboros/marketplace/isolated_deps.py +++ b/ouroboros/marketplace/isolated_deps.py @@ -85,20 +85,14 @@ def _installer_env(env_root: pathlib.Path, *, ecosystem: str = "") -> Dict[str, "CARGO_HOME": str(env_root / "cargo" / "home"), "CARGO_TARGET_DIR": str(env_root / "cargo" / "target"), }) if ecosystem == "node": - # Emergency-only: when the PATH node is missing/execution-probed broken - # and the healthy bundled node was selected (skill-family precedence in - # platform_layer), npm's `#!/usr/bin/env node` shebang must resolve the - # working runtime, so the curated PATH gains the bundled-node dir. On - # healthy systems the env stays byte-identical. npm itself is not - # bundled: an absent npm still fails honestly upstream, and an npm - # launcher with an ABSOLUTE node shebang ignoring PATH is a disclosed - # residual. - from ouroboros.platform_layer import skill_node_emergency_path_dir + # npm's `#!/usr/bin/env node` shebang must resolve a working runtime. + # node_runtime owns the emergency verdict, the byte-identical healthy + # case, and the disclosed residuals (npm itself is not bundled: an + # absent npm still fails honestly upstream, and an npm launcher with an + # ABSOLUTE node shebang ignores PATH). + from ouroboros.platform_layer import prepend_skill_node_emergency_path - prepend_dir = skill_node_emergency_path_dir() - if prepend_dir: - current = env.get("PATH", "") - env["PATH"] = os.pathsep.join([prepend_dir, current]) if current else prepend_dir + prepend_skill_node_emergency_path(env) return env diff --git a/ouroboros/node_runtime.py b/ouroboros/node_runtime.py index e3f397eb0..4b7916f7f 100644 --- a/ouroboros/node_runtime.py +++ b/ouroboros/node_runtime.py @@ -17,7 +17,7 @@ import os import pathlib import shutil import subprocess -from typing import Dict, List, NamedTuple, Tuple +from typing import Any, Dict, List, NamedTuple, Tuple from ouroboros import platform_layer as _platform @@ -224,3 +224,48 @@ def skill_node_emergency_path_dir(timeout_sec: float = 10) -> str: if _path_node_runtime_health(timeout_sec=timeout_sec).healthy: return "" return str(pathlib.Path(selected).parent) + + +def prepend_skill_node_emergency_path(env: Dict[str, str], *, fallback_path: str = "") -> None: + """Front-load a skill-family child's PATH with the emergency node dir. + + The APPLYING half of ``skill_node_emergency_path_dir``: on a healthy PATH + there is no emergency, ``env`` is left untouched and the child environment + stays byte-identical. Shared by the isolated-dep installer env and the + extension companion spawn env so the two cannot drift on which node an + `#!/usr/bin/env node` shebang resolves. ``fallback_path`` is the PATH the + child would otherwise inherit, for an ``env`` carrying none of its own. + """ + prepend_dir = skill_node_emergency_path_dir() + if not prepend_dir: + return + existing = env.get("PATH") or fallback_path + env["PATH"] = os.pathsep.join([prepend_dir, existing]) if existing else prepend_dir + + +def skill_manifest_owns_path(spec: Dict[str, Any]) -> bool: + """Whether a skill's manifest declares its own PATH for a child process. + + An explicit manifest PATH means the author owns the runtime lookup, so + neither the bundled argv rewrite below nor the emergency prepend above may + shadow it (T14). + """ + return any(str(key).upper() == "PATH" for key in (spec.get("env") or {})) + + +def skill_node_argv(spec: Dict[str, Any], declared_runtime: str, argv: List[str]) -> List[str]: + """A node-family child's argv, rewritten onto the selected node runtime. + + T14 symmetry with the python -> ``sys.executable`` rewrite: python skills + and companions already run on the embedded interpreter, so a node one runs + on the runtime ``select_skill_node_runtime`` picked (bundled-first, health + rollback included). ``npm`` is never rewritten — its launcher resolves node + through a ``#!/usr/bin/env node`` shebang, which the emergency PATH prepend + covers instead. On a healthy PATH argv stays byte-identical. + """ + if declared_runtime not in {"node", "npm"} or skill_manifest_owns_path(spec): + return argv + selected, _provenance = select_skill_node_runtime() + if selected and argv and argv[0] == "node": + return [selected, *argv[1:]] + return argv diff --git a/ouroboros/platform_layer.py b/ouroboros/platform_layer.py index 07e5946f1..71dfd7d32 100644 --- a/ouroboros/platform_layer.py +++ b/ouroboros/platform_layer.py @@ -1459,7 +1459,10 @@ _NODE_RUNTIME_REEXPORTS = ( "NodeRuntimeHealth", "node_runtime_health", "probe_node_version", + "prepend_skill_node_emergency_path", "select_skill_node_runtime", + "skill_manifest_owns_path", + "skill_node_argv", "skill_node_emergency_path_dir", ) diff --git a/ouroboros/skill_loader.py b/ouroboros/skill_loader.py index c7f95eae6..92d226dde 100644 --- a/ouroboros/skill_loader.py +++ b/ouroboros/skill_loader.py @@ -182,6 +182,11 @@ def skill_state_dir_path(drive_root: pathlib.Path, name: str) -> pathlib.Path: consumers (``load_review_state``, the RC auditor's admission reuse): reading state must never mutate the root it reads from. + The generation-fenced companion recovery resolves the state dir through + here BEFORE its fence (fix-round-6); creation belongs to the post-fence + attach in ``_publish_registrations``, so a stale refusal creates no + directories. ``skill_state_dir`` is the creating variant for load paths. + The name is normalized to its alnum-dashes shape before joining so a malicious manifest ``name: ../foo`` cannot escape the state root. """ diff --git a/tests/test_extension_companion.py b/tests/test_extension_companion.py index f5547ee2e..6b6dd99c7 100644 --- a/tests/test_extension_companion.py +++ b/tests/test_extension_companion.py @@ -388,7 +388,7 @@ def test_unload_completing_between_snapshot_and_publication_refuses_recovery( ) token_bytes = token_path.read_bytes() - real_state_path = extension_loader.skill_state_path + real_state_path = extension_loader.skill_state_dir_path def _unload_wins_the_race(drive_root_arg, skill_name_arg): # Deterministic interleave: the concurrent unload COMPLETES after @@ -396,7 +396,7 @@ def test_unload_completing_between_snapshot_and_publication_refuses_recovery( extension_loader.unload_extension(skill_name_arg) return real_state_path(drive_root_arg, skill_name_arg) - monkeypatch.setattr(extension_loader, "skill_state_path", _unload_wins_the_race) + monkeypatch.setattr(extension_loader, "skill_state_dir_path", _unload_wins_the_race) result = extension_loader.ensure_companions_running( loaded.name, drive_root, lambda: {}, repo_path=str(repo_root), ) @@ -444,21 +444,21 @@ def test_recovery_publication_refuses_on_generation_mismatch_without_effects( ) token_bytes = token_path.read_bytes() - real_state_path = extension_loader.skill_state_path + real_state_path = extension_loader.skill_state_dir_path def _reload_wins_the_race(drive_root_arg, skill_name_arg): # Deterministic interleave: the unload/reload COMPLETES after # recovery snapshotted the generation and before it publishes; # the bundle EXISTS again, under a fresh generation. Restore the # real resolver first — the reload itself resolves state dirs. - monkeypatch.setattr(extension_loader, "skill_state_path", real_state_path) + monkeypatch.setattr(extension_loader, "skill_state_dir_path", real_state_path) extension_loader.unload_extension(skill_name_arg) assert extension_loader.load_extension( loaded, lambda: {}, drive_root=drive_root_arg, skills=[loaded], ) is None return real_state_path(drive_root_arg, skill_name_arg) - monkeypatch.setattr(extension_loader, "skill_state_path", _reload_wins_the_race) + monkeypatch.setattr(extension_loader, "skill_state_dir_path", _reload_wins_the_race) result = extension_loader.ensure_companions_running( loaded.name, drive_root, lambda: {}, repo_path=str(repo_root), ) @@ -509,7 +509,7 @@ def test_stale_recovery_does_not_break_live_publication_authorization( (repo_root_v2 / loaded_v1.name / "scripts" / "daemon.py").write_text( "print('v2')\n", encoding="utf-8" ) - real_state_path = extension_loader.skill_state_path + real_state_path = extension_loader.skill_state_dir_path g2_token_bytes: dict = {} def _reload_v2_wins_the_race(drive_root_arg, skill_name_arg): @@ -517,7 +517,7 @@ def test_stale_recovery_does_not_break_live_publication_authorization( # recovery snapshotted generation/payload and before it publishes; # the G2 token is bound to the NEW root's content hash while the # stale recovery still holds the v1 snapshot at the old root. - monkeypatch.setattr(extension_loader, "skill_state_path", real_state_path) + monkeypatch.setattr(extension_loader, "skill_state_dir_path", real_state_path) extension_loader.unload_extension(skill_name_arg) loaded_v2 = find_skill( drive_root, skill_name_arg, repo_path=str(repo_root_v2), @@ -537,7 +537,7 @@ def test_stale_recovery_does_not_break_live_publication_authorization( g2_token_bytes["value"] = token_path.read_bytes() return real_state_path(drive_root_arg, skill_name_arg) - monkeypatch.setattr(extension_loader, "skill_state_path", _reload_v2_wins_the_race) + monkeypatch.setattr(extension_loader, "skill_state_dir_path", _reload_v2_wins_the_race) result = extension_loader.ensure_companions_running( loaded_v1.name, drive_root, lambda: {}, repo_path=str(repo_root), selected_skill=loaded_v1, # the stale v1 snapshot the recovery holds @@ -651,7 +651,7 @@ def test_stale_recovery_with_env_from_settings_has_zero_filesystem_effects( assert env["EXT_OVERLAY"] == "from-manifest" assert str(site_dir.resolve()) in (env.get("PYTHONPATH") or "") - real_state_path = extension_loader.skill_state_path + real_state_path = extension_loader.skill_state_dir_path snapshot: dict = {} def _unload_wins_and_snapshots(drive_root_arg, skill_name_arg): @@ -664,7 +664,7 @@ def test_stale_recovery_with_env_from_settings_has_zero_filesystem_effects( assert not lock_path.exists() return real_state_path(drive_root_arg, skill_name_arg) - monkeypatch.setattr(extension_loader, "skill_state_path", _unload_wins_and_snapshots) + monkeypatch.setattr(extension_loader, "skill_state_dir_path", _unload_wins_and_snapshots) result = extension_loader.ensure_companions_running( loaded.name, drive_root, lambda: {}, repo_path=str(repo_root), ) @@ -699,7 +699,7 @@ def test_transient_hash_error_never_rotates_a_valid_token( loaded, lambda: {}, drive_root=drive_root, skills=[loaded], ) is None env_token = fake.started[0].env["HOST_SERVICE_TOKEN"] - state_dir = extension_loader.skill_state_path(drive_root, loaded.name) + state_dir = extension_loader.skill_state_dir_path(drive_root, loaded.name) token_path = state_dir / host_service.AUTH_TOKEN_FILENAME before = token_path.read_bytes() @@ -813,11 +813,10 @@ def _node_companion_api(tmp_path: pathlib.Path, monkeypatch, *, runtime: str, co def _staged_companion_env(api: PluginAPIImpl) -> tuple: """The staged descriptor plus its env AFTER the post-fence materialization - (the campaign structure defers env to `materialize_companion_env`).""" - from ouroboros.extension_child_catalog import materialize_companion_env + (the campaign structure defers env to `PluginAPIImpl._companion_env`).""" spawn = api._staged.companion_spawns[0] - materialize_companion_env(api, spawn.descriptor, spawn.spec, "test-token") + spawn.descriptor.env.update(api._companion_env(spawn.spec, "test-token")) return spawn.descriptor, spawn.descriptor.env @@ -827,13 +826,13 @@ def test_companion_node_command_rewritten_via_policy_helper(tmp_path: pathlib.Pa emergency, so the child env stays byte-identical.""" import os - from ouroboros import platform_layer + from ouroboros import node_runtime selected = str(tmp_path / "bundle" / "bin" / "node") monkeypatch.setattr( - platform_layer, "select_skill_node_runtime", lambda timeout_sec=10: (selected, "bundled") + node_runtime, "select_skill_node_runtime", lambda timeout_sec=10: (selected, "bundled") ) - monkeypatch.setattr(platform_layer, "skill_node_emergency_path_dir", lambda timeout_sec=10: "") + monkeypatch.setattr(node_runtime, "skill_node_emergency_path_dir", lambda timeout_sec=10: "") api = _node_companion_api(tmp_path, monkeypatch, runtime="node", command=["node", "server.js"]) api.register_companion_process("daemon") @@ -849,16 +848,16 @@ def test_companion_npm_not_rewritten_but_emergency_path_prepended(tmp_path: path `#!/usr/bin/env node` shebang finds the working runtime.""" import os - from ouroboros import platform_layer + from ouroboros import node_runtime bundle_bin = str(tmp_path / "bundle" / "bin") monkeypatch.setattr( - platform_layer, + node_runtime, "select_skill_node_runtime", lambda timeout_sec=10: (str(pathlib.Path(bundle_bin) / "node"), "bundled"), ) monkeypatch.setattr( - platform_layer, "skill_node_emergency_path_dir", lambda timeout_sec=10: bundle_bin + node_runtime, "skill_node_emergency_path_dir", lambda timeout_sec=10: bundle_bin ) api = _node_companion_api(tmp_path, monkeypatch, runtime="npm", command=["npm", "run", "start"]) @@ -874,14 +873,14 @@ def test_companion_node_unusable_leaves_command_and_env_untouched(tmp_path: path """Nothing usable -> honest failure at spawn time, no rewrite, no env edit.""" import os - from ouroboros import platform_layer + from ouroboros import node_runtime monkeypatch.setattr( - platform_layer, + node_runtime, "select_skill_node_runtime", lambda timeout_sec=10: ("", "bundled:absent; path:broken:signal:SIGKILL"), ) - monkeypatch.setattr(platform_layer, "skill_node_emergency_path_dir", lambda timeout_sec=10: "") + monkeypatch.setattr(node_runtime, "skill_node_emergency_path_dir", lambda timeout_sec=10: "") api = _node_companion_api(tmp_path, monkeypatch, runtime="node", command=["node", "server.js"]) api.register_companion_process("daemon") diff --git a/tests/test_isolated_deps.py b/tests/test_isolated_deps.py index 408857c2e..069b2fe03 100644 --- a/tests/test_isolated_deps.py +++ b/tests/test_isolated_deps.py @@ -37,12 +37,12 @@ def test_installer_env_node_emergency_prepends_bundled_node_dir(monkeypatch, tmp """Emergency only (PATH node dead + healthy bundled selected): the node installer env PATH gains the bundled-node dir so npm's `#!/usr/bin/env node` shebang resolves the working runtime. Python installs never get it.""" - from ouroboros import platform_layer + from ouroboros import node_runtime monkeypatch.setenv("PATH", "/usr/bin") bundle_bin = str(tmp_path / "bundle" / "bin") monkeypatch.setattr( - platform_layer, "skill_node_emergency_path_dir", lambda timeout_sec=10: bundle_bin + node_runtime, "skill_node_emergency_path_dir", lambda timeout_sec=10: bundle_bin ) node_env = _installer_env(tmp_path / ".ouroboros_env", ecosystem="node") @@ -53,11 +53,11 @@ def test_installer_env_node_emergency_prepends_bundled_node_dir(monkeypatch, tmp def test_installer_env_node_healthy_system_is_byte_identical(monkeypatch, tmp_path): - from ouroboros import platform_layer + from ouroboros import node_runtime monkeypatch.setenv("PATH", "/usr/bin") monkeypatch.setattr( - platform_layer, "skill_node_emergency_path_dir", lambda timeout_sec=10: "" + node_runtime, "skill_node_emergency_path_dir", lambda timeout_sec=10: "" ) node_env = _installer_env(tmp_path / ".ouroboros_env", ecosystem="node")