extensions: companion env/argv policy goes home; one skill-state path resolver (CH-1, CH-4, owner 16A)

extension_child_catalog.py had become a cap-driven bucket: 100 of its lines
(skill_state_path, companion_manifest_path_override, companion_node_argv,
materialize_companion_env) were companion env/argv/state-dir policy that
belongs to PluginAPIImpl and the node-runtime owner, exiled only because
extension_plugin_api.py sat at its 1000-line band. The band is paid down by
simplification, not by another module: the route-method vocabulary is one
function (contracts.normalize_extension_route_methods, shared by
register_route and the child-catalog re-check so the two cannot drift), a
producer-less `disposers` list and its empty loop are gone, three helpers are
shorter. materialize_companion_env becomes PluginAPIImpl._companion_env, the
argv/PATH policy lives in node_runtime (its duplicate emergency-PATH logic
from isolated_deps merged), and the surface kinds are enumerated once
(SURFACE_KINDS). skill_state_path was a byte-for-byte re-implementation of
skill_loader.skill_state_dir_path through a private back-import; the SSOT
stays in skill_loader and the eleven monkeypatch sites point at it.
child_catalog 271 -> 154, plugin_api 999 -> 999 (band held), loader 991.
This commit is contained in:
Ouroboros 2026-09-02 00:53:43 +00:00
parent 8b596344d4
commit 464e1bb99d
11 changed files with 284 additions and 314 deletions

View file

@ -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",

View file

@ -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)

View file

@ -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

View file

@ -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

View file

@ -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] = {} # {"<skill>:<tab_id>": 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:

View file

@ -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

View file

@ -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

View file

@ -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",
)

View file

@ -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.
"""

View file

@ -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")

View file

@ -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")