diff --git a/docs/integrations/deepseek-harness-connector.md b/docs/integrations/deepseek-harness-connector.md index 1b73bf7a6b..18bd9e26b8 100644 --- a/docs/integrations/deepseek-harness-connector.md +++ b/docs/integrations/deepseek-harness-connector.md @@ -131,6 +131,46 @@ continuity or an outer wake/timer. See the adapter README for the home and classification precedence, plus the hermetic verification smoke (`examples/loopx-turn-dsh-builtin-host-e2e-smoke.py`). +## Host Selection And Managed Executor Readback + +The Turn host is **selected, never inferred**. `dsh` is the shipped default +because it is the managed execution unit the steward drives; `LOOPX_TURN_HOST` +re-points that default, and an explicit `--host` (or `--host-adapter-command-json`) +wins over both. A configured `DEEPSEEK_API_KEY` only *authenticates* the selected +host: discovering a credential never changes where a Turn runs. + +Both `loopx turn plan` and `loopx turn run-once` report a `managed_executor` +block, so a caller reads the planned executor instead of inferring it from a +host id: + +```json +{ + "schema_version": "managed_executor_binding_v0", + "executor": "dsh", + "executor_kind": "managed", + "credential_env": "DEEPSEEK_API_KEY", + "endpoint_env": "DEEPSEEK_BASE_URL", + "operator_credential_bound": true, + "available": true, + "unavailable_reason": null +} +``` + +`executor_kind` names where the Turn's model work is billed and bounded: +`managed` for a host bound to an operator credential, `individual` for a host +that runs on one person's own CLI login, and `generic` for a caller-supplied +adapter command. `operator_credential_bound` is the narrower claim: it is `true` +only when the operator credential or an explicit injected runner hook is +configured. `available` is `false` only when LoopX can prove the planned host +cannot launch here, and `null` for executors this projection does not probe +rather than an unproven claim. Only the credential variable *name* is reported; +the value is never read back. + +`run-once --execute` fails closed on that verdict: status `unavailable`, no host +invocation, no journal write, and no quota spend, with +`dsh_runtime_unavailable` or `operator_credential_unconfigured` naming the +missing fact. `plan` reports the same verdict without refusing. + ## Boundaries - LoopX keeps the durable goal, todo, claim, gate, quota, evidence, and diff --git a/docs/reference/protocols/loopx-turn-v0.md b/docs/reference/protocols/loopx-turn-v0.md index 1c14861471..62062c7f41 100644 --- a/docs/reference/protocols/loopx-turn-v0.md +++ b/docs/reference/protocols/loopx-turn-v0.md @@ -97,6 +97,45 @@ terminal failures reach the Turn Journal. The module/subprocess invocation with `--host generic-cli` remains the compatibility and rollback path. See [DeepSeek Harness connector](../../integrations/deepseek-harness-connector.md). +### Host Selection + +The Turn host is **selected, never inferred**. `loopx turn plan` and +`loopx turn run-once` default to the managed `dsh` host, the operator may +re-point that default with `LOOPX_TURN_HOST` or one explicit `--host`, and a +configured operator credential only *authenticates* the host that was already +selected. Discovering `DEEPSEEK_API_KEY` must never re-point a Turn by itself. + +| surface | value | +| --- | --- | +| shipped default host | `dsh` (managed executor) | +| explicit default selector | `LOOPX_TURN_HOST` | +| per-command override | `--host codex-cli\|claude-code\|dsh\|generic-cli` (plan), `codex-cli\|dsh\|generic-cli` (run-once) | +| authenticating credential | `DEEPSEEK_API_KEY`, optional endpoint `DEEPSEEK_BASE_URL` | + +This is a default behavior change for the affected lanes: `run-once` moved from +`generic-cli` to `dsh`, and `plan` from `codex-cli` to `dsh`. `--host +generic-cli` and `--host codex-cli` remain the explicit compatibility and +rollback paths, and a machine that wants the former default should set +`LOOPX_TURN_HOST=generic-cli` (or `codex-cli`) once instead of relying on the +ambient environment. + +`plan` and `run-once` payloads carry the executor readback `managed_executor` +(`managed_executor_binding_v0`): the executor and its kind (`managed`, +`individual`, `generic`), the credential env var *name* (never its value), the +endpoint env var name, whether the executor is operator-credential-bound, and +whether it can launch here. When it cannot, `available` is `false` and +`unavailable_reason` names the missing fact: + +| `unavailable_reason` | meaning | remediation | +| --- | --- | --- | +| `dsh_runtime_unavailable` | the DeepSeek Harness runtime is not importable and no explicit runner hook was supplied | install the released runtime, pass its runner hook, or select `--host codex-cli` | +| `operator_credential_unconfigured` | the managed host is selected but no operator credential or runner hook would authenticate it | set `DEEPSEEK_API_KEY`, or select `--host codex-cli` explicitly | + +`run-once --execute` fails closed on that verdict: status `unavailable`, no host +invocation, no Journal write, and no quota slot spend. An explicitly selected +individual host (`--host codex-cli`, `--host claude-code`) is billed to that +individual CLI login and makes no launchability claim (`available: null`). + ### Five Questions For Any Agent CLI Before wiring Trae CLI, Codex CLI, or another host, answer these five questions: diff --git a/examples/control_plane/cli-output-probe-runner.py b/examples/control_plane/cli-output-probe-runner.py index 1ad861fb4f..878187ccf1 100644 --- a/examples/control_plane/cli-output-probe-runner.py +++ b/examples/control_plane/cli-output-probe-runner.py @@ -108,6 +108,9 @@ def _receipt_row( "reward_memory_outcome_prompt_revision": ( semantics.reward_memory_outcome_prompt_revision(text) ), + "managed_executor_binding_revision": ( + semantics.managed_executor_binding_revision(text) + ), "guided_todo_delta_schema_versions": ( semantics.guided_todo_delta_schema_versions(payload) if isinstance(payload, dict) diff --git a/examples/loopx-turn-managed-default-flow-smoke.py b/examples/loopx-turn-managed-default-flow-smoke.py new file mode 100644 index 0000000000..0bb978854c --- /dev/null +++ b/examples/loopx-turn-managed-default-flow-smoke.py @@ -0,0 +1,642 @@ +#!/usr/bin/env python3 +"""Qualify the explicit operator default flow for one bounded managed Turn. + +The shipped operator rule is *selection first*: the default host is the managed +``dsh`` executor by product decision, the operator credential only authenticates +that selection, and discovering a credential never re-points a Turn. That rule +is only usable if the *default* command (no explicit ``--host``) actually starts +the managed Turn and reports what ran. + +This smoke is hermetic: a local mock OpenAI-compatible SSE server stands in for +the model endpoint, so no operator key and no individual CLI subscription is +consumed. It proves, through the public CLI only: + +1. no credential: the default host is still ``dsh``, reported as an unauthenticated + managed executor with the typed ``operator_credential_unconfigured`` reason; +2. credential: the same default host reports its credential environment, its + billing boundary, and its launchability before any work runs; +3. credential: ``turn run-once`` without ``--host`` starts the real dsh runtime, + commits one validated Turn, and reports the mode/executor/status readback; +4. credential but an unavailable managed runtime: the same default flow fails + closed with a typed reason and writes nothing; +5. an explicit ``--host codex-cli`` stays selected even while a credential is + configured, and an explicit ``--host dsh`` without one is unbound. +""" + +from __future__ import annotations + +import contextlib +import importlib.abc +import importlib.util +import io +import json +import os +import sys +import tempfile +import threading +from collections.abc import Iterator +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Any + + +REPO_ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(REPO_ROOT)) + +from loopx.cli import main as cli_main # noqa: E402 +from loopx.control_plane.turn_driver.host_binding import ( # noqa: E402 + DSH_RUNTIME_UNAVAILABLE, + EXECUTOR_KIND_INDIVIDUAL, + EXECUTOR_KIND_MANAGED, + OPERATOR_CREDENTIAL_UNCONFIGURED, +) + +GOAL_ID = "loopx-turn-managed-default-flow" +AGENT_ID = "codex-managed-default-flow" +TODO_ID = "todo_manageddefault01" +CREDENTIAL_ENV = "DEEPSEEK_API_KEY" +ENDPOINT_ENV = "DEEPSEEK_BASE_URL" +MARKER_NAME = "docs/managed-default-flow-marker.txt" +MARKER_VALUE = "loopx-turn-managed-default-flow-step-1" + + +def _write_fixture(root: Path) -> tuple[Path, Path, Path, Path]: + project = root / "project" + runtime = root / "runtime" + workspace = root / "workspace" + runtime.mkdir(parents=True) + (workspace / "docs").mkdir(parents=True) + + state = project / ".codex" / "goals" / GOAL_ID / "ACTIVE_GOAL_STATE.md" + state.parent.mkdir(parents=True) + state.write_text( + "\n".join( + [ + "---", + "status: active", + "updated_at: 2026-01-01T00:00:00+00:00", + "---", + "", + "# LoopX Managed Default Flow Fixture", + "", + "## Next Action", + "", + "Run one bounded managed Turn on the credential-resolved executor.", + "", + "## Agent Todo", + "", + "- [ ] [P1] Write the managed default flow marker.", + " ", + "", + "## User Todo", + "", + ] + ) + + "\n", + encoding="utf-8", + ) + registry = project / ".loopx" / "registry.json" + registry.parent.mkdir(parents=True) + registry.write_text( + json.dumps( + { + "schema_version": 1, + "common_runtime_root": str(runtime), + "goals": [ + { + "id": GOAL_ID, + "domain": "loopx-turn-public-fixture", + "status": "active", + "repo": str(project), + "state_file": str(state.relative_to(project)), + "adapter": { + "kind": "fixture_v0", + "status": "connected-delivery", + }, + "quota": {"compute": 1.0, "window_hours": 24}, + "coordination": { + "agent_model": "peer_v1", + "registered_agents": [AGENT_ID], + "agent_profiles": { + AGENT_ID: { + "schema_version": "agent_profile_v1", + "profile_role": "fixture", + "scope": "public qualification", + }, + }, + "write_scope": ["docs/**"], + }, + }, + ], + }, + indent=2, + sort_keys=True, + ) + + "\n", + encoding="utf-8", + ) + return project, runtime, workspace, registry + + +def _write_minimal_cordis(session_root: Path) -> Path: + """A dsh JSON-RPC composition that avoids node-pty/subprocess.""" + + path = session_root.parent / "no-pty.cordis.yml" + path.write_text( + "\n".join( + [ + "- id: sdk-jsonrpc-server", + " name: '@deepseek-ai/dsh-sdk-jsonrpc-server'", + "- id: agent-core", + " name: '@deepseek-ai/dsh-agent-spine-demo'", + " config:", + " workspaceContext:", + " maxBytes: 65536", + "- id: llm-deepseek", + " name: '@deepseek-ai/dsh-llm-deepseek'", + "- id: sessions", + " name: '@deepseek-ai/dsh-session-persistence-jsonl'", + " config:", + f" root: {session_root}", + "", + ] + ), + encoding="utf-8", + ) + return path + + +def _validator_command() -> list[str]: + program = ( + "import pathlib,sys,json; " + "json.load(sys.stdin); " + f"p=pathlib.Path({MARKER_NAME!r}); " + "raise SystemExit(0 if p.is_file() and " + f"p.read_text(encoding='utf-8').strip() == {MARKER_VALUE!r} else 9)" + ) + return [sys.executable, "-c", program] + + +def _run_cli(argv: list[str]) -> tuple[int, dict[str, Any]]: + output = io.StringIO() + with contextlib.redirect_stdout(output): + exit_code = cli_main(argv) + payload = json.loads(output.getvalue()) + assert isinstance(payload, dict), payload + return exit_code, payload + + +def _plan_argv( + registry: Path, + runtime: Path, + project: Path, + *, + host: str | None = None, + execution_mode: str | None = None, +) -> list[str]: + argv = [ + "--registry", + str(registry), + "--runtime-root", + str(runtime), + "--format", + "json", + "turn", + "plan", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--scan-root", + str(project), + ] + if host: + argv.extend(["--host", host]) + if execution_mode: + argv.extend(["--execution-mode", execution_mode]) + return argv + + +def _run_once_argv( + *, + registry: Path, + runtime: Path, + project: Path, + workspace: Path, + dsh_home: Path, + cordis: Path, + turn_instance_id: str, + runner_binding: bool = True, +) -> list[str]: + # No --host: the default must come from the explicit product selection, and + # the operator credential only authenticates it. That is the surface this + # smoke qualifies. + argv = [ + "--registry", + str(registry), + "--runtime-root", + str(runtime), + "--format", + "json", + "turn", + "run-once", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--turn-instance-id", + turn_instance_id, + "--project", + str(workspace), + ] + if runner_binding: + argv.extend( + [ + "--dsh-home", + str(dsh_home), + "--dsh-cordis", + str(cordis), + "--dsh-model", + "mock-model", + ] + ) + argv.extend( + [ + "--validation-command-json", + json.dumps(_validator_command()), + "--validation-failure-kind", + "repair_required", + "--scan-root", + str(project), + "--no-global-sync", + "--execute", + ] + ) + return argv + + +def _quota_spend_count(runtime: Path) -> int: + index = runtime / "goals" / GOAL_ID / "runs" / "index.jsonl" + if not index.is_file(): + return 0 + return sum( + 1 + for line in index.read_text(encoding="utf-8").splitlines() + if json.loads(line).get("classification") == "quota_slot_spent" + ) + + +class _UnavailableHarnessRuntime(importlib.abc.MetaPathFinder): + """Make the DeepSeek Harness runtime unimportable for one bounded call.""" + + def find_spec(self, fullname, path=None, target=None): + if fullname == "deepseek_harness" or fullname.startswith("deepseek_harness."): + raise ModuleNotFoundError(fullname) + return None + + +@contextlib.contextmanager +def _harness_runtime_unavailable() -> Iterator[None]: + """Hide an already-imported runtime for one bounded call. + + ``importlib.util.find_spec`` answers from ``sys.modules`` before consulting + ``sys.meta_path``, so a loader-only blocker cannot simulate an absent + runtime in a process that already imported the SDK for the skip check. + """ + + finder = _UnavailableHarnessRuntime() + hidden = { + name: module + for name, module in sys.modules.items() + if name == "deepseek_harness" or name.startswith("deepseek_harness.") + } + for name in hidden: + del sys.modules[name] + sys.meta_path.insert(0, finder) + try: + yield + finally: + sys.meta_path.remove(finder) + sys.modules.update(hidden) + + +@contextlib.contextmanager +def _operator_credential(value: str | None) -> Iterator[None]: + previous = os.environ.get(CREDENTIAL_ENV) + if value is None: + os.environ.pop(CREDENTIAL_ENV, None) + else: + os.environ[CREDENTIAL_ENV] = value + try: + yield + finally: + if previous is None: + os.environ.pop(CREDENTIAL_ENV, None) + else: + os.environ[CREDENTIAL_ENV] = previous + + +def main() -> int: + try: + harness_spec = importlib.util.find_spec("deepseek_harness") + except (ImportError, ValueError): + harness_spec = None + if harness_spec is None: + print("skip: deepseek-harness-sdk is not installed") + return 0 + + result_block = json.dumps( + { + "result_kind": "validated_progress", + "classification": "managed_default_flow_mock_llm", + "summary": "The default managed flow returned a typed result.", + "recommended_action": "Review the managed default flow marker.", + "next_action": "Inspect the marker and replay idempotently.", + "vision_unchanged_reason": "The objective path is unchanged.", + } + ) + + with tempfile.TemporaryDirectory(prefix="loopx-managed-default-flow-") as directory: + root = Path(directory) + project, runtime, workspace, registry = _write_fixture(root) + session_root = root / "sessions" + session_root.mkdir(parents=True) + dsh_home = root / "dsh-home" + cordis = _write_minimal_cordis(session_root) + marker_path = workspace / MARKER_NAME + + class MockHandler(BaseHTTPRequestHandler): + def do_POST(self) -> None: + length = int(self.headers.get("content-length", "0")) + self.rfile.read(length) + # The mock model is the only tool in this composition: it + # writes the marker that independent validation then checks. + marker_path.write_text(MARKER_VALUE, encoding="utf-8") + self.send_response(200) + self.send_header("content-type", "text/event-stream") + self.end_headers() + self.wfile.write( + ( + 'data: {"choices":[{"delta":{"role":"assistant","content":' + + json.dumps(result_block) + + "}}]}\n\n" + ).encode("utf-8") + ) + self.wfile.write( + b'data: {"choices":[{"delta":{"content":""},' + b'"finish_reason":"stop"}],"usage":{"prompt_tokens":2,' + b'"completion_tokens":1}}\n\n' + ) + self.wfile.write(b"data: [DONE]\n\n") + + def log_message(self, _format: str, *args: object) -> None: + return + + server = ThreadingHTTPServer(("127.0.0.1", 0), MockHandler) + thread = threading.Thread( + target=server.serve_forever, name="managed-default-mock-llm", daemon=True + ) + thread.start() + base_url = f"http://127.0.0.1:{server.server_address[1]}" + ambient = { + key: os.environ.get(key) + for key in ( + ENDPOINT_ENV, + CREDENTIAL_ENV, + "DSH_CWD", + "DSH_HOME", + "DSH_SESSION_ROOT", + ) + } + try: + os.environ[ENDPOINT_ENV] = base_url + for key in ("DSH_CWD", "DSH_HOME", "DSH_SESSION_ROOT"): + os.environ.pop(key, None) + + with _operator_credential(None): + unbound_default_exit, unbound_default_plan = _run_cli( + _plan_argv(registry, runtime, project) + ) + unbound_exit, unbound_plan = _run_cli( + _plan_argv( + registry, + runtime, + project, + host="dsh", + execution_mode="isolated-headless", + ) + ) + with _operator_credential("managed-default-flow-mock-key"): + managed_plan_exit, managed_plan = _run_cli( + _plan_argv(registry, runtime, project) + ) + individual_exit, individual_payload = _run_cli( + _plan_argv(registry, runtime, project, host="codex-cli") + ) + run_exit, run_payload = _run_cli( + _run_once_argv( + registry=registry, + runtime=runtime, + project=project, + workspace=workspace, + dsh_home=dsh_home, + cordis=cordis, + turn_instance_id="managed-default-flow-turn-1", + ) + ) + spend_count = _quota_spend_count(runtime) + with _harness_runtime_unavailable(): + unavailable_exit, unavailable_payload = _run_cli( + _run_once_argv( + registry=registry, + runtime=runtime, + project=project, + workspace=workspace, + dsh_home=dsh_home, + cordis=cordis, + turn_instance_id="managed-default-flow-turn-2", + runner_binding=False, + ) + ) + spend_after_refusal = _quota_spend_count(runtime) + marker_ok = ( + marker_path.is_file() + and marker_path.read_text(encoding="utf-8").strip() == MARKER_VALUE + ) + managed = managed_plan["managed_executor"] + unbound = unbound_plan.get("managed_executor") or {} + unbound_default = unbound_default_plan.get("managed_executor") or {} + individual = individual_payload["managed_executor"] + run_executor = run_payload.get("managed_executor") or {} + summary = { + "schema_version": "loopx_turn_managed_default_flow_v1", + "mock_llm_base_url": base_url, + "default_without_credential": { + "exit_code": unbound_default_exit, + "host_kind": unbound_default_plan.get("host", {}).get("kind"), + "execution_mode": unbound_default_plan.get("host", {}).get( + "execution_mode" + ), + "executor": unbound_default.get("executor"), + "executor_kind": unbound_default.get("executor_kind"), + "credential_env": unbound_default.get("credential_env"), + "operator_credential_bound": unbound_default.get( + "operator_credential_bound" + ), + "available": unbound_default.get("available"), + "unavailable_reason": unbound_default.get("unavailable_reason"), + }, + "explicit_individual_host": { + "exit_code": individual_exit, + "host_kind": individual_payload.get("host", {}).get("kind"), + "executor": individual.get("executor"), + "executor_kind": individual.get("executor_kind"), + "operator_credential_bound": individual.get( + "operator_credential_bound" + ), + "available": individual.get("available"), + }, + "managed_default_plan": { + "exit_code": managed_plan_exit, + "host_kind": managed_plan.get("host", {}).get("kind"), + "execution_mode": managed_plan.get("host", {}).get( + "execution_mode" + ), + "executor": managed.get("executor"), + "executor_kind": managed.get("executor_kind"), + "credential_env": managed.get("credential_env"), + "operator_credential_bound": managed.get( + "operator_credential_bound" + ), + "available": managed.get("available"), + "unavailable_reason": managed.get("unavailable_reason"), + }, + "explicit_host_without_credential": { + "exit_code": unbound_exit, + "host_kind": unbound_plan.get("host", {}).get("kind"), + "executor_kind": unbound.get("executor_kind"), + "credential_env": unbound.get("credential_env"), + "operator_credential_bound": unbound.get( + "operator_credential_bound" + ), + "available": unbound.get("available"), + "unavailable_reason": unbound.get("unavailable_reason"), + }, + "managed_default_run": { + "exit_code": run_exit, + "status": run_payload.get("status"), + "mode": run_payload.get("mode"), + "execution_mode": run_payload.get("execution_mode"), + "host_kind": run_payload.get("host", {}).get("kind"), + "executor": run_executor.get("executor"), + "executor_kind": run_executor.get("executor_kind"), + "operator_credential_bound": run_executor.get( + "operator_credential_bound" + ), + "result_kind": run_payload.get("result_kind"), + "validation_status": (run_payload.get("validation") or {}).get( + "status" + ), + "quota_slot_spend_count": run_payload.get("quota_slot_spend_count"), + "effects": run_payload.get("effects"), + "marker_valid": marker_ok, + }, + "quota_slot_spend_count": spend_count, + "quota_slot_spend_count_after_refusal": spend_after_refusal, + "managed_runtime_unavailable": { + "exit_code": unavailable_exit, + "status": unavailable_payload.get("status"), + "reason": unavailable_payload.get("reason"), + "effects": unavailable_payload.get("effects"), + "executor_kind": unavailable_payload.get( + "managed_executor", {} + ).get("executor_kind"), + }, + "global_registry_synced": False, + } + finally: + for key, value in ambient.items(): + if value is None: + os.environ.pop(key, None) + else: + os.environ[key] = value + server.shutdown() + server.server_close() + + print(json.dumps(summary, indent=2, sort_keys=True)) + + effects = summary["managed_default_run"]["effects"] or {} + ok = ( + unbound_default_exit == 0 + and summary["default_without_credential"]["host_kind"] == "dsh" + and summary["default_without_credential"]["execution_mode"] + == "isolated-headless" + and summary["default_without_credential"]["executor_kind"] + == EXECUTOR_KIND_MANAGED + and summary["default_without_credential"]["credential_env"] is None + and summary["default_without_credential"]["operator_credential_bound"] is False + and summary["default_without_credential"]["available"] is False + and summary["default_without_credential"]["unavailable_reason"] + == OPERATOR_CREDENTIAL_UNCONFIGURED + and individual_exit == 0 + and summary["explicit_individual_host"]["host_kind"] == "codex-cli" + and summary["explicit_individual_host"]["executor_kind"] + == EXECUTOR_KIND_INDIVIDUAL + and summary["explicit_individual_host"]["operator_credential_bound"] is False + and summary["explicit_individual_host"]["available"] is None + and managed_plan_exit == 0 + and summary["managed_default_plan"]["host_kind"] == "dsh" + and summary["managed_default_plan"]["execution_mode"] == "isolated-headless" + and summary["managed_default_plan"]["executor_kind"] == EXECUTOR_KIND_MANAGED + and summary["managed_default_plan"]["credential_env"] == CREDENTIAL_ENV + and summary["managed_default_plan"]["operator_credential_bound"] is True + and summary["explicit_host_without_credential"]["exit_code"] == 0 + and summary["explicit_host_without_credential"]["host_kind"] == "dsh" + and summary["explicit_host_without_credential"]["credential_env"] is None + and summary["explicit_host_without_credential"]["operator_credential_bound"] + is False + and summary["explicit_host_without_credential"]["available"] is False + and summary["explicit_host_without_credential"]["unavailable_reason"] + == OPERATOR_CREDENTIAL_UNCONFIGURED + and run_exit == 0 + and summary["managed_default_run"]["status"] == "committed" + and summary["managed_default_run"]["host_kind"] == "dsh" + and summary["managed_default_run"]["execution_mode"] == "isolated-headless" + and summary["managed_default_run"]["executor_kind"] == EXECUTOR_KIND_MANAGED + and summary["managed_default_run"]["operator_credential_bound"] is True + and summary["managed_default_run"]["quota_slot_spend_count"] == 1 + and summary["managed_default_run"]["validation_status"] == "passed" + and effects + == { + "host_invoked": True, + "state_written": True, + "quota_spent": True, + "scheduler_acknowledged": False, + } + and summary["managed_default_run"]["marker_valid"] is True + and spend_count == 1 + and unavailable_exit != 0 + and summary["managed_runtime_unavailable"]["status"] == "unavailable" + and summary["managed_runtime_unavailable"]["reason"] == DSH_RUNTIME_UNAVAILABLE + and summary["managed_runtime_unavailable"]["executor_kind"] + == EXECUTOR_KIND_MANAGED + and (summary["managed_runtime_unavailable"]["effects"] or {}) + == { + "host_invoked": False, + "state_written": False, + "quota_spent": False, + "scheduler_acknowledged": False, + } + and spend_after_refusal == 1 + ) + if not ok: + print("managed default flow smoke failed") + return 1 + print("managed default flow smoke passed") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/examples/loopx-turn-managed-executor-binding-smoke.py b/examples/loopx-turn-managed-executor-binding-smoke.py new file mode 100644 index 0000000000..56e71afbbc --- /dev/null +++ b/examples/loopx-turn-managed-executor-binding-smoke.py @@ -0,0 +1,365 @@ +#!/usr/bin/env python3 +"""Prove explicit Turn host selection, its readback, and the fail-closed start. + +The smoke drives the public CLI only. It never calls a provider and never reads +a credential value: the operator credential is a fixture string whose only role +is to authenticate the host the operator already selected. +""" + +from __future__ import annotations + +import contextlib +import importlib.abc +import importlib.machinery +import io +import json +import os +import sys +import tempfile +from collections.abc import Iterator +from pathlib import Path +from typing import Any + + +REPO_ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(REPO_ROOT)) + +from loopx.cli import main as cli_main # noqa: E402 +from loopx.control_plane.turn_driver import executor as turn_executor # noqa: E402 +from loopx.control_plane.turn_driver.host_binding import ( # noqa: E402 + DSH_RUNTIME_UNAVAILABLE, + EXECUTOR_KIND_INDIVIDUAL, + EXECUTOR_KIND_MANAGED, + OPERATOR_CREDENTIAL_UNCONFIGURED, +) + + +GOAL_ID = "loopx-turn-managed-executor-fixture" +AGENT_ID = "codex-managed-executor-fixture" +TODO_ID = "todo_managedexec01" +CREDENTIAL_ENV = "DEEPSEEK_API_KEY" +RUNTIME_MODULE = "deepseek_harness" + + +def _write_fixture(root: Path) -> tuple[Path, Path, Path, Path]: + project = root / "project" + runtime = root / "runtime" + workspace = root / "workspace" + runtime.mkdir(parents=True) + workspace.mkdir(parents=True) + + state = project / ".codex" / "goals" / GOAL_ID / "ACTIVE_GOAL_STATE.md" + state.parent.mkdir(parents=True) + state.write_text( + "\n".join( + [ + "---", + "status: active", + "updated_at: 2026-01-01T00:00:00+00:00", + "---", + "", + "# LoopX Managed Executor Fixture", + "", + "## Next Action", + "", + "Run one bounded managed Turn on the planned executor only.", + "", + "## Agent Todo", + "", + "- [ ] [P1] Run one bounded managed Turn without leaving the planned executor.", + ( + f" " + ), + "", + ] + ), + encoding="utf-8", + ) + registry = project / ".loopx" / "registry.json" + registry.parent.mkdir(parents=True) + registry.write_text( + json.dumps( + { + "schema_version": 1, + "common_runtime_root": str(runtime), + "goals": [ + { + "id": GOAL_ID, + "domain": "loopx-turn-public-fixture", + "status": "active", + "repo": str(project), + "state_file": str(state.relative_to(project)), + "adapter": { + "kind": "fixture_v0", + "status": "connected-delivery", + }, + "quota": {"compute": 10.0, "window_hours": 24}, + "coordination": { + "agent_model": "peer_v1", + "registered_agents": [AGENT_ID], + "agent_profiles": { + AGENT_ID: { + "schema_version": "agent_profile_v1", + "profile_role": "fixture", + "scope": "public qualification", + } + }, + "write_scope": ["docs/**"], + }, + } + ], + }, + indent=2, + sort_keys=True, + ) + + "\n", + encoding="utf-8", + ) + return project, runtime, workspace, registry + + +class _HarnessRuntimeFinder(importlib.abc.MetaPathFinder): + """Pin whether the DeepSeek Harness runtime resolves, for one bounded call.""" + + def __init__(self, available: bool) -> None: + self._available = available + + def find_spec(self, fullname, path=None, target=None): + if fullname != RUNTIME_MODULE and not fullname.startswith(f"{RUNTIME_MODULE}."): + return None + if not self._available: + raise ModuleNotFoundError(fullname) + return importlib.machinery.ModuleSpec( + fullname, loader=None, is_package=fullname == RUNTIME_MODULE + ) + + +@contextlib.contextmanager +def _harness_runtime(*, available: bool) -> Iterator[None]: + finder = _HarnessRuntimeFinder(available) + sys.meta_path.insert(0, finder) + sys.modules.pop(RUNTIME_MODULE, None) + try: + yield + finally: + sys.meta_path.remove(finder) + sys.modules.pop(RUNTIME_MODULE, None) + + +@contextlib.contextmanager +def _operator_credential(value: str | None) -> Iterator[None]: + previous = os.environ.get(CREDENTIAL_ENV) + if value is None: + os.environ.pop(CREDENTIAL_ENV, None) + else: + os.environ[CREDENTIAL_ENV] = value + try: + yield + finally: + if previous is None: + os.environ.pop(CREDENTIAL_ENV, None) + else: + os.environ[CREDENTIAL_ENV] = previous + + +def _run_cli(argv: list[str]) -> tuple[int, dict[str, Any]]: + output = io.StringIO() + with contextlib.redirect_stdout(output): + exit_code = cli_main(argv) + return exit_code, json.loads(output.getvalue()) + + +def _plan_command( + registry: Path, runtime: Path, project: Path, *, host: str | None = None +) -> list[str]: + argv = [ + "--registry", + str(registry), + "--runtime-root", + str(runtime), + "--format", + "json", + "turn", + "plan", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--scan-root", + str(project), + ] + return [*argv, "--host", host] if host else argv + + +def _run_once_command( + registry: Path, + runtime: Path, + project: Path, + workspace: Path, + *, + instance: str, + host: str | None = None, +) -> list[str]: + argv = [ + "--registry", + str(registry), + "--runtime-root", + str(runtime), + "--format", + "json", + "turn", + "run-once", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--turn-instance-id", + instance, + "--project", + str(workspace), + "--scan-root", + str(project), + "--no-global-sync", + "--execute", + ] + return [*argv, "--host", host] if host else argv + + +def _managed_binding(payload: dict[str, Any]) -> dict[str, Any]: + binding = payload["managed_executor"] + assert binding["schema_version"] == "managed_executor_binding_v0", binding + assert binding["executor"] == payload["host"]["kind"], binding + assert binding["executor_kind"] == EXECUTOR_KIND_MANAGED, binding + assert (binding["unavailable_reason"] is None) is binding["available"], binding + return binding + + +def main() -> int: + with tempfile.TemporaryDirectory( + prefix="loopx-turn-managed-executor-" + ) as directory: + root = Path(directory) + project, runtime, workspace, registry = _write_fixture(root) + + # 1. The shipped default is the managed host, and without the operator + # credential it refuses instead of borrowing a personal login. + with _operator_credential(None), _harness_runtime(available=True): + exit_code, payload = _run_cli(_plan_command(registry, runtime, project)) + assert exit_code == 0, payload + assert payload["host"]["kind"] == "dsh", payload + unbound = _managed_binding(payload) + assert unbound["operator_credential_bound"] is False, unbound + assert unbound["available"] is False, unbound + assert unbound["unavailable_reason"] == OPERATOR_CREDENTIAL_UNCONFIGURED, ( + unbound + ) + + # 2. Configuring the credential authenticates that same selection; it + # does not get to pick a different host. + with ( + _operator_credential("sk-fixture-operator"), + _harness_runtime(available=True), + ): + exit_code, payload = _run_cli(_plan_command(registry, runtime, project)) + assert exit_code == 0, payload + assert payload["host"]["kind"] == "dsh", payload + bound = _managed_binding(payload) + assert bound["credential_env"] == CREDENTIAL_ENV, bound + assert bound["operator_credential_bound"] is True, bound + assert bound["available"] is True, bound + + # 3. With the runtime genuinely missing the same plan reports the other + # typed reason rather than promising a launch. + with ( + _operator_credential("sk-fixture-operator"), + _harness_runtime(available=False), + ): + exit_code, payload = _run_cli(_plan_command(registry, runtime, project)) + assert exit_code == 0, payload + missing = _managed_binding(payload) + assert missing["available"] is False, missing + assert missing["unavailable_reason"] == DSH_RUNTIME_UNAVAILABLE, missing + + # 4. An explicit individual host stays selected even while the operator + # credential is configured, and makes no launch claim. + with ( + _operator_credential("sk-fixture-operator"), + _harness_runtime(available=True), + ): + exit_code, payload = _run_cli( + _plan_command(registry, runtime, project, host="codex-cli") + ) + assert exit_code == 0, payload + assert payload["host"]["kind"] == "codex-cli", payload + individual = payload["managed_executor"] + assert individual["executor_kind"] == EXECUTOR_KIND_INDIVIDUAL, individual + assert individual["available"] is None, individual + assert individual["operator_credential_bound"] is False, individual + + # 5. Executing the unauthenticated managed default fails closed: typed + # status, no host invocation, no journal, and no quota slot spend. + with _operator_credential(None), _harness_runtime(available=True): + exit_code, refusal = _run_cli( + _run_once_command( + registry, + runtime, + project, + workspace, + instance="managed-executor-unauthenticated", + ) + ) + assert exit_code == 1, refusal + assert refusal["ok"] is False, refusal + assert refusal["status"] == "unavailable", refusal + assert refusal["reason"] == OPERATOR_CREDENTIAL_UNCONFIGURED, refusal + _expect_no_effects(refusal) + _expect_no_journal(refusal, runtime) + + # 6. The same refusal covers a provably unlaunchable runtime. + with ( + _operator_credential("sk-fixture-operator"), + _harness_runtime(available=False), + ): + exit_code, refusal = _run_cli( + _run_once_command( + registry, + runtime, + project, + workspace, + instance="managed-executor-runtime-missing", + ) + ) + assert exit_code == 1, refusal + assert refusal["ok"] is False, refusal + assert refusal["status"] == "unavailable", refusal + assert refusal["reason"] == DSH_RUNTIME_UNAVAILABLE, refusal + _expect_no_effects(refusal) + _expect_no_journal(refusal, runtime) + + print("managed executor binding smoke passed") + return 0 + + +def _expect_no_effects(payload: dict[str, Any]) -> None: + assert payload["effects"] == { + "host_invoked": False, + "state_written": False, + "quota_spent": False, + "scheduler_acknowledged": False, + }, payload + assert payload["quota_slot_spend_count"] == 0, payload + + +def _expect_no_journal(payload: dict[str, Any], runtime: Path) -> None: + journal = turn_executor.turn_journal_path( + runtime, + goal_id=GOAL_ID, + turn_key=str(payload["resume_turn_key"]), + ) + assert journal.exists() is False, journal + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/examples/loopx-turn-path-delta-acceptance-smoke.py b/examples/loopx-turn-path-delta-acceptance-smoke.py index 4c5c97f4cd..e6193a6b9b 100644 --- a/examples/loopx-turn-path-delta-acceptance-smoke.py +++ b/examples/loopx-turn-path-delta-acceptance-smoke.py @@ -194,6 +194,8 @@ def _run_turn( "json", "turn", "run-once", + "--host", + "generic-cli", "--goal-id", GOAL_ID, "--agent-id", diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 9d3236d65f..f95247bff0 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -18,6 +18,7 @@ run_configured_turn_outcome_ingest_fail_open, ) from ..capabilities.periodic_report.cadence_runtime import extend_cadence_turn_start_dispatch +from ..capabilities.periodic_report.pending_intent import periodic_report_pending_intent_interaction_hook from ..control_plane.quota.live_decision import build_live_quota_should_run_decision from ..control_plane.agents.workspace_guard import capture_delivery_workspace from ..control_plane.quota.heartbeat_receipt import ( @@ -54,6 +55,7 @@ run_loopx_turn_once, selected_turn_todo, ) +from ..control_plane.turn_driver.host_binding import managed_executor_binding from ..quota import spend_quota_slot from ..state_refresh import refresh_state_run from ..status import AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, collect_status @@ -65,16 +67,11 @@ from .turn_dsh_host import build_dsh_host_runner from .turn_registration import register_turn_commands as register_turn_commands from .turn_inspection import handle_turn_journal_inspection -from .turn_decision import ( - apply_controller_advisory_primary, - build_turn_decision_builder, - collect_turn_status_payload, -) -from .turn_managed_step import handle_turn_managed_step from .turn_rendering import ( render_loopx_turn_execution_markdown as _render_loopx_turn_execution_markdown, render_loopx_turn_plan_markdown as _render_loopx_turn_plan_markdown, ) +from .turn_selection import turn_controller_advisory_primary from .turn_todo_writeback import ( write_turn_repair_update, write_turn_validated_completion, @@ -112,16 +109,10 @@ def handle_turn_command( ) if inspection_result is not None: return inspection_result - managed_step_result = handle_turn_managed_step( - args, - registry_path=registry_path, - runtime_root_arg=runtime_root_arg, - output_format=output_format, - print_payload=print_payload, - ) - if managed_step_result is not None: - return managed_step_result try: + scan_roots = [Path(item).expanduser() for item in args.scan_path] + if not scan_roots: + scan_roots = [Path(args.scan_root).expanduser()] runtime_root = resolve_status_projection_cache_runtime_root( registry_path=registry_path, runtime_root_override=runtime_root_arg, @@ -154,30 +145,60 @@ def handle_turn_command( operator_inbox_urgency_projector = build_lark_operator_inbox_urgency_projector( runtime_root_arg=runtime_root, ) - scan_roots = [Path(item).expanduser() for item in args.scan_path] - if not scan_roots: - scan_roots = [Path(args.scan_root).expanduser()] + status_payload = collect_status( + registry_path=registry_path, + runtime_root_override=runtime_root_arg, + scan_roots=scan_roots, + limit=max(max(0, args.limit), AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK), + goal_id=args.goal_id, + available_capabilities=args.available_capabilities, + ) scheduler_context = scheduler_execution_context_for_turn( host=args.host, execution_mode=args.execution_mode, scheduler_owner=args.scheduler_owner, ) - # Decision building is shared with the managed step so both owners - # always answer from the same control-plane projection. - decision = apply_controller_advisory_primary( - build_turn_decision_builder( - args, + def build_turn_decision( + *, requested_action_todo_id: str | None = None + ) -> dict[str, Any]: + return build_live_quota_should_run_decision( + status_payload, + goal_id=args.goal_id, + agent_id=args.agent_id, + available_capabilities=args.available_capabilities, + include_scheduler_detail=False, + codex_app_current_rrule=None, registry_path=registry_path, runtime_root=runtime_root, - runtime_root_arg=runtime_root_arg, - status_payload=collect_turn_status_payload( - args, - registry_path=registry_path, - runtime_root_arg=runtime_root_arg, + route_source="loopx_turn_plan", + scheduler_execution_context=scheduler_context, + operator_inbox_urgency_projector=operator_inbox_urgency_projector, + bounded_research_frontier_projector=( + project_live_explore_composition_frontier ), + requested_action_todo_id=requested_action_todo_id, turn_start_hook_dispatch=turn_start_hook_dispatch, + interaction_projection_hooks=(periodic_report_pending_intent_interaction_hook( + registry_path=registry_path, runtime_root=runtime_root, + goal_id=args.goal_id, agent_id=args.agent_id),), ) - ) + + decision = build_turn_decision() + controller_default = turn_controller_advisory_primary(decision) + if controller_default is not None: + primary_todo_id, advisory_portfolio = controller_default + decision = build_turn_decision( + requested_action_todo_id=primary_todo_id, + ) + selected_todo = decision.get("selected_todo") + if not isinstance(selected_todo, dict) or ( + selected_todo.get("todo_id") != primary_todo_id + ): + raise ValueError( + "Turn controller advisory primary failed current eligibility" + ) + selected_todo["selected_by"] = "turn_controller_advisory_primary" + decision["action_portfolio"] = advisory_portfolio resume_identity = { "goal_id": args.resume_goal_id, "agent_id": args.resume_agent_id, @@ -219,6 +240,14 @@ def handle_turn_command( turn_instance_id=args.turn_instance_id, iteration_context_policy=args.iteration_context.replace("-", "_"), ) + # The executor readback names where this Turn's model work runs and + # whether that host can launch here, so a caller never has to infer it + # from the host id. The explicit runner hook is the one launchability + # fact only this command layer knows. + payload["managed_executor"] = managed_executor_binding( + args.host, + dsh_runner_configured=bool(getattr(args, "dsh_runner", None)), + ) if ( args.turn_command == "run-once" and args.execute diff --git a/loopx/cli_commands/turn_registration.py b/loopx/cli_commands/turn_registration.py index 524693460b..304654f03d 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -5,8 +5,16 @@ import argparse from collections.abc import Callable +from ..control_plane.turn_driver.host_binding import ( + MANAGED_TURN_HOST, + resolve_default_turn_host, +) from ..paths import default_public_scan_root +# Explicit host choices stay per-command: planning may name any host the Turn +# driver routes, while run-once only ships built-in adapters for these three. +PLANNED_TURN_HOST_CHOICES = ["codex-cli", "claude-code", "dsh", "generic-cli"] +RUN_ONCE_TURN_HOST_CHOICES = ["codex-cli", "dsh", "generic-cli"] AddFormat = Callable[[argparse.ArgumentParser], None] @@ -42,7 +50,21 @@ def register_turn_commands( help="Build one typed read-only host decision without launching or writing.", ) add_subcommand_format(plan) - _add_turn_decision_arguments(plan, default_host="codex-cli") + # The default host and the default execution mode are one decision: the + # selected managed host runs bounded headless Turns, so pairing it with a + # visible interactive mode would produce a default plan that cannot be + # scheduled. The mode follows the *selected* host, never the environment. + resolved_default_host = resolve_default_turn_host() + _add_turn_decision_arguments( + plan, + default_host=resolved_default_host, + host_choices=list(PLANNED_TURN_HOST_CHOICES), + default_execution_mode=( + "isolated-headless" + if resolved_default_host == MANAGED_TURN_HOST + else "interactive-visible" + ), + ) plan.add_argument( "--include-transaction-detail", action="store_true", @@ -61,62 +83,6 @@ def register_turn_commands( ) plan.add_argument("--limit", type=int, default=5) - managed_step = command_sub.add_parser( - "managed-step", - help=( - "Decide one bounded same-Turn continuation for a failed Turn " - "without executing it." - ), - description=( - "Read one canonical Turn journal, rebuild its validated receipt, " - "and ask the pure Turn Loop Controller for a disposition against " - "the current decision. Grants no execution authority: it never " - "launches a host, writes state, or spends quota. The Turn journal " - "remains the authority for the attempt count and retry budget." - ), - ) - add_subcommand_format(managed_step) - _add_turn_decision_arguments( - managed_step, - default_host="dsh", - host_choices=["codex-cli", "dsh", "generic-cli"], - execution_mode_choices=["isolated-headless"], - default_execution_mode="isolated-headless", - ) - managed_step.add_argument( - "--turn-key", - required=True, - help="Exact sha256 Turn key of the failed Turn to decide about.", - ) - managed_step.add_argument( - "--observed-attempt", - type=int, - help=( - "Caller's observed attempt count, reconciled against the Turn " - "journal. A disagreement is refused rather than adopted." - ), - ) - managed_step.add_argument( - "--observed-max-attempts", - type=int, - help=( - "Caller's observed retry ceiling, reconciled against the Turn " - "journal retry policy." - ), - ) - managed_step.add_argument( - "--scan-root", - default=default_public_scan_root(), - help="Public files to scan for obvious private material.", - ) - managed_step.add_argument( - "--scan-path", - action="append", - default=[], - help="Specific public file or directory to scan. Repeatable.", - ) - managed_step.add_argument("--limit", type=int, default=5) - run_once = command_sub.add_parser( "run-once", help=( @@ -140,8 +106,8 @@ def register_turn_commands( add_subcommand_format(run_once) _add_turn_decision_arguments( run_once, - default_host="generic-cli", - host_choices=["codex-cli", "dsh", "generic-cli"], + default_host=resolve_default_turn_host(), + host_choices=list(RUN_ONCE_TURN_HOST_CHOICES), execution_mode_choices=["isolated-headless"], default_execution_mode="isolated-headless", ) diff --git a/loopx/control_plane/operator_credential.py b/loopx/control_plane/operator_credential.py new file mode 100644 index 0000000000..0d47e889f6 --- /dev/null +++ b/loopx/control_plane/operator_credential.py @@ -0,0 +1,52 @@ +"""Operator-supplied model credential facts shared by LoopX host surfaces. + +This module reports credential *facts* and nothing else. It never selects a +host, an endpoint, or a model, and it never reads a credential value. + +Selection is a separate, explicit decision owned by the surface that runs the +work: the governed Turn host comes from +``turn_driver.host_binding.selected_turn_host`` and the steward channel endpoint +comes from ``chat_manager.manager_channel_binding``. Both report the credential +facts quoted from here so their readback cannot drift apart, and both treat the +credential as authentication for the configuration the operator selected -- +never as a reason to change it. Discovering that a credential exists may help +the operator set a surface up, but it must not silently re-point a surface that +is already configured. +""" + +from __future__ import annotations + +import os +from collections.abc import Mapping + +# Credential env vars the operator-supplied provider already reads. Only the +# variable name is ever reported back; values stay in the process environment. +OPERATOR_CREDENTIAL_ENV_VARS = ("DEEPSEEK_API_KEY",) +OPERATOR_ENDPOINT_ENV_VAR = "DEEPSEEK_BASE_URL" + + +def env_text(name: str, environ: Mapping[str, str] | None = None) -> str | None: + """Return a stripped env value, or ``None`` when it is unset or blank.""" + + source = os.environ if environ is None else environ + value = str(source.get(name, "") or "").strip() + return value or None + + +def configured_operator_credential( + environ: Mapping[str, str] | None = None, +) -> str | None: + """Return the configured operator credential env var name, else ``None``.""" + + for name in OPERATOR_CREDENTIAL_ENV_VARS: + if env_text(name, environ) is not None: + return name + return None + + +def operator_credential_configured( + environ: Mapping[str, str] | None = None, +) -> bool: + """Whether any operator model credential is configured for this process.""" + + return configured_operator_credential(environ) is not None diff --git a/loopx/control_plane/testing/cli_output_differential.py b/loopx/control_plane/testing/cli_output_differential.py index b1315e5357..415cc2775e 100644 --- a/loopx/control_plane/testing/cli_output_differential.py +++ b/loopx/control_plane/testing/cli_output_differential.py @@ -23,9 +23,7 @@ ACTION_PORTFOLIO_SCHEMA_VERSION_V2 = "quota_action_portfolio_v2" PLANNING_HORIZON_SCHEMA_VERSION_V0 = "quota_planning_horizon_v0" GUIDED_TODO_DELTA_SCHEMA_VERSION_V0 = "loopx_guided_todo_delta_v0" -PLANNING_INVENTORY_DETAIL_SCHEMA_VERSION_V0 = ( - "todo_planning_inventory_detail_v0" -) +PLANNING_INVENTORY_DETAIL_SCHEMA_VERSION_V0 = "todo_planning_inventory_detail_v0" Metric = Literal["chars", "utf8_bytes", "lines", "compact_payload_chars"] @@ -177,6 +175,55 @@ class GrowthAllowance: "compact_payload_chars": 640, } +# Two reviewed causes grow the Turn plan readback once, and both are consequences +# of the same declared behavior change: +# +# 1. the plan adds one bounded managed-executor binding so a caller sees which +# executor a planned Turn would use and whether it can launch here, instead +# of inferring it from the host id; +# 2. the default host becomes the managed `dsh` host, so the plan's host +# projection, scheduler execution context, and controller-owned +# `next_cli_actions` rendering change with it (a bounded-headless plan hands +# the next steps back to the outer controller instead of to the agent CLI +# loop). +# +# The allowance is bound to the rejected-then-accepted +# none-to-v0 transition and to the Turn surfaces that quote it; quota, status, +# and every other agent-facing surface keep the ordinary budget. It is one +# time: once v0 is the baseline a v0-to-v0 change receives no allowance. +_TURN_HOST_AND_MANAGED_EXECUTOR_BINDING_V0_GROWTH_ALLOWANCE: dict[Metric, int] = { + "chars": 832, + "utf8_bytes": 832, + "lines": 14, + "compact_payload_chars": 768, +} + +_MANAGED_EXECUTOR_BINDING_SURFACES = frozenset( + { + "loopx_turn_plan", + "loopx_turn_plan_transaction_detail", + "loopx_turn_run_once_preview", + } +) + + +def _turn_host_and_managed_executor_binding_allowance( + row_id: str, + base: Mapping[str, Any], + candidate: Mapping[str, Any], + metric: Metric, +) -> int: + surface = row_id.partition("/")[2].partition("/")[0] + if ( + row_id.startswith(("surface/", "variant/")) + and surface in _MANAGED_EXECUTOR_BINDING_SURFACES + and base.get("managed_executor_binding_revision") is None + and candidate.get("managed_executor_binding_revision") + == "managed_executor_binding_v0" + ): + return _TURN_HOST_AND_MANAGED_EXECUTOR_BINDING_V0_GROWTH_ALLOWANCE[metric] + return 0 + def _reward_memory_outcome_prompt_allowance( row_id: str, @@ -201,6 +248,7 @@ def _reward_memory_outcome_prompt_allowance( return _REWARD_MEMORY_OUTCOME_PROMPT_V1_MIGRATION_ALLOWANCE[metric] return 0 + # loopx_guided_todo_delta_v0 adds the continuation-aware Todo authoring # decision contract (reuse/update/link_successor/add_new plus a bounded # runnable-frontier summary) to the guided start-goal packet when an @@ -330,9 +378,7 @@ def _planning_horizon_schema_migration( base: dict[str, Any], candidate: dict[str, Any] ) -> str | None: base_versions = tuple(base.get("planning_horizon_schema_versions") or []) - candidate_versions = tuple( - candidate.get("planning_horizon_schema_versions") or [] - ) + candidate_versions = tuple(candidate.get("planning_horizon_schema_versions") or []) if base_versions == () and candidate_versions == ( PLANNING_HORIZON_SCHEMA_VERSION_V0, ): @@ -344,9 +390,7 @@ def _guided_todo_delta_schema_migration( base: dict[str, Any], candidate: dict[str, Any] ) -> str | None: base_versions = tuple(base.get("guided_todo_delta_schema_versions") or []) - candidate_versions = tuple( - candidate.get("guided_todo_delta_schema_versions") or [] - ) + candidate_versions = tuple(candidate.get("guided_todo_delta_schema_versions") or []) if base_versions == () and candidate_versions == ( GUIDED_TODO_DELTA_SCHEMA_VERSION_V0, ): @@ -357,9 +401,7 @@ def _guided_todo_delta_schema_migration( def _planning_inventory_detail_schema_migration( base: dict[str, Any], candidate: dict[str, Any] ) -> str | None: - base_versions = tuple( - base.get("planning_inventory_detail_schema_versions") or [] - ) + base_versions = tuple(base.get("planning_inventory_detail_schema_versions") or []) candidate_versions = tuple( candidate.get("planning_inventory_detail_schema_versions") or [] ) @@ -394,8 +436,7 @@ def _schema_migration_state( ) -> _SchemaMigrationState: base_signature = base.get("action_signature_sha256") signature_changed = bool( - base_signature - and candidate.get("action_signature_sha256") != base_signature + base_signature and candidate.get("action_signature_sha256") != base_signature ) signature_migration = ( _action_signature_migration(base, candidate) if signature_changed else None @@ -430,10 +471,9 @@ def _schema_migration_state( if inventory_detail_schema_changed else None ) - guided_todo_delta_schema_changed = ( - tuple(base.get("guided_todo_delta_schema_versions") or []) - != tuple(candidate.get("guided_todo_delta_schema_versions") or []) - ) + guided_todo_delta_schema_changed = tuple( + base.get("guided_todo_delta_schema_versions") or [] + ) != tuple(candidate.get("guided_todo_delta_schema_versions") or []) guided_todo_delta_schema_migration = ( _guided_todo_delta_schema_migration(base, candidate) if guided_todo_delta_schema_changed @@ -545,23 +585,46 @@ def _compare_row(base: dict[str, Any], candidate: dict[str, Any]) -> dict[str, A candidate, metric, ), + _turn_host_and_managed_executor_binding_allowance( + row_id, + base, + candidate, + metric, + ), ) # Thin installed prompts contain bilingual lifecycle instructions. A # small character-level clarification can cost three bytes per CJK # character. Keep character, line and absolute output ceilings intact; # do not relax quota or other agent-facing surfaces with this allowance. - if row_id.startswith("surface/heartbeat_prompt_thin/") and metric == "utf8_bytes": + if ( + row_id.startswith("surface/heartbeat_prompt_thin/") + and metric == "utf8_bytes" + ): allowance = max(allowance, 192) - if (row_id.startswith(("surface/", "variant/")) - and row_id.partition("/")[2].partition("/")[0] in { - "heartbeat_prompt_thin", "heartbeat_prompt_brief", "heartbeat_prompt_compact"} - and base.get("host_prompt_static_safety_revision") is None - and candidate.get("host_prompt_static_safety_revision") == "host_prompt_static_safety_v1"): + if ( + row_id.startswith(("surface/", "variant/")) + and row_id.partition("/")[2].partition("/")[0] + in { + "heartbeat_prompt_thin", + "heartbeat_prompt_brief", + "heartbeat_prompt_compact", + } + and base.get("host_prompt_static_safety_revision") is None + and candidate.get("host_prompt_static_safety_revision") + == "host_prompt_static_safety_v1" + ): # Authorized static safety + executable shell bootstrap restoration. # Absolute ceilings stay enforced by the probe; once merged, v1->v1 # receives no allowance. Quota/status and other surfaces are excluded. - allowance = max(allowance, {"chars": 512, "utf8_bytes": 640, - "lines": 5, "compact_payload_chars": 512}[metric]) + allowance = max( + allowance, + { + "chars": 512, + "utf8_bytes": 640, + "lines": 5, + "compact_payload_chars": 512, + }[metric], + ) if migration.portfolio_growth_migration: allowance = max( allowance, @@ -579,9 +642,7 @@ def _compare_row(base: dict[str, Any], candidate: dict[str, Any]) -> dict[str, A if migration.inventory_detail_growth_migration: allowance = max( allowance, - _PLANNING_INVENTORY_DETAIL_V0_MIGRATION_GROWTH_ALLOWANCE[ - metric - ], + _PLANNING_INVENTORY_DETAIL_V0_MIGRATION_GROWTH_ALLOWANCE[metric], ) if migration.guided_todo_delta_growth_migration: allowance = max( diff --git a/loopx/control_plane/testing/cli_output_semantics.py b/loopx/control_plane/testing/cli_output_semantics.py index 412f3bea19..fa32f1f19f 100644 --- a/loopx/control_plane/testing/cli_output_semantics.py +++ b/loopx/control_plane/testing/cli_output_semantics.py @@ -41,6 +41,28 @@ def reward_memory_outcome_prompt_revision(text: str) -> str | None: else None ) + +def managed_executor_binding_revision(text: str) -> str | None: + """Attribute the managed-executor binding readback on a Turn surface. + + This is qualification evidence for the exact projection, never a runtime + classifier: the binding key alone would match prose, so the revision also + requires the executor identity, its launchability claim, and the typed + reason slot that only this readback renders. + """ + + required = ( + '"managed_executor"', + '"executor_kind"', + '"available"', + '"unavailable_reason"', + ) + return ( + "managed_executor_binding_v0" + if all(fragment in text for fragment in required) + else None + ) + _MARKDOWN_HEADING = re.compile(r"^#{1,6}\s+.+$") _RUNTIME_ROOT_COMMAND_ROUTE = re.compile( r"(?m)(?:^|[\"'`])[^\r\n\S]*loopx\s+--runtime-root\s+" diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 634b26a854..854586f32d 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -28,6 +28,10 @@ reward_memory_reflection_digest, ) from .driver import selected_turn_todo +from .host_binding import ( + managed_executor_payload_entry, + managed_executor_unavailable_payload, +) from .host_failure import BuiltInHostError, project_host_failure, record_host_failure from .journal_store import ( LOOPX_TURN_JOURNAL_SCHEMA_VERSION, @@ -778,6 +782,7 @@ def _execution_payload( "status": journal.get("status"), "execution_mode": planned_host.get("execution_mode"), "host": journal.get("host"), + **managed_executor_payload_entry(plan), "result_kind": journal.get("result_kind"), "validation": journal.get("task_validation"), "receipt": journal.get("receipt"), @@ -1306,6 +1311,15 @@ def run_loopx_turn_once( "quota_spent": False, "scheduler_acknowledged": False, } + fail_closed = managed_executor_unavailable_payload( + plan, execute=execute, host_projection=host_projection + ) + if fail_closed is not None: + # Fail closed on an executor LoopX can prove cannot launch: report the + # planned executor and stop before the journal, host, and quota. + return _execution_payload( + plan, fail_closed, execute=True, replayed=False, effects=empty_effects + ) if not execute: preview = { "schema_version": LOOPX_TURN_JOURNAL_SCHEMA_VERSION, diff --git a/loopx/control_plane/turn_driver/host_binding.py b/loopx/control_plane/turn_driver/host_binding.py new file mode 100644 index 0000000000..278f040301 --- /dev/null +++ b/loopx/control_plane/turn_driver/host_binding.py @@ -0,0 +1,196 @@ +"""Explicit Turn host selection and managed executor readback. + +The Turn host is **selected, never inferred**. LoopX ships one explicit product +default, the operator may override it explicitly, and a discovered credential +only authenticates the host that was already selected. The presence of +``DEEPSEEK_API_KEY`` therefore never changes where a Turn runs; setting it is +what makes the selected managed host authenticated. + +- ``MANAGED_DEFAULT_TURN_HOST`` (``dsh``) is the shipped default: the managed + execution unit the steward drives runs on the DeepSeek Harness host. +- ``LOOPX_TURN_HOST`` re-points that default without repeating ``--host``. +- an explicit ``--host`` always wins over both. + +``managed_executor_binding`` turns the selection plus the operator environment +into the readback a caller can act on before a Turn runs: which executor the +plan would use, how that executor is billed and bounded, whether it can launch +here, and -- when it cannot -- one typed reason naming the missing fact. The +Turn executor fails closed on that verdict, so a bounded Turn never drifts onto +a host the operator did not select. +""" + +from __future__ import annotations + +import importlib.util +from collections.abc import Mapping +from typing import Any, Callable + +from ..operator_credential import ( + OPERATOR_ENDPOINT_ENV_VAR, + configured_operator_credential, + env_text, +) + +# Explicit selection surfaces. The default is a product decision recorded here +# once; nothing in this module reads the environment to decide *which* host runs. +MANAGED_TURN_HOST = "dsh" +INDIVIDUAL_TURN_HOST = "codex-cli" +MANAGED_DEFAULT_TURN_HOST = MANAGED_TURN_HOST +TURN_HOST_ENV_VAR = "LOOPX_TURN_HOST" +TURN_HOST_SOURCE_PRODUCT_DEFAULT = "product_default" +TURN_HOST_SOURCE_EXPLICIT_CONFIG = "explicit_config" + +MANAGED_EXECUTOR_BINDING_SCHEMA_VERSION = "managed_executor_binding_v0" +# Executor kinds name where a Turn's model work is billed and bounded rather +# than which adapter is launched: a managed executor runs on an +# operator-supplied credential, an individual executor on one person's own CLI +# login, and a generic executor on a caller-supplied adapter command. +EXECUTOR_KIND_MANAGED = "managed" +EXECUTOR_KIND_INDIVIDUAL = "individual" +EXECUTOR_KIND_GENERIC = "generic" +INDIVIDUAL_CLI_HOSTS = frozenset({INDIVIDUAL_TURN_HOST, "claude-code"}) +MANAGED_HOST = MANAGED_TURN_HOST + +# The built-in dsh host launches the DeepSeek Harness runtime unless the caller +# supplies the explicit runner hook, so that module being importable is the +# launchability fact this projection checks without side effects. +DSH_RUNTIME_MODULE = "deepseek_harness" +DSH_RUNTIME_UNAVAILABLE = "dsh_runtime_unavailable" +# A managed host is billed to the operator's own endpoint. Without the operator +# credential (or an explicit injected runner) LoopX cannot authenticate that +# endpoint, so it refuses instead of letting the managed default consume +# whatever personal login happens to exist on the machine. +OPERATOR_CREDENTIAL_UNCONFIGURED = "operator_credential_unconfigured" + + +def selected_turn_host( + environ: Mapping[str, str] | None = None, +) -> tuple[str, str]: + """Return the selected default Turn host and the source that selected it. + + Selection is environment-independent: the shipped product default applies + until the operator re-points it explicitly with ``LOOPX_TURN_HOST``. A + configured credential is never a selection signal. + """ + + explicit = env_text(TURN_HOST_ENV_VAR, environ) + if explicit: + return explicit, TURN_HOST_SOURCE_EXPLICIT_CONFIG + return MANAGED_DEFAULT_TURN_HOST, TURN_HOST_SOURCE_PRODUCT_DEFAULT + + +def resolve_default_turn_host(environ: Mapping[str, str] | None = None) -> str: + """Return the selected default Turn host.""" + + return selected_turn_host(environ)[0] + + +def _configured_env_name(name: str, environ: Mapping[str, str] | None) -> str | None: + return env_text(name, environ) and name + + +def dsh_runtime_importable( + module_probe: Callable[[str], bool] | None = None, +) -> bool: + """Whether the DeepSeek Harness runtime the built-in dsh host launches exists.""" + + if module_probe is not None: + return bool(module_probe(DSH_RUNTIME_MODULE)) + try: + return importlib.util.find_spec(DSH_RUNTIME_MODULE) is not None + except (ImportError, ValueError): + return False + + +def managed_executor_binding( + host: str, + *, + environ: Mapping[str, str] | None = None, + dsh_runner_configured: bool = False, + module_probe: Callable[[str], bool] | None = None, +) -> dict[str, Any]: + """Project the executor one planned Turn would run on. + + ``available`` is ``False`` only when LoopX can prove the planned executor + cannot launch here, which is what a caller has to fail closed on. ``None`` + records that this projection does not probe that executor kind, so it makes + no claim rather than an unproven ``True``. + """ + + if host == MANAGED_HOST: + credential_env = configured_operator_credential(environ) + runtime_available = bool( + dsh_runner_configured or dsh_runtime_importable(module_probe) + ) + operator_credential_bound = bool(credential_env or dsh_runner_configured) + if not runtime_available: + unavailable_reason: str | None = DSH_RUNTIME_UNAVAILABLE + elif not operator_credential_bound: + unavailable_reason = OPERATOR_CREDENTIAL_UNCONFIGURED + else: + unavailable_reason = None + return { + "schema_version": MANAGED_EXECUTOR_BINDING_SCHEMA_VERSION, + "executor": host, + "executor_kind": EXECUTOR_KIND_MANAGED, + "credential_env": credential_env, + "endpoint_env": _configured_env_name(OPERATOR_ENDPOINT_ENV_VAR, environ), + # Billing boundary, stated instead of assumed: a managed executor is + # operator-credential-bound only when the credential or an explicit + # runner hook is configured here. + "operator_credential_bound": operator_credential_bound, + "available": unavailable_reason is None, + "unavailable_reason": unavailable_reason, + } + return { + "schema_version": MANAGED_EXECUTOR_BINDING_SCHEMA_VERSION, + "executor": host, + "executor_kind": ( + EXECUTOR_KIND_INDIVIDUAL + if host in INDIVIDUAL_CLI_HOSTS + else EXECUTOR_KIND_GENERIC + ), + "credential_env": None, + "endpoint_env": None, + "operator_credential_bound": False, + "available": None, + "unavailable_reason": None, + } + + +def managed_executor_payload_entry(plan: Mapping[str, Any]) -> dict[str, Any]: + """Return the execution payload's ``managed_executor`` entry, when planned. + + Reading the entry from the plan keeps one authority for the executor + identity: the payload quotes the binding the plan resolved instead of + re-deriving an executor from the launched host. + """ + + binding = plan.get("managed_executor") + return {"managed_executor": dict(binding)} if isinstance(binding, Mapping) else {} + + +def managed_executor_unavailable_payload( + plan: Mapping[str, Any], + *, + execute: bool, + host_projection: Mapping[str, Any], +) -> dict[str, Any] | None: + """Return the fail-closed execution payload for an unlaunchable executor. + + ``None`` means the plan makes no claim that its executor cannot launch, so + the Turn continues normally. A returned payload stops the Turn before the + journal, the host, and quota with no effect recorded, so a bounded Turn + cannot quietly move onto a different executor than the plan read back. + """ + + if not execute: + return None + binding = plan.get("managed_executor") + if not isinstance(binding, Mapping) or binding.get("available") is not False: + return None + return { + "status": "unavailable", + "host": dict(host_projection), + "reason": str(binding.get("unavailable_reason") or ""), + } diff --git a/tests/control_plane/test_cli_output_budget.py b/tests/control_plane/test_cli_output_budget.py index 9f65be4bb4..8c16665160 100644 --- a/tests/control_plane/test_cli_output_budget.py +++ b/tests/control_plane/test_cli_output_budget.py @@ -729,6 +729,8 @@ def _mode_variant_commands( + [ "turn", "run-once", + "--host", + "generic-cli", "--goal-id", GOAL_ID, "--agent-id", diff --git a/tests/control_plane/test_cli_output_differential.py b/tests/control_plane/test_cli_output_differential.py index 8cc785ca51..2a3a72b372 100644 --- a/tests/control_plane/test_cli_output_differential.py +++ b/tests/control_plane/test_cli_output_differential.py @@ -60,8 +60,15 @@ def _receipt(*rows: dict[str, object]) -> dict[str, object]: def test_thin_bilingual_byte_allowance_does_not_relax_character_or_quota_limits(): from loopx.control_plane.testing.cli_output_differential import _compare_row - base = _row(row_id="surface/heartbeat_prompt_thin/small/markdown", format="markdown", - chars=100, utf8_bytes=100, lines=1, compact_payload_chars=100) + + base = _row( + row_id="surface/heartbeat_prompt_thin/small/markdown", + format="markdown", + chars=100, + utf8_bytes=100, + lines=1, + compact_payload_chars=100, + ) candidate = {**base, "chars": 120, "utf8_bytes": 260} assert not _compare_row(base, candidate)["failures"] assert _compare_row(base, {**candidate, "chars": 133})["failures"] @@ -87,20 +94,42 @@ def test_sync_commit_uses_main_as_cli_output_base() -> None: @pytest.mark.parametrize("row_kind", ["surface", "variant"]) @pytest.mark.parametrize("mode", ["thin", "brief", "compact"]) -def test_host_safety_restoration_budget_is_one_time_bounded_and_prompt_only(row_kind, mode): +def test_host_safety_restoration_budget_is_one_time_bounded_and_prompt_only( + row_kind, mode +): from loopx.control_plane.testing.cli_output_differential import _compare_row - from loopx.control_plane.testing.cli_output_semantics import host_prompt_static_safety_revision + from loopx.control_plane.testing.cli_output_semantics import ( + host_prompt_static_safety_revision, + ) from loopx.control_plane.heartbeat.rules import HOST_LOOP_SAFETY_RULE - assert host_prompt_static_safety_revision(HOST_LOOP_SAFETY_RULE) == "host_prompt_static_safety_v1" - assert host_prompt_static_safety_revision(HOST_LOOP_SAFETY_RULE.replace("requires explicit authorization", "is always allowed")) is None + + assert ( + host_prompt_static_safety_revision(HOST_LOOP_SAFETY_RULE) + == "host_prompt_static_safety_v1" + ) + assert ( + host_prompt_static_safety_revision( + HOST_LOOP_SAFETY_RULE.replace( + "requires explicit authorization", "is always allowed" + ) + ) + is None + ) base = _row(row_id=f"{row_kind}/heartbeat_prompt_{mode}/small/json") - current = {**base, "chars": base["chars"] + 500, - "host_prompt_static_safety_revision": "host_prompt_static_safety_v1"} + current = { + **base, + "chars": base["chars"] + 500, + "host_prompt_static_safety_revision": "host_prompt_static_safety_v1", + } assert not _compare_row(base, current)["failures"] assert _compare_row(base, {**current, "chars": base["chars"] + 513})["failures"] - assert _compare_row(current, {**current, "chars": current["chars"] + 500})["failures"] - assert _compare_row({**base, "row_id": "surface/status/small/json"}, - {**current, "row_id": "surface/status/small/json"})["failures"] + assert _compare_row(current, {**current, "chars": current["chars"] + 500})[ + "failures" + ] + assert _compare_row( + {**base, "row_id": "surface/status/small/json"}, + {**current, "row_id": "surface/status/small/json"}, + )["failures"] @pytest.mark.parametrize("row_kind", ["surface", "variant"]) @@ -122,28 +151,75 @@ def test_reward_memory_outcome_prompt_budget_is_one_time_bounded_and_prompt_only reward_memory_outcome_prompt_revision(full_contract) == "reward_memory_outcome_prompt_v1" ) - assert reward_memory_outcome_prompt_revision( - full_contract.replace("zero provider calls", "best effort") - ) is None + assert ( + reward_memory_outcome_prompt_revision( + full_contract.replace("zero provider calls", "best effort") + ) + is None + ) base = _row(row_id=f"{row_kind}/heartbeat_prompt_{mode}/small/json") current = { **base, "chars": base["chars"] + 640, - "reward_memory_outcome_prompt_revision": ( - "reward_memory_outcome_prompt_v1" - ), + "reward_memory_outcome_prompt_revision": ("reward_memory_outcome_prompt_v1"), } assert not _compare_row(base, current)["failures"] - assert _compare_row(base, {**current, "chars": base["chars"] + 641})[ - "failures" - ] + assert _compare_row(base, {**current, "chars": base["chars"] + 641})["failures"] assert _compare_row(current, {**current, "chars": current["chars"] + 205})[ "failures" ] other = {**base, "row_id": "surface/status/small/json"} - assert _compare_row(other, {**current, "row_id": other["row_id"]})[ + assert _compare_row(other, {**current, "row_id": other["row_id"]})["failures"] + + +def test_managed_executor_binding_budget_is_one_time_bounded_and_turn_only() -> None: + from loopx.control_plane.testing.cli_output_differential import ( + _TURN_HOST_AND_MANAGED_EXECUTOR_BINDING_V0_GROWTH_ALLOWANCE as ALLOWANCE, + _compare_row, + ) + from loopx.control_plane.testing.cli_output_semantics import ( + managed_executor_binding_revision, + ) + + projection = ( + '{\n "managed_executor": {\n' + ' "executor_kind": "managed",\n' + ' "available": true,\n' + ' "unavailable_reason": null\n }\n}' + ) + assert ( + managed_executor_binding_revision(projection) == "managed_executor_binding_v0" + ) + assert ( + managed_executor_binding_revision( + projection.replace('"unavailable_reason"', '"reason"') + ) + is None + ) + + base = _row(row_id="variant/loopx_turn_run_once_preview/small/json") + current = { + **base, + "chars": base["chars"] + 254, + "utf8_bytes": base["utf8_bytes"] + 254, + "lines": base["lines"] + 9, + "compact_payload_chars": base["compact_payload_chars"] + 205, + "managed_executor_binding_revision": "managed_executor_binding_v0", + } + # The reviewed allowance covers the binding block plus the host-selection + # consequences on the Turn surfaces, and nothing beyond it. + assert not _compare_row(base, current)["failures"] + assert _compare_row( + base, {**current, "chars": base["chars"] + ALLOWANCE["chars"] + 1} + )["failures"] + # One time: the same growth against a baseline that already carries v0 is a + # regression, not a migration. + assert _compare_row(current, {**current, "chars": current["chars"] + 254})[ "failures" ] + # Surface scoped: the allowance never reaches a non-Turn surface. + other = {**base, "row_id": "surface/status/small/json"} + assert _compare_row(other, {**current, "row_id": other["row_id"]})["failures"] def test_regular_integration_pr_keeps_requested_cli_output_base() -> None: @@ -193,21 +269,15 @@ def payload(runtime_hash: str, source_hash: str) -> str: { "action_portfolio": {"schema_version": "quota_action_portfolio_v0"}, "nested": { - "action_portfolio": { - "schema_version": "quota_action_portfolio_v0" - } + "action_portfolio": {"schema_version": "quota_action_portfolio_v0"} }, } ) == ["quota_action_portfolio_v0"] assert planning_horizon_schema_versions( { - "planning_horizon": { - "schema_version": "quota_planning_horizon_v0" - }, + "planning_horizon": {"schema_version": "quota_planning_horizon_v0"}, "nested": { - "planning_horizon": { - "schema_version": "quota_planning_horizon_v0" - } + "planning_horizon": {"schema_version": "quota_planning_horizon_v0"} }, } ) == ["quota_planning_horizon_v0"] @@ -224,15 +294,7 @@ def payload(runtime_hash: str, source_hash: str) -> str: } ) == ["todo_planning_inventory_detail_v0"] assert guided_todo_delta_schema_versions( - { - "steps": [ - { - "todo_delta": { - "schema_version": "loopx_guided_todo_delta_v0" - } - } - ] - } + {"steps": [{"todo_delta": {"schema_version": "loopx_guided_todo_delta_v0"}}]} ) == ["loopx_guided_todo_delta_v0"] with_observability_field = json.loads(payload("third-runtime", "third-source")) @@ -423,9 +485,7 @@ def test_action_portfolio_migration_still_fails_above_bounded_growth() -> None: result = compare_cli_output_receipts(_receipt(_row()), _receipt(candidate)) assert result["ok"] is False - assert "chars grew by 1601; allowance is 1600" in ( - result["rows"][0]["failures"] - ) + assert "chars grew by 1601; allowance is 1600" in (result["rows"][0]["failures"]) def test_quota_action_portfolio_schema_migration_has_same_bounded_budget() -> None: @@ -512,9 +572,7 @@ def test_guided_todo_delta_migration_still_fails_above_bounded_growth() -> None: result = compare_cli_output_receipts(_receipt(_row()), _receipt(candidate)) assert result["ok"] is False - assert "chars grew by 513; allowance is 512" in ( - result["rows"][0]["failures"] - ) + assert "chars grew by 513; allowance is 512" in (result["rows"][0]["failures"]) def test_unknown_guided_todo_delta_schema_migration_fails_closed() -> None: @@ -538,9 +596,7 @@ def test_unknown_action_portfolio_schema_migration_fails_closed() -> None: result = compare_cli_output_receipts(_receipt(_row()), _receipt(candidate)) assert result["ok"] is False - assert result["rows"][0]["failures"] == [ - "action_portfolio schema coverage changed" - ] + assert result["rows"][0]["failures"] == ["action_portfolio schema coverage changed"] def test_unknown_action_signature_coverage_migration_fails_closed() -> None: @@ -552,9 +608,7 @@ def test_unknown_action_signature_coverage_migration_fails_closed() -> None: result = compare_cli_output_receipts(_receipt(_row()), _receipt(candidate)) assert result["ok"] is False - assert result["rows"][0]["failures"] == [ - "action_signature semantic digest changed" - ] + assert result["rows"][0]["failures"] == ["action_signature semantic digest changed"] def test_planning_horizon_v0_migration_has_one_bounded_growth_budget() -> None: @@ -590,9 +644,7 @@ def test_planning_horizon_v0_migration_fails_above_its_bounded_growth() -> None: result = compare_cli_output_receipts(_receipt(_row()), _receipt(candidate)) assert result["ok"] is False - assert "chars grew by 3201; allowance is 3200" in ( - result["rows"][0]["failures"] - ) + assert "chars grew by 3201; allowance is 3200" in (result["rows"][0]["failures"]) def test_unknown_planning_horizon_schema_migration_fails_closed() -> None: @@ -603,16 +655,12 @@ def test_unknown_planning_horizon_schema_migration_fails_closed() -> None: result = compare_cli_output_receipts(_receipt(_row()), _receipt(candidate)) assert result["ok"] is False - assert result["rows"][0]["failures"] == [ - "planning_horizon schema coverage changed" - ] + assert result["rows"][0]["failures"] == ["planning_horizon schema coverage changed"] def test_planning_inventory_detail_v0_has_one_bounded_growth_budget() -> None: candidate = _row( - planning_inventory_detail_schema_versions=[ - "todo_planning_inventory_detail_v0" - ], + planning_inventory_detail_schema_versions=["todo_planning_inventory_detail_v0"], chars=41_280, utf8_bytes=41_280, lines=1_036, @@ -631,15 +679,11 @@ def test_planning_inventory_detail_v0_has_one_bounded_growth_budget() -> None: def test_planning_inventory_detail_migration_is_bounded_and_fail_closed() -> None: oversized = _row( - planning_inventory_detail_schema_versions=[ - "todo_planning_inventory_detail_v0" - ], + planning_inventory_detail_schema_versions=["todo_planning_inventory_detail_v0"], chars=41_281, ) unknown = _row( - planning_inventory_detail_schema_versions=[ - "todo_planning_inventory_detail_v1" - ] + planning_inventory_detail_schema_versions=["todo_planning_inventory_detail_v1"] ) oversized_result = compare_cli_output_receipts( @@ -652,8 +696,9 @@ def test_planning_inventory_detail_migration_is_bounded_and_fail_closed() -> Non ) assert oversized_result["ok"] is False - assert "chars grew by 1281; allowance is 1280" in ( - oversized_result["rows"][0]["failures"] + assert ( + "chars grew by 1281; allowance is 1280" + in (oversized_result["rows"][0]["failures"]) ) assert unknown_result["rows"][0]["failures"] == [ "planning inventory detail schema coverage changed" @@ -749,7 +794,7 @@ def test_runtime_root_route_count_only_matches_executable_command_prefixes() -> text = ( " loopx --runtime-root /tmp/indented refresh-state\n" "loopx --runtime-root /tmp/runtime refresh-state\n" - "{\"command\": \"loopx --runtime-root '/tmp/runtime root' quota spend-slot\"}\n" + '{"command": "loopx --runtime-root \'/tmp/runtime root\' quota spend-slot"}\n' "- expanded: `loopx --runtime-root /tmp/runtime heartbeat-prompt`\n" "Use --runtime-root PATH to select a runtime.\n" "The command is loopx --runtime-root /tmp/runtime.\n" @@ -850,26 +895,39 @@ def test_fixture_contract_mismatch_fails_closed() -> None: @pytest.mark.parametrize("previous", range(4)) def test_agent_context_v4_migration_is_bounded_and_one_time(previous): - base = _row(action_signature_coverages=[f"turn_envelope_action_dimensions_v{previous}"]) - candidate = {**base, "action_signature_sha256": "agent-context-signature", - "action_signature_coverages": ["turn_envelope_action_dimensions_v4"], - "chars": 42_048, "utf8_bytes": 42_048, "lines": 1_048, - "compact_payload_chars": 21_664} + base = _row( + action_signature_coverages=[f"turn_envelope_action_dimensions_v{previous}"] + ) + candidate = { + **base, + "action_signature_sha256": "agent-context-signature", + "action_signature_coverages": ["turn_envelope_action_dimensions_v4"], + "chars": 42_048, + "utf8_bytes": 42_048, + "lines": 1_048, + "compact_payload_chars": 21_664, + } result = compare_cli_output_receipts(_receipt(base), _receipt(candidate)) assert result["ok"] and result["review_required"] assert result["rows"][0]["review_signals"] == [ f"action_signature coverage migrated: turn_envelope_action_dimensions_v{previous}" - " -> turn_envelope_action_dimensions_v4"] + " -> turn_envelope_action_dimensions_v4" + ] for metric in ("chars", "utf8_bytes", "lines", "compact_payload_chars"): too_large = {**candidate, metric: candidate[metric] + 1} - assert not compare_cli_output_receipts(_receipt(base), _receipt(too_large))["ok"] + assert not compare_cli_output_receipts(_receipt(base), _receipt(too_large))[ + "ok" + ] # After migration, neither growing again nor changing semantics is excused. grown = {**candidate, "chars": candidate["chars"] + 2_048} assert not compare_cli_output_receipts(_receipt(candidate), _receipt(grown))["ok"] changed = {**candidate, "action_signature_sha256": "unexpected-semantic-change"} assert not compare_cli_output_receipts(_receipt(candidate), _receipt(changed))["ok"] - reverse = {**candidate, "action_signature_coverages": base["action_signature_coverages"], - "action_signature_sha256": "reverse-signature"} + reverse = { + **candidate, + "action_signature_coverages": base["action_signature_coverages"], + "action_signature_sha256": "reverse-signature", + } assert not compare_cli_output_receipts(_receipt(candidate), _receipt(reverse))["ok"] @@ -879,10 +937,16 @@ def test_public_multi_subagent_probe_reaches_v4_producer(tmp_path): from tests.control_plane import test_cli_output_budget as probe from loopx.control_plane.testing import cli_output_semantics as semantics - runner = runpy.run_path(str(Path(__file__).resolve().parents[2] - / 'examples/control_plane/cli-output-probe-runner.py')) + runner = runpy.run_path( + str( + Path(__file__).resolve().parents[2] + / "examples/control_plane/cli-output-probe-runner.py" + ) + ) with probe._stable_budget_fixture_root(tmp_path) as root: - rows = runner['_multi_subagent_rows'](probe, semantics, root) + rows = runner["_multi_subagent_rows"](probe, semantics, root) assert len(rows) == 1 - assert rows[0]['action_signature_coverages'] == ['turn_envelope_action_dimensions_v4'] - assert any('agent_context' in path for path in rows[0]['json_shape_paths']) + assert rows[0]["action_signature_coverages"] == [ + "turn_envelope_action_dimensions_v4" + ] + assert any("agent_context" in path for path in rows[0]["json_shape_paths"]) diff --git a/tests/test_loopx_turn_driver.py b/tests/test_loopx_turn_driver.py index a55c52b954..bd98280dc2 100644 --- a/tests/test_loopx_turn_driver.py +++ b/tests/test_loopx_turn_driver.py @@ -1713,6 +1713,8 @@ def test_turn_run_once_cli_commits_validated_result_and_one_quota_slot( "json", "turn", "run-once", + "--host", + "generic-cli", "--goal-id", "loopx-turn-fixture", "--agent-id", @@ -1782,6 +1784,8 @@ def test_turn_run_once_cli_commits_validated_result_and_one_quota_slot( "json", "turn", "run-once", + "--host", + "generic-cli", "--goal-id", "loopx-turn-fixture", "--agent-id", @@ -1954,6 +1958,8 @@ def test_turn_run_once_cli_completes_selected_todo_after_validation( "json", "turn", "run-once", + "--host", + "generic-cli", "--goal-id", "loopx-turn-fixture", "--agent-id", @@ -2195,6 +2201,8 @@ def _turn_run_once_completion_argv( "json", "turn", "run-once", + "--host", + "generic-cli", "--goal-id", "loopx-turn-fixture", "--agent-id", @@ -2811,6 +2819,8 @@ def test_turn_run_once_cli_rejects_unproven_host_claim_before_writeback( "json", "turn", "run-once", + "--host", + "generic-cli", "--goal-id", "loopx-turn-fixture", "--agent-id", diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index 2bc97abaa5..5191f8b791 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -24,7 +24,9 @@ BuiltInHostError, LOOPX_TURN_JOURNAL_SCHEMA_VERSION, _task_validation_stage, + turn_journal_path, ) +from loopx.control_plane.turn_driver.host_binding import managed_executor_binding from loopx.control_plane.turn_driver.settlement import execute_turn_driver_settlement from loopx.control_plane.turn_driver.transaction import TRANSACTION_PHASES @@ -77,6 +79,25 @@ def _codex_plan() -> dict[str, object]: ) +def _managed_plan(*, runtime_available: bool) -> dict[str, object]: + """One dsh plan carrying the executor readback the command layer attaches.""" + + plan = _plan() + envelope = plan["turn_envelope"] + assert isinstance(envelope, dict) + managed = build_loopx_turn_plan( + envelope, + host="dsh", + execution_mode="isolated-headless", + ) + managed["managed_executor"] = managed_executor_binding( + "dsh", + environ={"DEEPSEEK_API_KEY": "fixture-operator-credential"}, + module_probe=lambda _module: runtime_available, + ) + return managed + + def _adaptive_observation_plan( *, required_write_scopes: list[str] | None = None, @@ -2501,3 +2522,74 @@ def scheduler(_spend: dict[str, object]) -> dict[str, object]: turn_key=str(transaction["turn_key"]), ) assert audited["last_recovery"] == resumed["recovery"] + + +def test_run_once_fails_closed_when_the_managed_executor_cannot_launch(tmp_path): + plan = _managed_plan(runtime_available=False) + transaction = plan["transaction"] + assert isinstance(transaction, dict) + runtime_root = tmp_path / "runtime" + journal = turn_journal_path( + runtime_root, + goal_id="fixture-goal", + turn_key=str(transaction["turn_key"]), + ) + + payload = run_loopx_turn_once( + plan, + host_runner=lambda _request: pytest.fail("an unavailable executor must not run"), + project=tmp_path, + runtime_root=runtime_root, + goal_id="fixture-goal", + timeout_seconds=5, + execute=True, + ) + + assert payload["ok"] is False + assert payload["status"] == "unavailable" + assert payload["reason"] == "dsh_runtime_unavailable" + assert payload["effects"] == { + "host_invoked": False, + "state_written": False, + "quota_spent": False, + "scheduler_acknowledged": False, + } + assert payload["quota_slot_spend_count"] == 0 + assert payload["managed_executor"] == plan["managed_executor"] + assert journal.exists() is False + + +def test_run_once_preview_reports_the_managed_executor_without_refusing(tmp_path): + plan = _managed_plan(runtime_available=False) + + payload = run_loopx_turn_once( + plan, + host_runner=lambda _request: pytest.fail("preview must not run the host"), + project=tmp_path, + runtime_root=tmp_path / "runtime", + goal_id="fixture-goal", + timeout_seconds=5, + execute=False, + ) + + assert payload["ok"] is True + assert payload["status"] == "preview" + assert payload["managed_executor"]["available"] is False + assert payload["managed_executor"]["unavailable_reason"] == "dsh_runtime_unavailable" + + +def test_run_once_does_not_refuse_a_launchable_managed_executor(tmp_path): + plan = _managed_plan(runtime_available=True) + + # The refusal is the only guard under test here: without writeback, spend, + # and scheduler callbacks the executor stops at its own contract instead. + with pytest.raises(ValueError, match="requires writeback, spend, and scheduler"): + run_loopx_turn_once( + plan, + host_runner=lambda _request: pytest.fail("host must not run without callbacks"), + project=tmp_path, + runtime_root=tmp_path / "runtime", + goal_id="fixture-goal", + timeout_seconds=5, + execute=True, + ) diff --git a/tests/test_turn_default_host_binding.py b/tests/test_turn_default_host_binding.py new file mode 100644 index 0000000000..115e9efeea --- /dev/null +++ b/tests/test_turn_default_host_binding.py @@ -0,0 +1,122 @@ +"""The Turn host is selected explicitly; a credential only authenticates it.""" + +from __future__ import annotations + +import pytest + +from loopx.cli import build_parser +from loopx.control_plane.operator_credential import configured_operator_credential +from loopx.control_plane.turn_driver.host_binding import ( + MANAGED_DEFAULT_TURN_HOST, + MANAGED_TURN_HOST, + TURN_HOST_ENV_VAR, + TURN_HOST_SOURCE_EXPLICIT_CONFIG, + TURN_HOST_SOURCE_PRODUCT_DEFAULT, + resolve_default_turn_host, + selected_turn_host, +) + + +def test_default_host_is_the_managed_product_default(): + assert MANAGED_DEFAULT_TURN_HOST == MANAGED_TURN_HOST == "dsh" + assert resolve_default_turn_host({}) == MANAGED_DEFAULT_TURN_HOST + assert selected_turn_host({}) == ( + MANAGED_DEFAULT_TURN_HOST, + TURN_HOST_SOURCE_PRODUCT_DEFAULT, + ) + + +@pytest.mark.parametrize( + "environ", + [ + {"DEEPSEEK_API_KEY": "sk-operator"}, + {"DEEPSEEK_API_KEY": ""}, + {"DEEPSEEK_API_KEY": " "}, + {"DEEPSEEK_BASE_URL": "https://example.invalid"}, + {"DEEPSEEK_API_KEY": "sk-operator", "DEEPSEEK_BASE_URL": "https://x.invalid"}, + ], +) +def test_a_credential_never_changes_the_selected_host(environ): + """Discovering a credential must not re-point a Turn by itself.""" + + assert resolve_default_turn_host(environ) == MANAGED_DEFAULT_TURN_HOST + assert selected_turn_host(environ)[1] == TURN_HOST_SOURCE_PRODUCT_DEFAULT + + +def test_explicit_config_repoints_the_default_host(): + environ = {TURN_HOST_ENV_VAR: "codex-cli", "DEEPSEEK_API_KEY": "sk-operator"} + + assert selected_turn_host(environ) == ( + "codex-cli", + TURN_HOST_SOURCE_EXPLICIT_CONFIG, + ) + assert resolve_default_turn_host(environ) == "codex-cli" + + +def test_configured_credential_names_the_env_var(): + assert ( + configured_operator_credential({"DEEPSEEK_API_KEY": "sk-operator"}) + == "DEEPSEEK_API_KEY" + ) + assert configured_operator_credential({}) is None + + +def _turn_argv(command: str) -> list[str]: + argv = ["turn", command, "--goal-id", "goal-x", "--agent-id", "agent-x"] + if command == "run-once": + argv.extend(["--project", "."]) + return argv + + +@pytest.mark.parametrize("command", ["plan", "run-once"]) +@pytest.mark.parametrize( + "environ", + [{}, {"DEEPSEEK_API_KEY": "sk-operator"}, {"DEEPSEEK_API_KEY": " "}], +) +def test_cli_defaults_to_the_selected_host_regardless_of_credentials( + command, environ, monkeypatch +): + for name in ("DEEPSEEK_API_KEY", TURN_HOST_ENV_VAR): + monkeypatch.delenv(name, raising=False) + for name, value in environ.items(): + monkeypatch.setenv(name, value) + + assert ( + build_parser().parse_args(_turn_argv(command)).host == MANAGED_DEFAULT_TURN_HOST + ) + + +@pytest.mark.parametrize("command", ["plan", "run-once"]) +def test_explicit_config_environment_repoints_the_cli_default(command, monkeypatch): + monkeypatch.setenv(TURN_HOST_ENV_VAR, "codex-cli") + + assert build_parser().parse_args(_turn_argv(command)).host == "codex-cli" + + +def test_explicit_host_flag_wins_over_the_default(monkeypatch): + monkeypatch.setenv("DEEPSEEK_API_KEY", "sk-operator") + monkeypatch.setenv(TURN_HOST_ENV_VAR, "dsh") + + args = build_parser().parse_args([*_turn_argv("run-once"), "--host", "generic-cli"]) + + assert args.host == "generic-cli" + + +@pytest.mark.parametrize("command", ["plan", "run-once"]) +def test_default_execution_mode_follows_the_selected_host(command, monkeypatch): + for name in ("DEEPSEEK_API_KEY", TURN_HOST_ENV_VAR): + monkeypatch.delenv(name, raising=False) + managed = build_parser().parse_args(_turn_argv(command)) + + monkeypatch.setenv(TURN_HOST_ENV_VAR, "codex-cli") + individual = build_parser().parse_args(_turn_argv(command)) + + # The selected managed host runs bounded headless Turns; pairing it with a + # visible interactive mode would make the shipped default unschedulable. + # run-once ships only the isolated-headless mode, so it keeps that either way. + assert managed.host == MANAGED_DEFAULT_TURN_HOST + assert managed.execution_mode == "isolated-headless" + assert individual.host == "codex-cli" + assert individual.execution_mode == ( + "interactive-visible" if command == "plan" else "isolated-headless" + ) diff --git a/tests/test_turn_managed_executor_binding.py b/tests/test_turn_managed_executor_binding.py new file mode 100644 index 0000000000..fddd1b9e99 --- /dev/null +++ b/tests/test_turn_managed_executor_binding.py @@ -0,0 +1,178 @@ +"""The managed executor readback names the executor and whether it can launch.""" + +from __future__ import annotations + +import pytest + +from loopx.control_plane.turn_driver.host_binding import ( + DSH_RUNTIME_UNAVAILABLE, + EXECUTOR_KIND_GENERIC, + EXECUTOR_KIND_INDIVIDUAL, + EXECUTOR_KIND_MANAGED, + MANAGED_EXECUTOR_BINDING_SCHEMA_VERSION, + MANAGED_TURN_HOST, + OPERATOR_CREDENTIAL_UNCONFIGURED, + managed_executor_binding, + resolve_default_turn_host, +) + +_NO_RUNTIME = lambda _module: False # noqa: E731 - tiny probe fixture +_RUNTIME = lambda _module: True # noqa: E731 - tiny probe fixture + + +def test_managed_executor_reports_the_operator_credential_and_endpoint(): + binding = managed_executor_binding( + "dsh", + environ={ + "DEEPSEEK_API_KEY": "sk-operator", + "DEEPSEEK_BASE_URL": "https://example.invalid", + }, + module_probe=_RUNTIME, + ) + + assert binding == { + "schema_version": MANAGED_EXECUTOR_BINDING_SCHEMA_VERSION, + "executor": "dsh", + "executor_kind": EXECUTOR_KIND_MANAGED, + "credential_env": "DEEPSEEK_API_KEY", + "endpoint_env": "DEEPSEEK_BASE_URL", + "operator_credential_bound": True, + "available": True, + "unavailable_reason": None, + } + + +def test_managed_executor_fails_closed_when_the_runtime_is_missing(): + binding = managed_executor_binding( + "dsh", + environ={"DEEPSEEK_API_KEY": "sk-operator"}, + module_probe=_NO_RUNTIME, + ) + + assert binding["available"] is False + assert binding["unavailable_reason"] == DSH_RUNTIME_UNAVAILABLE + + +def test_configured_runner_hook_makes_the_managed_host_launchable(): + binding = managed_executor_binding( + "dsh", + environ={"DEEPSEEK_API_KEY": "sk-operator"}, + dsh_runner_configured=True, + module_probe=_NO_RUNTIME, + ) + + assert binding["available"] is True + assert binding["unavailable_reason"] is None + + +def test_managed_executor_reports_an_unconfigured_credential_without_inventing_one(): + binding = managed_executor_binding( + "dsh", + environ={}, + dsh_runner_configured=True, + module_probe=_RUNTIME, + ) + + assert binding["credential_env"] is None + assert binding["endpoint_env"] is None + + +def test_managed_selection_without_the_operator_credential_fails_closed(): + """The managed host authenticates with the operator credential or not at all.""" + + binding = managed_executor_binding("dsh", environ={}, module_probe=_RUNTIME) + + assert binding["executor_kind"] == EXECUTOR_KIND_MANAGED + assert binding["operator_credential_bound"] is False + assert binding["available"] is False + assert binding["unavailable_reason"] == OPERATOR_CREDENTIAL_UNCONFIGURED + + +@pytest.mark.parametrize("blank", ["", " "]) +def test_blank_credential_counts_as_unconfigured(blank): + binding = managed_executor_binding( + "dsh", environ={"DEEPSEEK_API_KEY": blank}, module_probe=_RUNTIME + ) + + assert binding["credential_env"] is None + assert binding["available"] is False + assert binding["unavailable_reason"] == OPERATOR_CREDENTIAL_UNCONFIGURED + + +def test_configured_runner_hook_counts_as_an_operator_credential_boundary(): + binding = managed_executor_binding( + "dsh", + environ={}, + dsh_runner_configured=True, + module_probe=_NO_RUNTIME, + ) + + assert binding["operator_credential_bound"] is True + assert binding["available"] is True + + +def test_individual_and_generic_executors_are_not_operator_credential_bound(): + for host in ("codex-cli", "generic-cli"): + binding = managed_executor_binding(host, environ={"DEEPSEEK_API_KEY": "sk-x"}) + + assert binding["operator_credential_bound"] is False, binding + + +@pytest.mark.parametrize( + ("host", "expected_kind"), + [ + ("codex-cli", EXECUTOR_KIND_INDIVIDUAL), + ("claude-code", EXECUTOR_KIND_INDIVIDUAL), + ("generic-cli", EXECUTOR_KIND_GENERIC), + ], +) +def test_other_hosts_make_no_launch_claim_and_carry_no_operator_env( + host, expected_kind +): + binding = managed_executor_binding( + host, + environ={"DEEPSEEK_API_KEY": "sk-operator"}, + module_probe=_RUNTIME, + ) + + assert binding["executor_kind"] == expected_kind + assert binding["available"] is None + assert binding["unavailable_reason"] is None + assert binding["credential_env"] is None + assert binding["endpoint_env"] is None + assert binding["operator_credential_bound"] is False + + +@pytest.mark.parametrize( + "environ", + [ + {}, + {"DEEPSEEK_API_KEY": "sk-operator"}, + {"DEEPSEEK_API_KEY": ""}, + {"DEEPSEEK_API_KEY": " "}, + ], +) +def test_default_resolution_always_names_the_managed_executor(environ): + default_host = resolve_default_turn_host(environ) + binding = managed_executor_binding( + default_host, + environ=environ, + module_probe=_RUNTIME, + ) + + # The default host comes from the product default, not from the credential, + # so the readback always describes the managed executor. Whether it may run + # is a separate, explicitly projected fact. + assert default_host == MANAGED_TURN_HOST + assert binding["executor_kind"] == EXECUTOR_KIND_MANAGED + assert binding["available"] is ( + "DEEPSEEK_API_KEY" in environ and bool(environ["DEEPSEEK_API_KEY"].strip()) + ) + + +def test_endpoint_without_credential_is_reported_but_does_not_switch_host(): + environ = {"DEEPSEEK_BASE_URL": "https://example.invalid"} + binding = managed_executor_binding("codex-cli", environ=environ) + + assert resolve_default_turn_host(environ) == MANAGED_TURN_HOST + assert binding["endpoint_env"] is None