From 13daeee8749b0f9e01ab7e84389a8127d0c338eb Mon Sep 17 00:00:00 2001 From: witbrock Date: Sun, 13 Sep 2026 15:53:28 -0400 Subject: [PATCH 1/7] Add scoped prepared coding-agent instance lifecycle --- docs/design_index.md | 1 + docs/engineering/codex_dgx_worker.md | 5 + docs/engineering/coding_agent_instances.md | 157 ++++++++ scripts/codex_von_capacity.py | 53 +++ scripts/codex_von_inbox.py | 10 +- scripts/codex_von_inbox_prompt.md | 2 +- scripts/codex_von_instance.py | 283 ++++++++++++++ scripts/codex_von_instance_schedule.py | 104 ++++++ scripts/codex_von_worker.py | 55 ++- scripts/codex_von_worker_prompt.md | 9 +- .../integrations/internal_mcp/catalogue.py | 2 + .../coding_agent_instance_tools.py | 50 +++ .../services/coding_agent_instance_service.py | 211 +++++++++++ tests/backend/test_codex_von_worker.py | 23 +- tests/backend/test_coding_agent_instances.py | 349 ++++++++++++++++++ tests/test_codex_von_capacity.py | 66 ++++ 16 files changed, 1352 insertions(+), 28 deletions(-) create mode 100644 docs/engineering/coding_agent_instances.md create mode 100644 scripts/codex_von_capacity.py create mode 100644 scripts/codex_von_instance.py create mode 100644 scripts/codex_von_instance_schedule.py create mode 100644 src/backend/integrations/internal_mcp/coding_agent_instance_tools.py create mode 100644 src/backend/services/coding_agent_instance_service.py create mode 100644 tests/backend/test_coding_agent_instances.py create mode 100644 tests/test_codex_von_capacity.py diff --git a/docs/design_index.md b/docs/design_index.md index b4770a8a..267b91ee 100644 --- a/docs/design_index.md +++ b/docs/design_index.md @@ -83,6 +83,7 @@ live. | Build/deployment receipts and task work products | [Build and deployment evidence](engineering/deployment_evidence.md) | | Task product selection and continuing a responsibility from Tasks | [Task responsibility continuation](engineering/task_responsibility_continuation.md) | | Subscription-backed coding task pickup on the DGX | [Codex DGX worker pilot](engineering/codex_dgx_worker.md) | +| Provisioning prepared coding-agent instances | [Instance enrolment and lifecycle](engineering/coding_agent_instances.md) | | Coding-agent progress and completion messages in Von | [Interactive coding-agent messages](engineering/codex_dgx_worker.md#interactive-coding-agent-messages); operator-specific delivery configuration in personal `AGENTS.md` | | Enduring role competence, organisational continuity, active acquisition, learned reuse or transfer | [Role-learning convergence](engineering/role_learning_convergence.md), with the [staged development rehearsal](engineering/role_convergence_rehearsal.md); live Jira owns delivery status and the next integrating increment | | Material security exposure: authentication/authorisation, private or cross-namespace data, secrets, untrusted content with tool authority, effects outside ordinary bounded and recoverable delegation, deployment, or administrator surfaces | [Security considerations](engineering/security_considerations.md) | diff --git a/docs/engineering/codex_dgx_worker.md b/docs/engineering/codex_dgx_worker.md index 0d7a8489..006dd4fb 100644 --- a/docs/engineering/codex_dgx_worker.md +++ b/docs/engineering/codex_dgx_worker.md @@ -159,6 +159,11 @@ operator profile; it is not a hardened boundary against a hostile host user. ## Installation +For multiple prepared accounts/organisations, use the +[instance lifecycle](coding_agent_instances.md), including shared host capacity. +Do not duplicate an existing identity's poller. The legacy installation below +remains valid for an existing single consumer. + Use a dedicated Codex home with ChatGPT subscription login, an explicit non-Sol model, and a working native Linux sandbox. Keep the API-key Codex home separate. The launcher should unset `OPENAI_API_KEY` and `CODEX_API_KEY`, set diff --git a/docs/engineering/coding_agent_instances.md b/docs/engineering/coding_agent_instances.md new file mode 100644 index 00000000..da83daa0 --- /dev/null +++ b/docs/engineering/coding_agent_instances.md @@ -0,0 +1,157 @@ +# Prepared coding-agent instances + +This capability packages the existing Codex task/inbox worker for several +independently attributed instances on a prepared Linux host. It does not create +Unix accounts, enrol authentication, grant arbitrary host access, or certify +isolation between mutually untrusted workloads. Same-account workers share that +account's authority. + +The callable surface is `coding_agent_instances`: `list`, `status`, +`reconcile`, `pause`, and `resume`. A normal Von agent calls it under its +authenticated user's organisation context. The user must still have live +`MANAGE_MEMBERS` permission in that exact organisation. Payload actor IDs are +not accepted. Separate coding-controller identities do not automatically acquire +that user's provisioning authority. + +## Specification and host enrolment + +An operator enrols concrete slots in the private JSON registry selected by +`VON_CODING_AGENT_REGISTRY`. It must be root/service-owned and not group/world +writable. Keep it and all credential references outside Git and model context. +One registry covers all approved host consumers; enrol existing identities as +`reserved: true` so another slot cannot duplicate them. Identity duplication +within the registry is rejected. Independent registries or manually launched +unenrolled workers are outside this topology guarantee. + +Each `instances` entry has an exact stable `instance_id`, `agent_id`, +`agent_name`, `role`, `delegator_id`, `organisation_id`, `endpoint`, +`host`, `unix_account`, `allow_identity_creation`, and +`provisioner_command`. The command is an operator-approved argv, for example: + +```json +{ + "instances": { + "research-helper": { + "instance_id": "research-helper", + "agent_id": "#V#research_helper", + "agent_name": "Research Helper", + "role": "Coding tasks assigned by the enrolled researcher", + "delegator_id": "#V#researcher", + "organisation_id": "#V#research_lab", + "endpoint": "https://von.example.invalid", + "host": "prepared-host", + "unix_account": "research-agent", + "allow_identity_creation": true, + "provisioner_command": ["/operator/fixed-research-helper"] + } + } +} +``` + +The fixed executable invokes the reviewed `scripts/codex_von_instance.py` +with its operator-chosen `--binding` path under the prepared account. Its stdin +contains only an action and optional model/reasoning settings. Bind remote SSH +or cross-account execution to that one executable/configuration; do not grant a +general root shell or let a caller supply argv/account/path flags. Access to the +helper itself must be restricted to the trusted service/operator. The host +independently checks its hostname, account and selector configuration, and its +receipt must match the server's binding. + +The private host binding repeats the six identity/host/account fields and adds: + +- `instance_root`: a dedicated private directory, never an existing worker root; +- `release_commit`: full reviewed SHA, available from the approved repository; +- `worker_config`: the existing worker's config, including the matching + `agent_id`, `agent_name`, `delegator_id`, `organisation_id`, + `source_repo`, `codex_command`, `python_environment`, and + `state_root=/state`; +- existing task-source image delegation, repository authentication references, + sandbox/profile and permitted deployment route, only where authorised; +- `systemd`: `unit_directory` for this account's user units, + `memory_max_bytes`, optional `poll_interval_seconds` (default 300), and + `environment_files` containing approved **paths**, never secret values; +- `capacity_enrolment_verified` and `resource_limits_verified`: operator + attestations after checking enrolment and resource controls. These are + prerequisites, not independent evidence of live aggregate safety. + +Keep the canonical Von endpoint, actual database/environment binding, Unix +identity and Codex login separate. The operator must verify that the enrolled +environment reaches that endpoint's intended database before enabling work. +The existing worker uses canonical Python services, not an HTTP endpoint +parameter. Changing an endpoint string does not redirect that database. + +The normal installer writes deterministic `von-coding-INSTANCE.service/timer` +units and refuses to overwrite a unit without its own instance marker. It +installs `MemoryMax`, disables swap for the service and polls on the configured +cadence. Alternatively an existing scheduler can be retained with fixed +`schedule_commands.enable/disable/status` argv bindings; status must return +JSON with boolean `enabled`. That adapter owns its verified resource ceiling. +User lingering, prepared Python dependencies, a functioning sandbox, per-account +Codex subscription authentication and repository credential enrolment remain +operator prerequisites for a new account. + +## Normal use and recovery + +1. Call `coding_agent_instances(action="list")` to discover only slots for the + current authenticated delegator and organisation. +2. Call `reconcile` for an enrolled slot, optionally passing + `settings={"model":"gpt-6-astra","reasoning_effort":"high"}`. Canonical concept + and membership APIs create/reuse the preapproved identity; an unrelated + existing concept cannot be repurposed. Task preferences still override these + defaults. Tokens are preserved, not assumed to prove provider availability. +3. A new instance is installed **paused**. Reconcile retains an existing pause. + Call `resume` explicitly to enable that instance's timer. Call `pause` to + stop future polling without terminating an active run. +4. Inspect `status`: selected release/settings, schedule state and retained + runtime evidence are separate. A runtime receipt can be historical. Actual + readiness additionally requires canonical task pickup, outcome and intended + recipient read-back; provisioning never asserts those. + +Provisioning records its step before effects, serialises per-instance changes, +uses the existing coherent release installer, and holds the worker's existing +lock during reconfiguration. An active run returns `waiting_for_active_run`. +After an interrupted installation, reconcile the same slot. After uncertain +scheduling, inspect status then repeat the same pause/resume operation. It +targets the same unit, identity and state, not a new consumer. Do not remove +loaded releases or state to retry. + +For upgrade/rollback, the operator selects a reviewed release SHA in the binding, +pauses the timer, waits for the active worker lock, and reconciles the same +instance. Its task/inbox state survives. Verify canonical runtime and recipient +receipts before claiming the new version is serving work. + +## Aggregate capacity + +Every participating worker config must contain the same `host_capacity`: + +```json +{ + "lock_path": "/operator/host-capacity.lock", + "reserve_bytes": 17179869184, + "launch_headroom_bytes": 25769803776 +} +``` + +Those example memory values are not a universal recommendation. Calibrate them +against the host and concurrent non-coding services. The operator creates one +stable regular inode, not replaceable by consumers, readable across prepared +accounts and not group/world writable. All consumers, including the existing +worker, must be enrolled before enabling another one. The admission module +serialises coding work host-wide and checks Linux `MemAvailable` while holding +the lock. It reports `waiting_for_capacity` without selecting a new task. +Coding children inherit the descriptor so controller death does not release +admission while the child still runs. Per-service memory ceilings remain +necessary; this lock alone does not constrain non-coding processes. + +## Evidence boundary for the initial candidate + +Targeted tests exercise scoped discovery, forged actor rejection, live role +loss, bad settings, canonical identity/membership adapter calls, repeated and +interrupted reconciliation, retained pause, active-run protection, scheduler +ownership and two actual competing lock holders including an inherited child. +Existing worker/inbox tests cover task selection and attachment authority. +Mocked service/scheduler coverage is not a live Von deployment or another Unix +account. Before calling the package ready, run the task's isolated canary through +the normal agent-callable path and retain task, attachment and recipient +read-back. The implementation task's current delivery explicitly prohibits +starting another worker, so that acceptance step remains outstanding. diff --git a/scripts/codex_von_capacity.py b/scripts/codex_von_capacity.py new file mode 100644 index 00000000..8ae27bfb --- /dev/null +++ b/scripts/codex_von_capacity.py @@ -0,0 +1,53 @@ +"""One host-wide admission slot shared by prepared Unix accounts. + +The operator creates a stable, read-only-to-consumers lock inode. All enrolled +workers, including the existing worker, must use it. The lock is retained by +coding subprocesses so controller interruption cannot admit another launch. +This intentionally serialises coding work; it is not a fleet scheduler. +""" + +import fcntl +import os +import stat +from contextlib import contextmanager +from pathlib import Path + + +def available_memory_bytes(): + for line in Path("/proc/meminfo").read_text().splitlines(): + if line.startswith("MemAvailable:"): + return int(line.split()[1]) * 1024 + raise RuntimeError("Cannot establish host available memory") + + +@contextmanager +def admission(config): + binding = config.get("host_capacity") + if not binding: + yield () + return + path = Path(binding["lock_path"]) + fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW) + try: + info = os.fstat(fd) + if not stat.S_ISREG(info.st_mode) or info.st_mode & 0o022: + raise PermissionError( + "Host capacity lock must have a stable operator-owned inode" + ) + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + yield None + return + reserve = int(binding["reserve_bytes"]) + launch = int(binding["launch_headroom_bytes"]) + if reserve < 0 or launch <= 0: + raise ValueError("Invalid host capacity memory budget") + if available_memory_bytes() < reserve + launch: + yield None + return + yield (fd,) + finally: + # Do not explicitly LOCK_UN: an inherited descriptor must retain the + # lock if the controller dies while its coding child is still running. + os.close(fd) diff --git a/scripts/codex_von_inbox.py b/scripts/codex_von_inbox.py index 5ef5de64..5e791732 100644 --- a/scripts/codex_von_inbox.py +++ b/scripts/codex_von_inbox.py @@ -314,7 +314,7 @@ def launch(config, state, lock_fd): env=env, stdout=events, stderr=errors, - pass_fds=(lock_fd,), + pass_fds=(lock_fd, *config.get("_capacity_fds", ())), check=False, ) write_json(run / "exit.json", {"returncode": completed.returncode}) @@ -470,7 +470,7 @@ def finish(config, api, state, path): state["resume_applied"] = True write_json(path, state) content += "\n\nThe task is queued for the additional work." - note = f"Question from {config['delegator_id']}:\n{message['content']}\n\nCodex DGX answer:\n{content}" + note = f"Question from {config['delegator_id']}:\n{message['content']}\n\nCoding-agent answer:\n{content}" state["effect_stage"] = "task_note" comments = task_comments(api, task_id) existing = next( @@ -511,10 +511,10 @@ def finish(config, api, state, path): organisation_concept_id=config["organisation_id"], thread_id=task_id or message.get("thread_id"), reply_to_id=message["message_id"], - subject="Codex DGX reply", + subject=config.get("agent_name", config["agent_id"]) + " reply", content=state["report_content"], metadata={ - "attribution": "Sent by the Codex DGX coding worker", + "attribution": "Sent by " + config["agent_id"], "source_message_id": message["message_id"], }, ) @@ -570,7 +570,7 @@ def finish_or_retain(config, api, state, path): recipient_ids=[config["delegator_id"]], organisation_concept_id=config["organisation_id"], reply_to_id=message["message_id"], - subject="Codex DGX recovery pending", + subject=config.get("agent_name", config["agent_id"]) + " recovery pending", content=state.setdefault("failure_report_content", content), ) row = api.messages.get_message_for_user( diff --git a/scripts/codex_von_inbox_prompt.md b/scripts/codex_von_inbox_prompt.md index ce8caa88..f0b040f8 100644 --- a/scripts/codex_von_inbox_prompt.md +++ b/scripts/codex_von_inbox_prompt.md @@ -1,4 +1,4 @@ -You are Codex DGX answering a new message from your configured delegator. +You are the configured coding agent answering a new message from your configured delegator. The controller supplies the exact source message, bounded recent direct-message history in its organisation and as-of time, and canonical accessible task records. This is a bounded projection, not a complete canonical inventory. Consult diff --git a/scripts/codex_von_instance.py b/scripts/codex_von_instance.py new file mode 100644 index 00000000..62267676 --- /dev/null +++ b/scripts/codex_von_instance.py @@ -0,0 +1,283 @@ +"""Trusted host provisioner for one operator-enrolled coding-agent slot. + +Use a fixed command/account/config binding, with the request on stdin. A caller +cannot choose a config path, account, command or credential in its request. +The operator prepares the Unix account, Python/Codex authentication, repository, +service unit and shared capacity lock. No account creation or general shell. +""" + +from __future__ import annotations + +import argparse +import ast +import fcntl +import hashlib +import json +import os +import pwd +import re +import socket +import stat +import subprocess +import sys +from pathlib import Path + +try: + from . import codex_von_instance_schedule as schedule + from . import codex_von_release as release + from .codex_von_worker import ( + bind_backend_root, + resolve_execution_settings, + write_json, + ) +except ImportError: + import codex_von_instance_schedule as schedule + import codex_von_release as release + from codex_von_worker import ( + bind_backend_root, + resolve_execution_settings, + write_json, + ) + + +def run(argv): + return subprocess.run( + argv, check=True, text=True, capture_output=True + ).stdout.strip() + + +def schedule_action(spec, action): + if "systemd" in spec: + return schedule.control(spec, action) + return run(spec["schedule_commands"][action]) + + +def verify_instance_release(target): + """Reject older workers that would silently ignore aggregate admission.""" + target = Path(target) + worker = ast.parse((target / "scripts/codex_von_worker.py").read_text()) + supported = any( + isinstance(node, ast.Assign) + and any( + isinstance(name, ast.Name) and name.id == "INSTANCE_PROTOCOL_VERSION" + for name in node.targets + ) + and isinstance(node.value, ast.Constant) + and node.value.value == 1 + for node in worker.body + ) + if not supported or not (target / "scripts/codex_von_capacity.py").is_file(): + raise ValueError("Release does not support instance pause and host admission") + + +def binding(path): + path = Path(path) + info = path.stat() + if ( + not stat.S_ISREG(info.st_mode) + or info.st_mode & 0o022 + or info.st_uid not in {0, os.geteuid()} + ): + raise PermissionError("Host instance binding is not operator-owned") + value = json.loads(path.read_text()) + if value["unix_account"] != pwd.getpwuid(os.geteuid()).pw_name: + raise PermissionError("Prepared Unix account does not match") + if value["host"] != socket.gethostname(): + raise PermissionError("Prepared host does not match") + if not re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", value["instance_id"]): + raise ValueError("Invalid enrolled instance ID") + return value + + +def preflight(spec): + config = spec["worker_config"] + for key in ("agent_id", "delegator_id", "organisation_id"): + if config[key] != spec[key]: + raise ValueError("Worker selector differs from enrolled binding") + if config.get("report_recipient_id", spec["delegator_id"]) != spec["delegator_id"]: + raise ValueError("Report recipient must be the authorised delegator") + root = Path(spec["instance_root"]).resolve() + if Path(config["state_root"]).resolve() != root / "state": + raise ValueError("Instance state must be private to its enrolled root") + if not Path(config["source_repo"]).is_dir(): + raise ValueError("Repository is not prepared") + for field in ("codex_command", "python_environment"): + if not Path(config[field]).exists(): + raise ValueError("Codex launcher or Python environment is not prepared") + capacity = config["host_capacity"] + if not Path(capacity["lock_path"]).is_file(): + raise ValueError("Host-wide capacity lock is not prepared") + if not spec.get("capacity_enrolment_verified"): + raise ValueError("Operator must enrol all host consumers in shared admission") + if not spec.get("resource_limits_verified"): + raise ValueError("Operator must prepare the service resource limits") + if not re.fullmatch(r"[0-9a-f]{40}", spec["release_commit"]): + raise ValueError("A reviewed full release SHA is required") + if "systemd" not in spec: + for action in ("enable", "disable", "status"): + argv = spec["schedule_commands"][action] + if ( + not isinstance(argv, list) + or not argv + or not all(isinstance(x, str) for x in argv) + ): + raise ValueError("Invalid fixed schedule command") + resolve_execution_settings(config, {}) + return root + + +def lifecycle(spec, request): + if set(request) - {"action", "settings"}: + raise ValueError("Unrecognised host request fields") + action = request["action"] + if action not in {"preflight", "reconcile", "resume", "pause", "status"}: + raise ValueError("Unknown instance action") + settings = request.get("settings", {}) + if not isinstance(settings, dict) or set(settings) - {"model", "reasoning_effort"}: + raise ValueError("Invalid execution settings") + if settings and action not in {"preflight", "reconcile"}: + raise ValueError("Settings are only accepted for reconciliation") + root = preflight(spec) + identity = { + key: spec[key] + for key in ( + "instance_id", + "agent_id", + "delegator_id", + "organisation_id", + "host", + "unix_account", + ) + } + if action == "preflight": + return {**identity, "status": "prepared_host"} + root.mkdir(parents=True, exist_ok=True, mode=0o700) + os.chmod(root, 0o700) + with (root / "provision.lock").open("a") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + state_path = root / "instance.json" + state = ( + json.loads(state_path.read_text()) + if state_path.exists() + else { + **identity, + "desired_state": "paused", + "phase": "new", + } + ) + if any(state.get(key) != value for key, value in identity.items()): + raise ValueError("Retained instance belongs to another binding") + if action == "reconcile": + config = dict(spec["worker_config"]) + config["agent_name"] = spec["agent_name"] + config["instance_state_path"] = str(state_path) + settings = {**state.get("execution_overrides", {}), **settings} + for key, dest in ( + ("model", "model"), + ("reasoning_effort", "model_reasoning_effort"), + ): + if key in settings: + config[dest] = settings[key] + selected = resolve_execution_settings(config, {}) + digest = hashlib.sha256( + json.dumps(config, sort_keys=True).encode() + ).hexdigest() + # Never change configuration under an active run. Pausing the timer + # does not kill a running consumer, and the worker lock is independent. + state_root = Path(config["state_root"]) + state_root.mkdir(parents=True, exist_ok=True) + with (state_root / "worker.lock").open("a") as worker_lock: + try: + fcntl.flock(worker_lock, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + return {**identity, "status": "waiting_for_active_run"} + state.update( + phase="installing", + selected_release=spec["release_commit"], + selected_settings=selected, + execution_overrides=settings, + ) + write_json(state_path, state) + target = release.prepare( + config["source_repo"], root / "releases", spec["release_commit"] + ) + verify_instance_release(target) + write_json(root / "worker.json", config) + release.activate(root / "releases", target) + if "systemd" in spec: + schedule.install(spec) + state.update(phase="installed", config_sha256=digest) + write_json(state_path, state) + elif action == "resume": + if state.get("phase") not in {"installed", "scheduled"}: + raise ValueError("Reconcile the instance before resume") + release.current(root / "releases") + if not (root / "worker.json").is_file(): + raise ValueError("Instance configuration is not installed") + # Retain intent before schedule effects; retry reconciles a fixed unit. + state["desired_state"] = "running" + write_json(state_path, state) + elif action == "pause": + state["desired_state"] = "paused" + write_json(state_path, state) + if action != "status": + command = "enable" if state["desired_state"] == "running" else "disable" + schedule_action(spec, command) + state["phase"] = ( + "scheduled" + if command == "enable" + else ("installed" if (root / "worker.json").is_file() else "new") + ) + write_json(state_path, state) + schedule_state = json.loads(schedule_action(spec, "status")) + if not isinstance(schedule_state, dict) or not isinstance( + schedule_state.get("enabled"), bool + ): + raise TypeError("Invalid schedule status receipt") + if action != "status" and schedule_state["enabled"] != ( + state["desired_state"] == "running" + ): + raise RuntimeError("Schedule read-back does not match retained intent") + runtime_path = root / "state/runtime.json" + runtime = json.loads(runtime_path.read_text()) if runtime_path.exists() else {} + runtime = { + key: runtime[key] + for key in ( + "observed_at", + "pid", + "commit", + "backend_root", + "worker_script", + "task_service_module", + ) + if key in runtime + } + schedule_state = { + key: schedule_state[key] + for key in ("enabled", "active", "unit") + if key in schedule_state + } + # Explicitly separate scheduling from a worker's actual runtime receipt. + return { + **identity, + "status": state["phase"], + "desired_state": state["desired_state"], + "selected_release": state.get("selected_release"), + "selected_settings": state.get("selected_settings"), + "schedule": schedule_state, + "runtime_observed": runtime, + "task_and_reply_readiness": "not_established_by_provisioning", + } + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--binding", required=True) + args = parser.parse_args() + os.umask(0o077) + bind_backend_root(Path(__file__).resolve().parents[1]) + print(json.dumps(lifecycle(binding(args.binding), json.load(sys.stdin)))) + + +if __name__ == "__main__": + main() diff --git a/scripts/codex_von_instance_schedule.py b/scripts/codex_von_instance_schedule.py new file mode 100644 index 00000000..33b03e4e --- /dev/null +++ b/scripts/codex_von_instance_schedule.py @@ -0,0 +1,104 @@ +"""Deterministic per-instance systemd user units; no model-authored commands.""" + +import json +import subprocess +from pathlib import Path + + +def _quote(value, *, command=False): + value = str(value) + if any(character in value for character in ("\n", "\r", "\x00")): + raise ValueError("Invalid systemd argument") + if command: + value = value.replace("$", "$$") + return ( + '"' + value.replace("\\", "\\\\").replace('"', '\\"').replace("%", "%%") + '"' + ) + + +def install(spec): + settings = spec["systemd"] + root = Path(spec["instance_root"]) + unit_dir = Path(settings["unit_directory"]) + unit_dir.mkdir(parents=True, exist_ok=True) + name = "von-coding-" + spec["instance_id"] + marker = "# Managed coding instance " + spec["instance_id"] + " at " + str(root) + memory = int(settings["memory_max_bytes"]) + interval = int(settings.get("poll_interval_seconds", 300)) + if memory <= 0 or interval < 30: + raise ValueError("Invalid systemd resource/cadence binding") + current = root / "releases/current" + argv = [ + Path(spec["worker_config"]["python_environment"]) / "bin/python", + current / "scripts/codex_von_worker.py", + "--backend-root", + current, + "--config", + root / "worker.json", + ] + service = [ + marker, + "[Unit]", + "Description=Von coding instance " + spec["instance_id"], + "[Service]", + "Type=oneshot", + "UMask=0077", + "KillMode=control-group", + "MemoryMax=" + str(memory), + "MemorySwapMax=0", + "TimeoutStartSec=infinity", + "ExecStart=" + " ".join(_quote(value, command=True) for value in argv), + ] + for path in settings.get("environment_files", []): + if not Path(path).is_file(): + raise ValueError("Enrolled environment file is not prepared") + service.append("EnvironmentFile=" + _quote(path)) + timer = [ + marker, + "[Unit]", + "Description=Poll Von coding instance " + spec["instance_id"], + "[Timer]", + "OnActiveSec=60", + "OnUnitInactiveSec=" + str(interval), + "Unit=" + name + ".service", + "[Install]", + "WantedBy=timers.target", + ] + for suffix, lines in (("service", service), ("timer", timer)): + path = unit_dir / (name + "." + suffix) + if path.exists() and not path.read_text().startswith(marker + "\n"): + raise PermissionError("Refusing to replace an unenrolled service unit") + temporary = path.with_suffix(path.suffix + ".tmp") + temporary.write_text("\n".join(lines) + "\n") + temporary.replace(path) + subprocess.run( + ["systemctl", "--user", "daemon-reload"], check=True, capture_output=True + ) + + +def control(spec, action): + name = "von-coding-" + spec["instance_id"] + ".timer" + if action in {"enable", "disable"}: + subprocess.run( + ["systemctl", "--user", action, "--now", name], + check=True, + capture_output=True, + ) + return "" + enabled = subprocess.run( + ["systemctl", "--user", "is-enabled", "--quiet", name], + capture_output=True, + check=False, + ) + active = subprocess.run( + ["systemctl", "--user", "is-active", "--quiet", name], + capture_output=True, + check=False, + ) + return json.dumps( + { + "enabled": enabled.returncode == 0, + "active": active.returncode == 0, + "unit": name, + } + ) diff --git a/scripts/codex_von_worker.py b/scripts/codex_von_worker.py index a9dfd523..7274a5cb 100755 --- a/scripts/codex_von_worker.py +++ b/scripts/codex_von_worker.py @@ -24,6 +24,7 @@ from pathlib import Path ACTIVITY_CHECKPOINT_SECONDS = 15 +INSTANCE_PROTOCOL_VERSION = 1 ACTIVE = {"pending", "in_progress", "blocked"} RESULT_SCHEMA = { @@ -75,6 +76,7 @@ def capability_context(config, *, read_only, worktree=None): "delegator_id": config["delegator_id"], "organisation_id": config["organisation_id"], "source_repo": config.get("source_repo"), + "approved_repository_roots": config.get("approved_repository_roots", []), "worktree": str(worktree) if worktree else None, "controller_runtime": config.get("runtime_evidence", {"available": False}), "public_served_revision": {"available": False, "reason": "Not probed here"}, @@ -711,9 +713,9 @@ def send(self, task, key, content): recipient_ids=[c["delegator_id"]], organisation_concept_id=c["organisation_id"], thread_id=task["task_concept_id"], - subject=f"Codex DGX: {task['title']}", + subject=f"{c.get('agent_name', c['agent_id'])}: {task['title']}", content=content, - metadata={"attribution": "Sent by the Codex DGX coding worker"}, + metadata={"attribution": "Sent by " + c["agent_id"]}, ) message_id = receipt.message["concept_id"] readback = self.messages.get_message_for_user(message_id, c["agent_id"]) @@ -741,7 +743,7 @@ def finish(self, state): content = ( f"{result['summary']}\n\n{result['evidence']}\n\n" f"{result['question']}\n\nTask: {state['task_id']}\n" - f"DGX worktree: {state.get('worktree', 'not started')}\n" + f"Coding worktree: {state.get('worktree', 'not started')}\n" f"Run: {state['attempt']}" ).strip() if state.get("execution_settings"): @@ -776,7 +778,7 @@ def finish(self, state): evidence_addition = ( f"{result['evidence']}\nVon result: {message_id}\n" f"Execution archive: {capture.get('archive_url', 'unavailable')}\n" - f"DGX worktree: {state.get('worktree', 'not started')}" + f"Coding worktree: {state.get('worktree', 'not started')}" ) evidence = task.get("evidence") or "" if evidence_addition not in evidence: @@ -1030,7 +1032,7 @@ def launch(config, state, inputs, conversation, state_path, lock_fd): env=environment, stdout=events, stderr=errors, - pass_fds=(lock_fd,), + pass_fds=(lock_fd, *config.get("_capacity_fds", ())), ) as process: process.stdin.write(prompt) process.stdin.close() @@ -1136,9 +1138,9 @@ def apply_deployment(config, api, state, state_path, lock_fd): receipt ) if success: - result[ - "summary" - ] += f" Deployed and verified the public server at {commit[:12]}." + result["summary"] += ( + f" Deployed and verified the public server at {commit[:12]}." + ) else: result["status"] = "blocked" result["summary"] += ( @@ -1149,9 +1151,9 @@ def apply_deployment(config, api, state, state_path, lock_fd): ) if recover_only: result["status"] = "blocked" - result[ - "summary" - ] += " The task changed; only the interrupted deployment was reconciled." + result["summary"] += ( + " The task changed; only the interrupted deployment was reconciled." + ) write_json(state_path, state) @@ -1249,7 +1251,7 @@ def tick(config, api, lock_fd): api.send( task, fingerprint([task_id, input_fingerprint(inputs)]) + ":started", - f"I am picking up this task on the DGX: {task['title']}.\nTask: {task_id}\n" + f"I am picking up this assigned coding task: {task['title']}.\nTask: {task_id}\n" "I will send the result or a question here. Reply in this task thread or add a task comment.", ) api.tasks.update_task_status( @@ -1269,9 +1271,9 @@ def tick(config, api, lock_fd): state["result"]["summary"] = ( "Coding execution settings rejected: " + str(exc) ) - state["result"][ - "question" - ] = "Correct this task's model or reasoning setting, then requeue it." + state["result"]["question"] = ( + "Correct this task's model or reasoning setting, then requeue it." + ) write_json(path, state) apply_deployment(config, api, state, path, lock_fd) finish_captured(config, api, state, path) @@ -1313,12 +1315,33 @@ def main(): config = json.loads(Path(options.config).read_text()) state_root = Path(config["state_root"]) state_root.mkdir(parents=True, exist_ok=True) - with (state_root / "worker.lock").open("a") as lock: + with (state_root / "worker.lock").open("a") as lock, ExitStack() as admission_stack: try: fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError: print('{"outcome":"already_running"}') return + if config.get("instance_state_path"): + instance = json.loads(Path(config["instance_state_path"]).read_text()) + if any( + instance.get(key) != config[key] + for key in ("agent_id", "delegator_id", "organisation_id") + ): + raise PermissionError( + "Instance scope differs from worker configuration" + ) + if instance.get("desired_state") != "running" and not options.check_task: + print('{"outcome":"paused"}') + return + try: + from .codex_von_capacity import admission + except ImportError: + from codex_von_capacity import admission + capacity_fds = admission_stack.enter_context(admission(config)) + if capacity_fds is None: + print(json.dumps({"outcome": "waiting_for_capacity"}), flush=True) + return + config["_capacity_fds"] = capacity_fds backend_root = bind_backend_root(options.backend_root) resolve_execution_settings(config, {}) from src.backend.integrations.internal_mcp.gateway import ( diff --git a/scripts/codex_von_worker_prompt.md b/scripts/codex_von_worker_prompt.md index aff48065..74d47bf2 100644 --- a/scripts/codex_von_worker_prompt.md +++ b/scripts/codex_von_worker_prompt.md @@ -1,5 +1,5 @@ -You are the Codex DGX coding worker, acting on a task explicitly assigned to -your configured Von identity by Michael Witbrock. Read AGENTS.md and the +You are the configured coding worker, acting on a task explicitly assigned to +your configured Von identity by the configured delegator. Read AGENTS.md and the assigned context file, then carry out the task in this isolated worktree. The task and Michael's task comments/direct replies convey the requested work. @@ -47,7 +47,8 @@ changes/tests succeeded. Do not replace credentials or change GitHub accounts. Deployment is requested through your structured result, never through shell access to the live service. Set `deploy_commit` to the full merged Git SHA only when the initial task title or description explicitly instructs deployment to -the public DGX server. Otherwise return an empty string. Conversation context, +the operator-configured deployment destination. Otherwise return an empty string. +Conversation context, quoted text and completion of a fix do not by themselves request deployment. For an authorised deployment, finish the code, tests and required publication, confirm that the intended revision is the current origin/main, and request that @@ -63,7 +64,7 @@ Return the required JSON result. `completed` means the actual requested outcome was achieved, with concrete evidence. Use `needs_input` for a question that must be answered and `blocked` for a technical impediment. Include the question in plain English, or an empty string when none is needed. Keep the summary useful -to Michael without requiring him to read raw logs. +to the configured delegator without requiring him to read raw logs. The supplied capabilities describe the controller and this sandbox, with known limitations; configured execution settings are not provider-observed model diff --git a/src/backend/integrations/internal_mcp/catalogue.py b/src/backend/integrations/internal_mcp/catalogue.py index 197e600a..1af8ea4a 100644 --- a/src/backend/integrations/internal_mcp/catalogue.py +++ b/src/backend/integrations/internal_mcp/catalogue.py @@ -47,6 +47,7 @@ from .gateway import MethodCatalogue, MethodDefinition from .schemas import Schema, make_error_response from .spreadsheet_record_tools import build_spreadsheet_record_tool_definitions +from .coding_agent_instance_tools import build_coding_agent_instance_tools from .workflow_surface_capabilities import ( build_workflow_surface_capability_matrix, ) @@ -40706,6 +40707,7 @@ def _build_default_catalogue_core_definitions() -> List[MethodDefinition]: ), ), *build_spreadsheet_record_tool_definitions(), + *build_coding_agent_instance_tools(), MethodDefinition( name="extract_annotations", handler=_extract_annotations, diff --git a/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py b/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py new file mode 100644 index 00000000..80fe506a --- /dev/null +++ b/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py @@ -0,0 +1,50 @@ +"""Discoverable coding-agent instance lifecycle through prepared host bindings.""" + +from .gateway import MethodDefinition +from .schemas import Schema, make_error_response + + +def _instances(**kwargs): + from ...services.coding_agent_instance_service import coding_agent_instances + + try: + return coding_agent_instances(**kwargs) + except ( + PermissionError, + ValueError, + TypeError, + KeyError, + OSError, + RuntimeError, + ) as exc: + return make_error_response( + "coding_agent_instance_unavailable", + type(exc).__name__ + ": " + str(exc) + if isinstance(exc, (PermissionError, ValueError)) + else "Instance operation unavailable; inspect operator readiness and private logs.", + ) + + +def build_coding_agent_instance_tools(): + return [ + MethodDefinition( + name="coding_agent_instances", + handler=_instances, + input_schema=Schema( + required={}, + optional={"action": str, "instance_id": str, "settings": dict}, + allow_unknown=False, + ), + output_schema=None, + category="write", + description=( + "List, inspect, provision/reconcile, pause or resume coding-agent instances " + "for the authenticated delegator in the current organisation. Uses existing " + "operator-enrolled host/account/identity/credential bindings; cannot create " + "host accounts or grant new access. Reconcile prepares a paused instance; " + "resume explicitly enables its enrolled schedule. Settings accepts model " + "and reasoning_effort; task overrides remain supported. Readiness requires " + "actual task and recipient evidence, not a selected release or active timer." + ), + ) + ] diff --git a/src/backend/services/coding_agent_instance_service.py b/src/backend/services/coding_agent_instance_service.py new file mode 100644 index 00000000..19927702 --- /dev/null +++ b/src/backend/services/coding_agent_instance_service.py @@ -0,0 +1,211 @@ +"""Actor-bound access to operator-enrolled coding-agent instance slots. + +Host/account/credential/identity bindings are operator configuration, never tool +arguments. Normal Von agent calls use the authenticated delegator's authority. +The host independently checks the slot and installs through the existing worker +release mechanism. This module never writes Vontology through a database handle. +""" + +from __future__ import annotations + +import json +import os +import stat +import subprocess +from pathlib import Path + + +def load_registry(): + path = Path(os.environ["VON_CODING_AGENT_REGISTRY"]) + info = path.stat() + if not stat.S_ISREG(info.st_mode) or info.st_mode & 0o022: + raise PermissionError( + "Instance registry must be operator-owned and not writable by group/other" + ) + if info.st_uid not in {0, os.geteuid()}: + raise PermissionError("Unexpected instance registry owner") + registry = json.loads(path.read_text()) + slots = registry["instances"] + # One configured consumer per represented identity, including existing + # consumers enrolled as reserved slots. Endpoint does not create identity. + identities = [slot["agent_id"] for slot in slots.values()] + if len(identities) != len(set(identities)): + raise ValueError("Duplicate coding-agent identity in registry") + return registry + + +def authorised_slots(): + from ..security.access_control import get_effective_organisation_concept_id + from ..security.role_resolver import get_effective_permissions + from .organisation_membership_governance_service import _trusted_actor + from .organisation_membership_service import resolve_user_organisation_membership + + actor, error = _trusted_actor() + if error: + raise PermissionError(error["error_code"]) + organisation = get_effective_organisation_concept_id() + membership = resolve_user_organisation_membership(actor, organisation) + if not membership or "MANAGE_MEMBERS" not in get_effective_permissions( + membership["role"] + ): + raise PermissionError( + "Live organisation membership-management authority required" + ) + registry = load_registry() + slots = { + key: slot + for key, slot in registry["instances"].items() + if slot["delegator_id"] == actor and slot["organisation_id"] == organisation + } + return actor, organisation, slots + + +def _host(slot, action, settings): + # The entire argv is an operator binding. The model cannot append flags, + # choose a host/account/path, substitute an endpoint or supply a shell. + argv = slot["provisioner_command"] + if ( + not isinstance(argv, list) + or not argv + or not all(isinstance(x, str) for x in argv) + ): + raise ValueError("Invalid operator provisioner command") + request = {"action": action, "settings": settings} + process = subprocess.run( + argv, + input=json.dumps(request), + text=True, + capture_output=True, + check=False, + ) + if process.returncode: + # A host error may contain credential paths or command output. + raise RuntimeError("Host provisioner failed; inspect its private operator log") + result = json.loads(process.stdout) + if not isinstance(result, dict): + raise TypeError("Invalid host receipt") + for key in ( + "instance_id", + "agent_id", + "delegator_id", + "organisation_id", + "host", + "unix_account", + ): + if result.get(key) != slot[key]: + raise RuntimeError( + "Host receipt does not match the authorised instance binding" + ) + return result + + +def _identity(slot, actor, organisation): + from . import concept_service + from .organisation_membership_governance_service import ( + manage_organisation_membership, + ) + from .organisation_membership_service import resolve_user_organisation_membership + + agent_id = slot["agent_id"] + concept = concept_service.get_concept_by_concept_id(agent_id) + if concept is None: + if not slot.get("allow_identity_creation", False): + raise PermissionError( + "The enrolled agent identity must be prepared by its operator" + ) + concept_service.create_concept( + name=slot["agent_name"], + concept_id=agent_id, + description=slot["role"], + parent_concept_ids=["#V#coding_agent"], + create_as_instance=True, + created_by_concept_id=actor, + organisation_concept_id=organisation, + visibility_scope_mode="organisation_general", + ) + concept = concept_service.get_concept_by_concept_id(agent_id) + if not concept or "#V#coding_agent" not in concept.get("relationships", {}).get( + "is_an_instance_of", [] + ): + raise PermissionError( + "Enrolled identity is not a canonical coding-agent instance" + ) + if not resolve_user_organisation_membership(agent_id, organisation): + receipt = manage_organisation_membership( + action="add", + user_concept_id=agent_id, + organisation_concept_id=organisation, + role="member", + request_id="coding-agent-instance:" + slot["instance_id"], + reason="Provision operator-enrolled coding-agent instance", + ) + if not receipt.get("success"): + raise PermissionError("Canonical agent membership creation did not succeed") + if not resolve_user_organisation_membership(agent_id, organisation): + raise RuntimeError("Agent membership read-back failed") + + +def coding_agent_instances(*, action="list", instance_id=None, settings=None): + """Discover or reconcile an enrolled instance without granting host access.""" + actor, organisation, slots = authorised_slots() + if action == "list": + return { + "instances": [ + { + key: slot[key] + for key in ( + "instance_id", + "agent_id", + "agent_name", + "role", + "delegator_id", + "organisation_id", + "endpoint", + "host", + "unix_account", + ) + } + for slot in slots.values() + ] + } + if instance_id not in slots: + raise PermissionError( + "Instance is not enrolled for the current actor and organisation" + ) + if action not in {"status", "reconcile", "pause", "resume"}: + raise ValueError("Use list, status, reconcile, pause or resume") + slot = slots[instance_id] + if slot["instance_id"] != instance_id or slot.get("reserved"): + raise PermissionError( + "Existing consumer is reserved; use its own operator route" + ) + values = settings or {} + if not isinstance(values, dict) or set(values) - {"model", "reasoning_effort"}: + raise ValueError("Only model and reasoning_effort are selectable settings") + if action != "reconcile" and values: + raise ValueError("Settings can only change during reconcile") + from ..utils.task_execution_preferences import normalise_task_execution_preference + + for key, value in values.items(): + normalised = normalise_task_execution_preference( + value, field_name="requested_" + key + ) + if normalised is None or normalised != value: + raise ValueError("Specify an exact execution-setting token") + if "sol" in str(values.get("model", "")).lower(): + raise ValueError("Sol-family models are disabled") + # Check host readiness/binding before represented identity effects. Reconcile + # remains paused until an explicit resume; replay never silently resumes. + if action in {"reconcile", "resume"}: + preflight = _host(slot, "preflight", values) + if preflight.get("status") != "prepared_host": + return preflight + _identity(slot, actor, organisation) + result = _host(slot, action, values) + return { + **result, + "instance_id": instance_id, + "agent_id": slot["agent_id"], + "delegator_id": actor, + "organisation_id": organisation, + } diff --git a/tests/backend/test_codex_von_worker.py b/tests/backend/test_codex_von_worker.py index 3f2307bb..e050e51b 100644 --- a/tests/backend/test_codex_von_worker.py +++ b/tests/backend/test_codex_von_worker.py @@ -69,6 +69,22 @@ def test_empty_poll_has_no_model_or_message(config, monkeypatch): worker.tick(config, SimpleNamespace(pending=list), 0) +def test_paused_instance_stops_before_loading_backend(config, monkeypatch, capsys): + config_path = Path(config["state_root"]) / "config.json" + state_path = Path(config["state_root"]) / "instance.json" + config["instance_state_path"] = str(state_path) + state_path.write_text(json.dumps({**config, "desired_state": "paused"})) + config_path.write_text(json.dumps(config)) + monkeypatch.setattr(sys, "argv", ["worker", "--config", str(config_path)]) + monkeypatch.setattr( + worker, + "bind_backend_root", + lambda *_: pytest.fail("Paused worker loaded backend"), + ) + worker.main() + assert json.loads(capsys.readouterr().out)["outcome"] == "paused" + + @pytest.mark.parametrize( "requested,expected", [ @@ -420,7 +436,9 @@ def test_interrupted_clone_recovers_with_independent_metadata_and_no_service_sec ) run(["git", "-C", str(source), "remote", "add", "origin", str(source)]) fake = tmp_path / "fake-codex" - fake.write_text(f"#!{sys.executable}\n" + """import json, os, sys + fake.write_text( + f"#!{sys.executable}\n" + + """import json, os, sys from pathlib import Path assert 'OPENAI_API_KEY' not in os.environ assert 'MONGO_URI' not in os.environ @@ -430,7 +448,8 @@ def test_interrupted_clone_recovers_with_independent_metadata_and_no_service_sec assert 'model_reasoning_effort="high"' in sys.argv result = Path(sys.argv[sys.argv.index('--output-last-message') + 1]) result.write_text(json.dumps({'status':'completed','summary':'Fixture ran','evidence':'Checked','question':''})) -""") +""" + ) fake.chmod(0o700) monkeypatch.setenv("OPENAI_API_KEY", "test-sentinel") monkeypatch.setenv("MONGO_URI", "test-sentinel") diff --git a/tests/backend/test_coding_agent_instances.py b/tests/backend/test_coding_agent_instances.py new file mode 100644 index 00000000..b29bd4b9 --- /dev/null +++ b/tests/backend/test_coding_agent_instances.py @@ -0,0 +1,349 @@ +"""Prepared instance lifecycle and its actor/host authority boundary.""" + +import json +from pathlib import Path +from types import SimpleNamespace + +import pytest + +pytest.importorskip("fcntl") + +from scripts import codex_von_instance as host +from src.backend.services import coding_agent_instance_service as service + + +@pytest.fixture +def slot(tmp_path): + return { + "instance_id": "test-agent", + "agent_id": "#V#test_agent", + "agent_name": "Test Agent", + "role": "Test coding role", + "delegator_id": "#V#owner", + "organisation_id": "#V#org", + "endpoint": "http://localhost:5000", + "host": "test-host", + "unix_account": "test-account", + "provisioner_command": ["/trusted/provisioner"], + "instance_root": str(tmp_path / "instance"), + "release_commit": "a" * 40, + "capacity_enrolment_verified": True, + "resource_limits_verified": True, + "schedule_commands": { + x: ["schedule", x] for x in ("enable", "disable", "status") + }, + "worker_config": { + "agent_id": "#V#test_agent", + "delegator_id": "#V#owner", + "organisation_id": "#V#org", + "source_repo": str(tmp_path), + "state_root": str(tmp_path / "instance/state"), + "codex_command": str(tmp_path), + "python_environment": str(tmp_path), + "host_capacity": {"lock_path": str(tmp_path / "host.lock")}, + }, + } + + +@pytest.fixture +def actor(monkeypatch, slot): + from src.backend.security import access_control + from src.backend.services import ( + organisation_membership_governance_service as governance, + ) + from src.backend.services import organisation_membership_service as memberships + + monkeypatch.setattr(governance, "_trusted_actor", lambda: ("#V#owner", None)) + monkeypatch.setattr( + access_control, "get_effective_organisation_concept_id", lambda: "#V#org" + ) + monkeypatch.setattr( + memberships, + "resolve_user_organisation_membership", + lambda *_: {"role": "owner"}, + ) + monkeypatch.setattr( + service, "load_registry", lambda: {"instances": {"test-agent": slot}} + ) + + +def test_normal_tool_path_lists_only_scoped_bindings(actor, slot): + from src.backend.integrations.internal_mcp.coding_agent_instance_tools import ( + _instances, + ) + + result = _instances(action="list") + assert result["instances"][0]["agent_id"] == "#V#test_agent" + assert "provisioner_command" not in result["instances"][0] + slot["organisation_id"] = "#V#other" + assert _instances(action="list") == {"instances": []} + assert _instances(action="resume", instance_id="test-agent")["success"] is False + + +def test_payload_actor_cannot_create_authority(actor, monkeypatch): + from src.backend.services import ( + organisation_membership_governance_service as governance, + ) + + monkeypatch.setattr( + governance, + "_trusted_actor", + lambda: (None, {"error_code": "client_supplied_identity_is_not_authority"}), + ) + with pytest.raises(PermissionError, match="client_supplied"): + service.coding_agent_instances(action="list") + + +def test_live_membership_required_even_for_enrolled_owner(actor, monkeypatch): + from src.backend.services import organisation_membership_service as memberships + + monkeypatch.setattr( + memberships, + "resolve_user_organisation_membership", + lambda *_: {"role": "member"}, + ) + with pytest.raises(PermissionError, match="Live organisation"): + service.coding_agent_instances(action="resume", instance_id="test-agent") + + +@pytest.mark.parametrize( + "settings", + [ + {"command": "evil"}, + {"model": "gpt-5.6-sol"}, + {"reasoning_effort": "extra high"}, + ], +) +def test_tool_rejects_commands_and_bad_settings_before_effects( + actor, monkeypatch, settings +): + monkeypatch.setattr(service, "_host", lambda *_: pytest.fail("host effect")) + with pytest.raises(ValueError): + service.coding_agent_instances( + action="reconcile", instance_id="test-agent", settings=settings + ) + + +def test_host_receipt_must_match_all_bindings(slot, monkeypatch): + response = { + k: slot[k] + for k in ( + "instance_id", + "agent_id", + "delegator_id", + "organisation_id", + "host", + "unix_account", + ) + } + response["organisation_id"] = "#V#wrong" + monkeypatch.setattr( + service.subprocess, + "run", + lambda *a, **kw: SimpleNamespace(returncode=0, stdout=json.dumps(response)), + ) + with pytest.raises(RuntimeError, match="does not match"): + service._host(slot, "status", {}) + + +def test_registry_rejects_duplicate_identity(tmp_path, slot, monkeypatch): + path = tmp_path / "registry.json" + path.write_text(json.dumps({"instances": {"one": slot, "two": slot}})) + path.chmod(0o600) + monkeypatch.setenv("VON_CODING_AGENT_REGISTRY", str(path)) + with pytest.raises(ValueError, match="Duplicate"): + service.load_registry() + path.chmod(0o666) + with pytest.raises(PermissionError): + service.load_registry() + + +@pytest.fixture +def prepared(slot, monkeypatch): + Path(slot["worker_config"]["host_capacity"]["lock_path"]).touch() + state = {"enabled": False, "prepared": 0} + + def schedule(argv): + if argv[-1] == "enable": + state["enabled"] = True + elif argv[-1] == "disable": + state["enabled"] = False + return json.dumps({"enabled": state["enabled"]}) + + def prepare(primary, root, sha): + state["prepared"] += 1 + target = root / sha + target.mkdir(parents=True, exist_ok=True) + (target / "scripts").mkdir(exist_ok=True) + (target / "scripts/codex_von_worker.py").write_text( + "INSTANCE_PROTOCOL_VERSION = 1" + ) + (target / "scripts/codex_von_capacity.py").touch() + return target + + def activate(root, target): + link = root / "current" + link.unlink(missing_ok=True) + link.symlink_to(target) + + monkeypatch.setattr(host, "run", schedule) + monkeypatch.setattr(host.release, "prepare", prepare) + monkeypatch.setattr(host.release, "activate", activate) + return state + + +def test_reconcile_repeat_pause_resume_and_interrupted_activation( + slot, prepared, monkeypatch +): + result = host.lifecycle( + slot, {"action": "reconcile", "settings": {"reasoning_effort": "high"}} + ) + assert result["desired_state"] == "paused" + assert not prepared["enabled"] + assert result["selected_settings"]["reasoning_effort"] == "high" + assert result["task_and_reply_readiness"] == "not_established_by_provisioning" + host.lifecycle(slot, {"action": "reconcile"}) + root = Path(slot["instance_root"]) + assert ( + json.loads((root / "worker.json").read_text())["model_reasoning_effort"] + == "high" + ) + assert len(list((root / "releases").iterdir())) == 2 # one release and current link + actual_run = host.run + monkeypatch.setattr( + host, "run", lambda *_: (_ for _ in ()).throw(OSError("interrupted")) + ) + with pytest.raises(OSError): + host.lifecycle(slot, {"action": "resume"}) + assert ( + json.loads((root / "instance.json").read_text())["desired_state"] == "running" + ) + monkeypatch.setattr(host, "run", actual_run) + assert host.lifecycle(slot, {"action": "resume"})["schedule"]["enabled"] + assert not host.lifecycle(slot, {"action": "pause"})["schedule"]["enabled"] + assert not host.lifecycle(slot, {"action": "reconcile"})["schedule"]["enabled"] + + +def test_active_instance_not_reconfigured(slot, prepared): + import fcntl + + root = Path(slot["worker_config"]["state_root"]) + root.mkdir(parents=True) + with (root / "worker.lock").open("a") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + result = host.lifecycle(slot, {"action": "reconcile"}) + assert result["status"] == "waiting_for_active_run" + assert prepared["prepared"] == 0 + + +def test_pause_before_install_cannot_activate(slot, prepared): + assert host.lifecycle(slot, {"action": "pause"})["status"] == "new" + with pytest.raises(ValueError, match="Reconcile"): + host.lifecycle(slot, {"action": "resume"}) + + +def test_old_worker_release_cannot_bypass_host_admission(tmp_path): + (tmp_path / "scripts").mkdir() + (tmp_path / "scripts/codex_von_worker.py").write_text("OLD_WORKER = True") + with pytest.raises(ValueError, match="does not support"): + host.verify_instance_release(tmp_path) + + +def test_host_rejects_account_and_selector_changes(slot, prepared, tmp_path): + slot["worker_config"]["organisation_id"] = "#V#wrong" + with pytest.raises(ValueError, match="selector"): + host.lifecycle(slot, {"action": "reconcile"}) + path = tmp_path / "binding.json" + path.write_text(json.dumps(slot)) + path.chmod(0o600) + with pytest.raises(PermissionError, match="account"): + host.binding(path) + + +def test_tool_reconcile_uses_canonical_identity_and_host_receipts( + actor, slot, prepared, monkeypatch +): + from src.backend.integrations.internal_mcp.coding_agent_instance_tools import ( + _instances, + ) + from src.backend.services import concept_service + from src.backend.services import ( + organisation_membership_governance_service as governance, + ) + from src.backend.services import organisation_membership_service as memberships + + concepts, members = {}, {"#V#owner"} + slot["allow_identity_creation"] = True + monkeypatch.setattr(concept_service, "get_concept_by_concept_id", concepts.get) + + def create(**kwargs): + assert kwargs["organisation_concept_id"] == "#V#org" + assert kwargs["created_by_concept_id"] == "#V#owner" + concepts[kwargs["concept_id"]] = { + "relationships": {"is_an_instance_of": ["#V#coding_agent"]}, + } + + def membership(**kwargs): + assert kwargs["organisation_concept_id"] == "#V#org" + assert kwargs["role"] == "member" + members.add(kwargs["user_concept_id"]) + return {"success": True} + + monkeypatch.setattr(concept_service, "create_concept", create) + monkeypatch.setattr(governance, "manage_organisation_membership", membership) + monkeypatch.setattr( + memberships, + "resolve_user_organisation_membership", + lambda user, org: ( + {"role": "owner" if user == "#V#owner" else "member"} + if org == "#V#org" and user in members + else None + ), + ) + monkeypatch.setattr( + service, + "_host", + lambda s, a, values: host.lifecycle(s, {"action": a, "settings": values}), + ) + result = _instances( + action="reconcile", + instance_id="test-agent", + settings={"model": "gpt-6-astra", "reasoning_effort": "high"}, + ) + assert result["status"] == "installed" + assert result["agent_id"] == "#V#test_agent" + assert result["desired_state"] == "paused" + assert len(concepts) == 1 + assert ( + _instances(action="reconcile", instance_id="test-agent")["status"] + == "installed" + ) + assert len(concepts) == 1 + slot["reserved"] = True + assert _instances(action="resume", instance_id="test-agent")["success"] is False + + +def test_systemd_units_are_reused_without_clobbering_other_units( + slot, tmp_path, monkeypatch +): + from scripts import codex_von_instance_schedule as schedule + + commands = [] + monkeypatch.setattr( + schedule.subprocess, "run", lambda argv, **kw: commands.append(argv) + ) + slot["systemd"] = { + "unit_directory": str(tmp_path / "units"), + "memory_max_bytes": 4 * 1024**3, + } + schedule.install(slot) + schedule.install(slot) + files = sorted((tmp_path / "units").iterdir()) + assert len(files) == 2 + service_unit = next(p for p in files if p.suffix == ".service") + assert "MemoryMax=4294967296" in service_unit.read_text() + assert "releases/current/scripts/codex_von_worker.py" in service_unit.read_text() + assert all(command[-1] == "daemon-reload" for command in commands) + service_unit.write_text("unrelated unit") + with pytest.raises(PermissionError, match="unenrolled"): + schedule.install(slot) diff --git a/tests/test_codex_von_capacity.py b/tests/test_codex_von_capacity.py new file mode 100644 index 00000000..9d59d8f1 --- /dev/null +++ b/tests/test_codex_von_capacity.py @@ -0,0 +1,66 @@ +"""Actual lock contenders; no model calls or coding workers are launched.""" + +import subprocess +import sys + +import pytest + +pytest.importorskip("fcntl") +from scripts import codex_von_capacity as capacity + + +def test_two_contenders_and_inherited_lock_survive_controller_close( + tmp_path, monkeypatch +): + lock = tmp_path / "host.lock" + lock.touch(mode=0o644) + config = { + "host_capacity": { + "lock_path": str(lock), + "reserve_bytes": 10, + "launch_headroom_bytes": 20, + } + } + monkeypatch.setattr(capacity, "available_memory_bytes", lambda: 100) + with capacity.admission(config) as fds: + assert fds is not None + child = subprocess.Popen( + [ + sys.executable, + "-c", + "import sys; print('holding', flush=True); sys.stdin.read(1)", + ], + pass_fds=fds, + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + text=True, + ) + assert child.stdout.readline().strip() == "holding" + with capacity.admission(config) as second: + assert second is None + try: + with capacity.admission(config) as second: + assert second is None # child retains the original open description + finally: + child.communicate("x", timeout=5) + with capacity.admission(config) as second: + assert second is not None + monkeypatch.setattr(capacity, "available_memory_bytes", lambda: 29) + with capacity.admission(config) as second: + assert second is None + + +def test_unsafe_or_missing_lock_denies(tmp_path): + path = tmp_path / "lock" + config = {"host_capacity": {"lock_path": str(path)}} + with pytest.raises(FileNotFoundError), capacity.admission(config): + pass + path.touch() + path.chmod(0o666) + with pytest.raises(PermissionError), capacity.admission(config): + pass + + +def test_legacy_unenrolled_worker_is_not_silently_migrated(): + with capacity.admission({}) as fds: + assert fds == () From 7cd1d721830cb5b7f458f7c623536c09fe29618e Mon Sep 17 00:00:00 2001 From: witbrock Date: Sun, 13 Sep 2026 18:25:07 -0400 Subject: [PATCH 2/7] Add optional task-scoped Android test-host candidate --- .../coding_agent_android_testing.md | 108 ++++ docs/engineering/coding_agent_instances.md | 7 + scripts/codex_von_android.py | 595 ++++++++++++++++++ scripts/codex_von_android_control.py | 131 ++++ scripts/codex_von_android_transport.py | 70 +++ scripts/codex_von_capacity.py | 11 + scripts/codex_von_worker.py | 76 ++- .../coding_agent_instance_tools.py | 57 +- .../services/coding_agent_android_service.py | 137 ++++ tests/backend/test_coding_agent_android.py | 284 +++++++++ 10 files changed, 1460 insertions(+), 16 deletions(-) create mode 100644 docs/engineering/coding_agent_android_testing.md create mode 100644 scripts/codex_von_android.py create mode 100644 scripts/codex_von_android_control.py create mode 100644 scripts/codex_von_android_transport.py create mode 100644 src/backend/services/coding_agent_android_service.py create mode 100644 tests/backend/test_coding_agent_android.py diff --git a/docs/engineering/coding_agent_android_testing.md b/docs/engineering/coding_agent_android_testing.md new file mode 100644 index 00000000..f5c78b97 --- /dev/null +++ b/docs/engineering/coding_agent_android_testing.md @@ -0,0 +1,108 @@ +# Optional Android testing for coding-agent instances + +This is a **candidate follow-on** to the prepared instance package. It is not +production enabled or accepted yet. It must not delay the core package. + +The intended job is native Chrome/Gboard investigation without a connected +phone. The simplest baseline is the existing operator-owned Mac SDK/image, +with new disposable AVD data for each task/run. No emulator or account is +created just because a coding instance is provisioned. + +## Authority and routes + +`coding_agent_android` uses the existing `VON_CODING_AGENT_REGISTRY`. An +instance can have an optional `android` entry with a fixed `provisioner_command`, +`host` and `unix_account`. The test host can differ from the coding host. The +service accepts only the authenticated enrolled agent or delegator, in the +exact organisation, with current membership and a current task created by the +enrolled delegator and assigned to that agent. No client actor/org/path/command +is accepted. Reserved coding identities can use this capability without +reprovisioning their worker. + +The fixed command invokes `scripts/codex_von_android.py --binding PRIVATE_JSON`. +The binding uses the core provisioner's host/account/identity validation and +adds `android`: `platform` (OS/architecture pair), `sdk_root`, `system_image` +(relative to SDK), `abi`, `arch`, `state_root`, `ram_mib`, `cores`, `adb_port`, +`emulator_port`, `host_capacity`, `capacity_enrolment_verified`, and a map of +named `fixtures` to operator-approved disposable URLs. Defaults bound failed +boot occupancy to 600 seconds and idle occupancy to 1800 seconds; legitimate +operations renew idle activity. These protect shared test-host resources, not +whole coding-task duration. + +The SDK and image are immutable shared inputs. Mutable AVD data and control +sockets are private to the prepared account/binding and task/run. All test-host +consumers must use the same operator-owned capacity inode; the core package's +single-slot admission serialises consumers. Guest RAM/CPU are explicit launch +arguments; headroom includes host overhead. A reserve check runs during boot +and use. This is not a hard total-RSS guarantee or a defence against arbitrary +same-account processes. Separate untrusted tenants require prepared accounts +and suitable OS confinement. + +For a Mac without incoming SSH, `codex_von_android_transport.py` supports a +private Unix listener and client. An operator can use SSH StreamLocal reverse +forwarding from the Mac to a private DGX socket. Both parent directories must be +owner-only. The DGX registry binds the client command to that exact socket; the +Mac listener binds a fixed host configuration. No TCP listener, remote-login +setting, credentials in model context or extra coding worker is needed. The +trusted service/operator owns the socket. Do not expose it directly to an +untrusted model process, since the host trusts the canonical service's prior +authorisation. Transport lifecycle remains operator-managed; this candidate +does not install another scheduler or silently enable a persistent tunnel. + +## Operations and evidence + +Call `discover` with instance/task IDs. It returns platform, tool versions, +acceleration and configured budgets; it is not proof of available capacity or +a running emulator. `start` also requires a stable `run_id`. Repeating an active +or terminal run returns retained state, rather than launching again. Read +`status` until ready, waiting for capacity, or stopped. A fresh attempt after a +terminal run uses a different run ID, after inspecting retained evidence. + +Ready sessions support native `tap`, bounded plain ASCII `text`, +`open_fixture` by enrolled name, `screenshot`, `report`, and `stop`. Screenshot +and run-report bytes travel to the controller, which checks the digest, +rechecks canonical assignment, and uses canonical task attachment services. +Reports identify actual Android/build/Chrome/Gboard versions and native-emulator +coverage. They do not assert physical-device behaviour. The older Android17 / +Chrome145 / Gboard17.2 evidence does not establish Chrome152 acceptance. + +`android_instance_id` in a worker configuration exposes the optional capability +in task context. The existing controller polls a task/run-specific request file +inside the coding workspace and writes responses there. The coding child gets +no Mongo credential or general host command. Requests cannot select another +task, actor or host. Private controller receipts deduplicate requests; an +interrupted effect remains uncertain and is not automatically repeated. No +additional worker or poller is installed. Task authority is rechecked by the +canonical service for each operation. Mailbox reads/writes reject symlink +redirection from the coding workspace. + +The current single-slot capacity baseline means a coding worker cannot acquire +a second slot for an emulator on its own host while retaining the first. The +intended DGX-to-Mac route uses separate host budgets. Same-host concurrent +coding/emulation needs explicit shared reservation support before activation. + +## Current delivery decision and remaining acceptance + +Merge decision: **not ready**. This branch is based on unmerged PR #682. +124 targeted tests pass, covering task/org/assignment denial, duplicate launch, command +injection rejection, and the existing shared lock's real contention/inheritance. +The normal canonical service path on DGX successfully discovered the real Mac +SDK through a temporary private SSH socket forward. A subsequent launch was +refused by capacity admission: 10.17 GiB estimated available at the final check, below +8 GiB reserve plus 6 GiB launch headroom. No emulator was launched and the guard +was not reduced to obtain a passing run. + +Still required before delivery: + +- real launch, environment read-back, native interaction and canonical attachment + read-back through the supported controller route; +- real interruption/stop and verified cleanup, including stale supervisor/socket + recovery (targeted coverage does not establish live emulator cleanup); +- demonstrable callable access from the existing coding worker/controller, + rather than only an operator's canonical service invocation; +- core package integration after its own acceptance, without inheriting this + follow-on as a gate for the core package. + +Broader architectures, Chrome152 and physical devices remain disclosed coverage +extensions. No production deployment or additional coding worker is authorised +by this candidate's test setup. diff --git a/docs/engineering/coding_agent_instances.md b/docs/engineering/coding_agent_instances.md index da83daa0..9c43e3d0 100644 --- a/docs/engineering/coding_agent_instances.md +++ b/docs/engineering/coding_agent_instances.md @@ -155,3 +155,10 @@ account. Before calling the package ready, run the task's isolated canary throug the normal agent-callable path and retain task, attachment and recipient read-back. The implementation task's current delivery explicitly prohibits starting another worker, so that acceptance step remains outstanding. + +## Optional native Android testing + +The separate [Android test-host candidate](coding_agent_android_testing.md) reuses +instance bindings and shared admission. It is optional and is not a release gate +for this core package. Read its current evidence and remaining acceptance before +claiming that a coding instance can use it autonomously. diff --git a/scripts/codex_von_android.py b/scripts/codex_von_android.py new file mode 100644 index 00000000..39668640 --- /dev/null +++ b/scripts/codex_von_android.py @@ -0,0 +1,595 @@ +"""Optional disposable Android sessions behind a fixed operator host binding. + +Only the trusted controller may invoke this executable. Requests never select +host paths, SDKs, accounts or shell commands. Same-account Unix access is not a +sandbox; separate untrusted tenants need separate prepared accounts. +""" + +from __future__ import annotations + +import argparse +import base64 +import fcntl +import hashlib +import json +import os +import platform +import re +import shutil +import signal +import socket +import subprocess +import sys +import time +from pathlib import Path + +try: + from .codex_von_instance import binding + from .codex_von_capacity import admission, available_memory_bytes +except ImportError: + from codex_von_instance import binding + from codex_von_capacity import admission, available_memory_bytes + +IDENTITY = ( + "instance_id", + "agent_id", + "delegator_id", + "organisation_id", + "host", + "unix_account", +) + + +def write(path, value): + temporary = path.with_suffix(".tmp") + temporary.write_text(json.dumps(value, indent=2)) + temporary.chmod(0o600) + temporary.replace(path) + + +def validate(spec): + cfg = spec["android"] + if [platform.system(), platform.machine()] != cfg["platform"]: + raise ValueError("Unsupported test-host platform") + if not cfg.get("capacity_enrolment_verified"): + raise ValueError("Enrol test-host consumers in shared capacity first") + if cfg["ram_mib"] < 1024 or cfg["cores"] < 1: + raise ValueError("Invalid emulator resource budget") + if cfg["cores"] > (os.cpu_count() or 1): + raise ValueError("Insufficient host CPU capacity") + if cfg["host_capacity"]["launch_headroom_bytes"] < cfg["ram_mib"] * 1024**2: + raise ValueError("Headroom must include guest RAM and calibrated host overhead") + sdk = Path(cfg["sdk_root"]).resolve() + image = (sdk / cfg["system_image"]).resolve() + image.relative_to(sdk) + for p in ( + sdk / "emulator/emulator", + sdk / "platform-tools/adb", + image / "system.img", + ): + if not p.is_file(): + raise ValueError("Prepared SDK or image unavailable") + root = Path(cfg["state_root"]) + root.mkdir(parents=True, exist_ok=True, mode=0o700) + root.chmod(0o700) + return cfg, root + + +def command(cfg, *args, binary=False): + result = subprocess.run( + [ + str(Path(cfg["sdk_root"]) / "platform-tools/adb"), + "-P", + str(cfg["adb_port"]), + "-s", + "emulator-" + str(cfg["emulator_port"]), + *args, + ], + capture_output=True, + timeout=30, + check=True, + ) + return result.stdout if binary else result.stdout.decode().strip() + + +def ports_free(cfg): + for port in (cfg["adb_port"], cfg["emulator_port"], cfg["emulator_port"] + 1): + with socket.socket() as probe: + probe.bind(("127.0.0.1", port)) + + +def probe(cfg): + sdk = Path(cfg["sdk_root"]) + acceleration = subprocess.run( + [str(sdk / "emulator/emulator"), "-accel-check"], + capture_output=True, + text=True, + timeout=30, + ) + return { + "platform": [platform.system(), platform.machine()], + "acceleration_available": acceleration.returncode == 0, + "emulator_version": subprocess.check_output( + [str(sdk / "emulator/emulator"), "-version"], text=True + ).splitlines()[0], + "adb_version": subprocess.check_output( + [str(sdk / "platform-tools/adb"), "version"], text=True + ).splitlines()[:2], + "ram_mib": cfg["ram_mib"], + "cores": cfg["cores"], + "coverage": "native_android_emulator", + "limitations": [ + "Not physical-device evidence; browser and keyboard versions must be read per run.", + "Disposable data; no enrolled Google account.", + ], + } + + +def current(root): + path = root / "session.json" + return json.loads(path.read_text()) if path.exists() else None + + +def own_session(state, request): + if ( + not state + or state["task_id"] != request["task_id"] + or state["run_id"] != request["run_id"] + ): + raise PermissionError("Session does not belong to this task and run") + + +def live_process(state): + if not state or not state.get("emulator_pid"): + return False + p = subprocess.run( + ["ps", "-p", str(state["emulator_pid"]), "-o", "command="], + text=True, + capture_output=True, + ) + return bool( + p.returncode == 0 and state["avd_name"] in p.stdout and "qemu" in p.stdout + ) + + +def stop_owned(cfg, root, state): + # The dedicated ADB ports are operator-reserved and checked before launch. + # Do not touch another session if the retained process identity does not match. + if live_process(state): + os.kill(state["emulator_pid"], signal.SIGTERM) + for _ in range(100): + if not live_process(state): + break + time.sleep(0.1) + if live_process(state): + raise RuntimeError("Owned emulator did not stop; capacity remains occupied") + adb_pid = state.get("adb_pid") + if adb_pid: + p = subprocess.run( + ["ps", "-p", str(adb_pid), "-o", "command="], text=True, capture_output=True + ) + if p.returncode == 0 and "adb" in p.stdout and str(cfg["adb_port"]) in p.stdout: + os.kill(adb_pid, signal.SIGTERM) + run_root = root / state["avd_name"] + if run_root.is_dir(): + shutil.rmtree(run_root) + (root / "control.sock").unlink(missing_ok=True) + state.update(status="stopped", stopped_at=time.time()) + write(root / "session.json", state) + + +def serve(spec, request): + cfg, root = validate(spec) + state = current(root) + own_session(state, request) + stop_request = root / ("stop-" + request["run_id"] + ".json") + with admission({"host_capacity": cfg["host_capacity"]}) as fds: + if fds is None: + state.update(status="waiting_for_capacity") + write(root / "session.json", state) + return + ports_free(cfg) + run_root = root / state["avd_name"] + run_root.mkdir(mode=0o700) + avds = run_root / "avd" + avds.mkdir() + avd = avds / (state["avd_name"] + ".avd") + avd.mkdir() + config = { + "abi.type": cfg["abi"], + "hw.cpu.arch": cfg["arch"], + "hw.cpu.ncore": cfg["cores"], + "hw.ramSize": cfg["ram_mib"], + "hw.lcd.width": 1080, + "hw.lcd.height": 1920, + "hw.lcd.density": 420, + "hw.keyboard": "no", + "hw.mainKeys": "no", + "hw.gpu.enabled": "yes", + "hw.gpu.mode": "host", + "hw.camera.back": "none", + "hw.camera.front": "none", + "image.sysdir.1": str(Path(cfg["sdk_root"]) / cfg["system_image"]) + "/", + "disk.dataPartition.size": "6G", + "showDeviceFrame": "no", + } + (avd / "config.ini").write_text( + "\n".join(f"{k}={v}" for k, v in config.items()) + ) + (avds / (state["avd_name"] + ".ini")).write_text(f"path={avd}\n") + env = { + **os.environ, + "ANDROID_SDK_ROOT": cfg["sdk_root"], + "ANDROID_AVD_HOME": str(avds), + "ANDROID_ADB_SERVER_PORT": str(cfg["adb_port"]), + } + env.pop("ANDROID_SERIAL", None) + log = (run_root / "host.log").open("ab") + adb = emulator = None + + def interrupted(*_): + raise InterruptedError("Test-host session interrupted") + + signal.signal(signal.SIGTERM, interrupted) + signal.signal(signal.SIGINT, interrupted) + try: + adb = subprocess.Popen( + [ + str(Path(cfg["sdk_root"]) / "platform-tools/adb"), + "-L", + f"tcp:127.0.0.1:{cfg['adb_port']}", + "nodaemon", + "server", + ], + env=env, + stdout=log, + stderr=log, + pass_fds=fds, + ) + state["adb_pid"] = adb.pid + write(root / "session.json", state) + time.sleep(1) + if adb.poll() is not None: + raise RuntimeError("Private ADB server failed") + emulator = subprocess.Popen( + [ + str(Path(cfg["sdk_root"]) / "emulator/emulator"), + "-avd", + state["avd_name"], + "-memory", + str(cfg["ram_mib"]), + "-cores", + str(cfg["cores"]), + "-port", + str(cfg["emulator_port"]), + "-no-snapshot", + "-no-audio", + "-no-boot-anim", + "-gpu", + "host", + ], + env=env, + stdout=log, + stderr=log, + pass_fds=fds, + ) + state["emulator_pid"] = emulator.pid + write(root / "session.json", state) + # Protect indefinite failed boot; this bounds occupancy of the test host. + deadline = time.monotonic() + cfg.get("boot_timeout_seconds", 600) + while time.monotonic() < deadline: + if stop_request.exists(): + raise InterruptedError("Task requested stop during boot") + if available_memory_bytes() < cfg["host_capacity"]["reserve_bytes"]: + raise RuntimeError("Host memory reserve crossed during boot") + if emulator.poll() is not None: + raise RuntimeError("Emulator exited before boot") + try: + if command(cfg, "shell", "getprop", "sys.boot_completed") == "1": + break + except subprocess.SubprocessError: + pass + time.sleep(2) + else: + raise RuntimeError("Emulator boot exceeded host occupancy bound") + state.update(status="ready", environment=probe(cfg)) + state["environment"].update( + android=command(cfg, "shell", "getprop", "ro.build.version.release"), + sdk=command(cfg, "shell", "getprop", "ro.build.version.sdk"), + build=command(cfg, "shell", "getprop", "ro.build.id"), + ) + for name, package in [ + ("chrome", "com.android.chrome"), + ("gboard", "com.google.android.inputmethod.latin"), + ]: + versions = command(cfg, "shell", "dumpsys", "package", package) + state["environment"][name] = re.findall( + r"versionName=([^\s]+)", versions + ) + write(root / "session.json", state) + endpoint = root / "control.sock" + endpoint.unlink(missing_ok=True) + with socket.socket(socket.AF_UNIX) as listener: + listener.bind(str(endpoint)) + endpoint.chmod(0o600) + listener.listen(2) + listener.settimeout(1) + active = time.monotonic() + while emulator.poll() is None: + if stop_request.exists(): + break + if available_memory_bytes() < cfg["host_capacity"]["reserve_bytes"]: + state["stop_reason"] = "host_memory_reserve_crossed" + break + if time.monotonic() - active > cfg.get( + "idle_timeout_seconds", 1800 + ): + break + try: + client, _ = listener.accept() + except socket.timeout: + continue + with client: + client.settimeout(5) + data = client.recv(65537) + try: + req = json.loads(data) + own_session(state, req) + result = operate(cfg, state, req) + active = time.monotonic() + client.sendall(json.dumps(result).encode()) + if req["action"] == "stop": + break + except Exception as exc: + try: + client.sendall( + json.dumps({"error": type(exc).__name__}).encode() + ) + except OSError: + pass + finally: + signal.signal(signal.SIGTERM, signal.SIG_IGN) + signal.signal(signal.SIGINT, signal.SIG_IGN) + for child in (emulator, adb): + if child and child.poll() is None: + child.terminate() + try: + child.wait(timeout=20) + except subprocess.TimeoutExpired: + child.kill() + child.wait() + (root / "control.sock").unlink(missing_ok=True) + state.update(status="stopped", stopped_at=time.time()) + write(root / "session.json", state) + shutil.rmtree(run_root) + log.close() + + +def operate(cfg, state, req): + action = req["action"] + if action == "status" or action == "stop": + return dict(state) + if action == "report": + raw = json.dumps(state, indent=2).encode() + return { + **state, + "artifact": { + "filename": "android-run-report.json", + "media_type": "application/json", + "sha256": hashlib.sha256(raw).hexdigest(), + "data_base64": base64.b64encode(raw).decode(), + }, + } + if action == "screenshot": + raw = command(cfg, "exec-out", "screencap", "-p", binary=True) + return { + **state, + "artifact": { + "filename": "android-screen.png", + "media_type": "image/png", + "sha256": hashlib.sha256(raw).hexdigest(), + "data_base64": base64.b64encode(raw).decode(), + }, + } + if action == "tap": + x, y = req["x"], req["y"] + if ( + type(x) is not int + or type(y) is not int + or not (0 <= x < 1080 and 0 <= y < 1920) + ): + raise ValueError("Tap outside display") + command(cfg, "shell", "input", "tap", str(x), str(y)) + elif action == "text": + value = req["text"] + # adb shell joins arguments; only a literal, safely quoted text value. + if ( + not isinstance(value, str) + or len(value) > 2000 + or not re.fullmatch(r"[a-zA-Z0-9 .,!?@_\-]+", value) + ): + raise ValueError( + "Native text input supports plain ASCII letters/digits/punctuation; use fixture tools for richer input" + ) + command(cfg, "shell", "input", "text", "'" + value.replace(" ", "%s") + "'") + elif action == "open_fixture": + url = cfg["fixtures"][req["fixture"]] + if not re.fullmatch(r"https?://[a-zA-Z0-9.:/_?=&%\-]+", url): + raise ValueError("Invalid operator fixture URL") + command( + cfg, + "shell", + "am", + "start", + "-a", + "android.intent.action.VIEW", + "-d", + "'" + url + "'", + "com.android.chrome", + ) + else: + raise ValueError("Unsupported Android action") + state.setdefault("actions", []).append( + {"action": action, "observed_at": time.time()} + ) + state["actions"] = state["actions"][-200:] + return {**state, "performed": action} + + +def lifecycle(spec, request, binding_path): + cfg, root = validate(spec) + identity = {key: spec[key] for key in IDENTITY} + allowed = {"action", "task_id", "run_id", "x", "y", "text", "fixture"} + if request.get("action") not in { + "discover", + "start", + "status", + "stop", + "screenshot", + "report", + "tap", + "text", + "open_fixture", + }: + raise ValueError("Unsupported Android action") + if set(request) - allowed: + raise ValueError("Unrecognised Android request fields") + if request["action"] == "discover": + return { + **identity, + "status": "available", + **probe(cfg), + "fixtures": list(cfg.get("fixtures", {})), + } + if not isinstance(request.get("task_id"), str) or not request["task_id"].startswith( + "#V#" + ): + raise ValueError("Canonical task required") + if not re.fullmatch(r"[a-zA-Z0-9_-]{1,64}", request.get("run_id", "")): + raise ValueError("Stable run ID required") + with (root / "operation.lock").open("a") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + state = current(root) + if request["action"] == "start": + if state and state["status"] not in {"stopped", "waiting_for_capacity"}: + own_session(state, request) + return { + **identity, + **state, + } # lost acknowledgement never launches twice + if ( + state + and state["run_id"] == request["run_id"] + and state["status"] == "stopped" + ): + own_session(state, request) + return { + **identity, + **state, + } # explicit new run needed after termination + digest = hashlib.sha256( + ( + spec["organisation_id"] + request["task_id"] + request["run_id"] + ).encode() + ).hexdigest()[:24] + state = { + "task_id": request["task_id"], + "run_id": request["run_id"], + "avd_name": "von_" + digest, + "status": "starting", + "started_at": time.time(), + } + write(root / "session.json", state) + log = (root / "supervisor.log").open("ab") + child = subprocess.Popen( + [ + sys.executable, + str(Path(__file__).resolve()), + "--binding", + str(binding_path), + "--serve", + ], + stdin=subprocess.PIPE, + stdout=log, + stderr=log, + start_new_session=True, + ) + child.stdin.write(json.dumps(request).encode()) + child.stdin.close() + state["supervisor_pid"] = child.pid + # Supervisor records emulator state; do not overwrite its concurrent receipt. + (root / "supervisor.pid").write_text(str(child.pid)) + log.close() + return {**identity, **state} + own_session(state, request) + endpoint = root / "control.sock" + if request["action"] == "stop": + write(root / ("stop-" + request["run_id"] + ".json"), request) + if not endpoint.exists(): + if request["action"] == "stop": + pid_path = root / "supervisor.pid" + if pid_path.exists(): + pid = int(pid_path.read_text()) + process = subprocess.run( + ["ps", "-p", str(pid), "-o", "command="], + text=True, + capture_output=True, + ) + if ( + "--serve" in process.stdout + and str(binding_path) in process.stdout + ): + os.kill(pid, signal.SIGTERM) + return {**identity, **state, "status": "stopping"} + stop_owned(cfg, root, state) + return {**identity, **state, "supervisor_ready": False} + with socket.socket(socket.AF_UNIX) as client: + client.settimeout(40) + try: + client.connect(str(endpoint)) + except (ConnectionRefusedError, FileNotFoundError): + if request["action"] == "stop": + stop_owned(cfg, root, state) + return {**identity, **state} + return { + **identity, + **state, + "supervisor_ready": False, + "recovery": "Controller disconnected; request stop before a new run", + } + client.sendall(json.dumps(request).encode()) + parts = [] + while data := client.recv(65536): + parts.append(data) + if sum(map(len, parts)) > 32 * 1024**2: + raise ValueError("Host evidence exceeds response bound") + result = json.loads(b"".join(parts)) + if "error" in result: + raise RuntimeError(result["error"]) + return {**identity, **result} + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--binding", required=True) + parser.add_argument("--serve", action="store_true") + args = parser.parse_args() + os.umask(0o077) + spec = binding(args.binding) + request = json.load(sys.stdin) + if args.serve: + try: + serve(spec, request) + except Exception as exc: + root = Path(spec["android"]["state_root"]) + state = current(root) + if state and state["run_id"] == request["run_id"]: + state.update(status="stopped", error=type(exc).__name__) + write(root / "session.json", state) + raise + else: + print(json.dumps(lifecycle(spec, request, args.binding))) + + +if __name__ == "__main__": + main() diff --git a/scripts/codex_von_android_control.py b/scripts/codex_von_android_control.py new file mode 100644 index 00000000..ee80fb99 --- /dev/null +++ b/scripts/codex_von_android_control.py @@ -0,0 +1,131 @@ +"""Existing coding controller services one task-bound Android request file. + +The coding child writes a request in its workspace; credentials remain in the +controller. Retained pending intents never replay uncertain native gestures. +""" + +import hashlib +import json +import os +import uuid +from contextlib import contextmanager +from pathlib import Path + +try: + from .codex_von_worker import write_json +except ImportError: + from codex_von_worker import write_json + + +def prepare(config, state, worktree): + if not config.get("android_instance_id"): + return None + root = Path(worktree) / ".run" / ("android-" + state["attempt"]) + root.mkdir(parents=True, exist_ok=True) + return root + + +@contextmanager +def mailbox(root): + # Walk workspace/.run/mailbox without following child-controlled symlinks. + fd = os.open(root.parents[1], os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + for component in (root.parent.name, root.name): + following = os.open( + component, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=fd + ) + os.close(fd) + fd = following + yield fd + finally: + os.close(fd) + + +def respond(root, result): + with mailbox(root) as directory: + temporary = ".response-" + uuid.uuid4().hex + fd = os.open( + temporary, + os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, + 0o600, + dir_fd=directory, + ) + with os.fdopen(fd, "w") as stream: + json.dump(result, stream) + os.rename( + temporary, "response.json", src_dir_fd=directory, dst_dir_fd=directory + ) + + +def poll(config, state, root, receipts): + with mailbox(root) as directory: + try: + fd = os.open( + "request.json", + os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK, + dir_fd=directory, + ) + except FileNotFoundError: + return + with os.fdopen(fd, "rb") as stream: + import stat + + if not stat.S_ISREG(os.fstat(stream.fileno()).st_mode): + raise ValueError("Android request must be a regular file") + raw = stream.read(65537) + if len(raw) > 65536: + raise ValueError("Android request too large") + value = json.loads(raw) + key = hashlib.sha256(raw).hexdigest() + receipt_path = receipts / (key + ".json") + if receipt_path.exists(): + result = json.loads(receipt_path.read_text()) + respond(root, result) + return + # Persist exact intent before effects. Controller interruption cannot replay + # taps or text. The agent can inspect status and request explicit recovery. + result = { + "request_sha256": key, + "status": "effect_uncertain", + "reason": "Intent retained; inspect session before issuing another gesture.", + } + write_json(receipt_path, result) + try: + if not isinstance(value, dict) or set(value) - { + "action", + "x", + "y", + "text", + "fixture", + "request_id", + }: + raise ValueError("Only native operation fields are accepted") + from src.backend.services.coding_agent_android_service import ( + coding_agent_android, + ) + + args = {k: v for k, v in value.items() if k != "request_id"} + result = { + "request_sha256": key, + "result": coding_agent_android( + instance_id=config["android_instance_id"], + task_id=state["task_id"], + run_id=state["attempt"], + **args, + ), + } + except Exception as exc: + result.update(status="failed_or_uncertain", error=type(exc).__name__) + write_json(receipt_path, result) + respond(root, result) + + +def stop(config, state): + from src.backend.services.coding_agent_android_service import coding_agent_android + + return coding_agent_android( + instance_id=config["android_instance_id"], + task_id=state["task_id"], + run_id=state["attempt"], + action="stop", + ) diff --git a/scripts/codex_von_android_transport.py b/scripts/codex_von_android_transport.py new file mode 100644 index 00000000..7774103a --- /dev/null +++ b/scripts/codex_von_android_transport.py @@ -0,0 +1,70 @@ +"""Private Unix-socket transport for an operator-bound Android host command. + +Use an SSH StreamLocal reverse forward for a test host without inbound SSH. +Both socket parents must be private; only the trusted canonical controller gets +access. No TCP listener, client-selected command, credentials or new scheduler. +""" + +import argparse +import json +import os +import socket +import sys +from pathlib import Path + + +def receive(sock, limit): + parts = [] + length = 0 + while data := sock.recv(65536): + parts.append(data) + length += len(data) + if length > limit: + raise ValueError("Transport payload exceeds bound") + return json.loads(b"".join(parts)) + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--socket", required=True) + parser.add_argument("--binding") + args = parser.parse_args() + os.umask(0o077) + path = Path(args.socket) + info = path.parent.stat() + if info.st_uid != os.geteuid() or info.st_mode & 0o077: + raise PermissionError("Transport socket directory must be owner-only") + if not args.binding: + with socket.socket(socket.AF_UNIX) as client: + client.settimeout(120) + client.connect(str(path)) + client.sendall(json.dumps(json.load(sys.stdin)).encode()) + client.shutdown(socket.SHUT_WR) + print(json.dumps(receive(client, 32 * 1024**2))) + return + try: + from .codex_von_android import binding, lifecycle + except ImportError: + from codex_von_android import binding, lifecycle + spec = binding(args.binding) + # An existing socket may be an active listener: never replace it blindly. + with socket.socket(socket.AF_UNIX) as listener: + listener.bind(str(path)) + path.chmod(0o600) + listener.listen(4) + try: + while True: + client, _ = listener.accept() + with client: + client.settimeout(120) + try: + result = lifecycle(spec, receive(client, 65536), args.binding) + except Exception as exc: + result = {"error": type(exc).__name__} + client.sendall(json.dumps(result).encode()) + finally: + path.unlink(missing_ok=True) + + +if __name__ == "__main__": + main() diff --git a/scripts/codex_von_capacity.py b/scripts/codex_von_capacity.py index 8ae27bfb..65effc8d 100644 --- a/scripts/codex_von_capacity.py +++ b/scripts/codex_von_capacity.py @@ -8,12 +8,23 @@ import fcntl import os +import platform +import re +import subprocess import stat from contextlib import contextmanager from pathlib import Path def available_memory_bytes(): + if platform.system() == "Darwin": + output = subprocess.check_output(["/usr/bin/vm_stat"], text=True) + page_size = int(re.search(r"page size of (\d+) bytes", output).group(1)) + pages = sum( + int(re.search(rf"{name}:\s+(\d+)", output).group(1)) + for name in ("Pages free", "Pages inactive", "Pages speculative") + ) + return pages * page_size for line in Path("/proc/meminfo").read_text().splitlines(): if line.startswith("MemAvailable:"): return int(line.split()[1]) * 1024 diff --git a/scripts/codex_von_worker.py b/scripts/codex_von_worker.py index 7274a5cb..6b532b34 100755 --- a/scripts/codex_von_worker.py +++ b/scripts/codex_von_worker.py @@ -18,6 +18,7 @@ import shlex import subprocess import sys +import time import uuid from contextlib import ExitStack from datetime import UTC, datetime @@ -80,6 +81,14 @@ def capability_context(config, *, read_only, worktree=None): "worktree": str(worktree) if worktree else None, "controller_runtime": config.get("runtime_evidence", {"available": False}), "public_served_revision": {"available": False, "reason": "Not probed here"}, + "android_test_host": { + "enrolled": bool(config.get("android_instance_id")), + "instance_id": config.get("android_instance_id"), + "controller_tool": "coding_agent_android", + "readiness": "Must discover through the trusted controller for the current task", + "coverage": "Native emulator, not physical-device evidence", + "limitations": "This declaration does not itself expose a tool or grant host access.", + }, "read_only_host_inspection": True, "run_can_edit_checkout": not read_only, "von_mcp": "disabled; canonical reads/effects belong to the controller", @@ -1020,6 +1029,22 @@ def launch(config, state, inputs, conversation, state_path, lock_fd): ] else: args[2:2] = ["--sandbox", "workspace-write"] + try: + from . import codex_von_android_control as android_control + except ImportError: + try: + import codex_von_android_control as android_control + except ImportError: + from scripts import codex_von_android_control as android_control + android_root = android_control.prepare(config, state, worktree) + if android_root: + context["capabilities"]["android_test_host"].update( + request_path=str(android_root / "request.json"), + response_path=str(android_root / "response.json"), + instruction="Atomically write an operation object with action and a fresh request_id; read the matching response. No task/actor/host fields. Do not repeat uncertain gestures.", + ) + write_json(context_path, context) + next_archive = time.monotonic() + ACTIVITY_CHECKPOINT_SECONDS with ( # noqa: SIM117 - streams must outlive the process context (run_dir / "events.jsonl").open("a") as events, (run_dir / "stderr.log").open("w") as errors, @@ -1038,11 +1063,38 @@ def launch(config, state, inputs, conversation, state_path, lock_fd): process.stdin.close() while True: try: - process.wait(timeout=ACTIVITY_CHECKPOINT_SECONDS) + process.wait( + timeout=2 if android_root else ACTIVITY_CHECKPOINT_SECONDS + ) break except subprocess.TimeoutExpired: - archive_checkpoint(config, state) - write_json(state_path, state) + if android_root: + try: + android_control.poll( + config, + state, + android_root, + run_dir / "android-receipts", + ) + except (OSError, ValueError): + pass # malformed/partially-written request cannot stop coding + if time.monotonic() >= next_archive: + archive_checkpoint(config, state) + write_json(state_path, state) + next_archive = time.monotonic() + ACTIVITY_CHECKPOINT_SECONDS + if android_root: + try: + write_json( + run_dir / "android-stop.json", android_control.stop(config, state) + ) + except Exception as exc: + write_json( + run_dir / "android-stop.json", + { + "error": type(exc).__name__, + "cleanup": "Unverified; test-host idle recovery remains active", + }, + ) state["finished_at"] = datetime.now(UTC).isoformat() write_json( run_dir / "exit.json", @@ -1138,9 +1190,9 @@ def apply_deployment(config, api, state, state_path, lock_fd): receipt ) if success: - result["summary"] += ( - f" Deployed and verified the public server at {commit[:12]}." - ) + result[ + "summary" + ] += f" Deployed and verified the public server at {commit[:12]}." else: result["status"] = "blocked" result["summary"] += ( @@ -1151,9 +1203,9 @@ def apply_deployment(config, api, state, state_path, lock_fd): ) if recover_only: result["status"] = "blocked" - result["summary"] += ( - " The task changed; only the interrupted deployment was reconciled." - ) + result[ + "summary" + ] += " The task changed; only the interrupted deployment was reconciled." write_json(state_path, state) @@ -1271,9 +1323,9 @@ def tick(config, api, lock_fd): state["result"]["summary"] = ( "Coding execution settings rejected: " + str(exc) ) - state["result"]["question"] = ( - "Correct this task's model or reasoning setting, then requeue it." - ) + state["result"][ + "question" + ] = "Correct this task's model or reasoning setting, then requeue it." write_json(path, state) apply_deployment(config, api, state, path, lock_fd) finish_captured(config, api, state, path) diff --git a/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py b/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py index 80fe506a..d2706cca 100644 --- a/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py +++ b/src/backend/integrations/internal_mcp/coding_agent_instance_tools.py @@ -19,14 +19,40 @@ def _instances(**kwargs): ) as exc: return make_error_response( "coding_agent_instance_unavailable", - type(exc).__name__ + ": " + str(exc) - if isinstance(exc, (PermissionError, ValueError)) - else "Instance operation unavailable; inspect operator readiness and private logs.", + ( + type(exc).__name__ + ": " + str(exc) + if isinstance(exc, (PermissionError, ValueError)) + else "Instance operation unavailable; inspect operator readiness and private logs." + ), ) def build_coding_agent_instance_tools(): return [ + MethodDefinition( + name="coding_agent_android", + handler=_android, + input_schema=Schema( + required={"instance_id": str, "task_id": str}, + optional={ + "action": str, + "run_id": str, + "x": int, + "y": int, + "text": str, + "fixture": str, + }, + allow_unknown=False, + ), + output_schema=None, + category="write", + description=( + "Optional native Android test-host capability for a current assigned task. " + "Discover actual host prerequisites; start/status/stop a disposable run with a stable run_id; " + "tap/text/open_fixture or attach a screenshot. Uses an operator-enrolled host route, " + "never creates another coding worker. Emulator evidence is not physical-device coverage." + ), + ), MethodDefinition( name="coding_agent_instances", handler=_instances, @@ -46,5 +72,28 @@ def build_coding_agent_instance_tools(): "and reasoning_effort; task overrides remain supported. Readiness requires " "actual task and recipient evidence, not a selected release or active timer." ), - ) + ), ] + + +def _android(**kwargs): + from ...services.coding_agent_android_service import coding_agent_android + + try: + return coding_agent_android(**kwargs) + except ( + PermissionError, + ValueError, + TypeError, + KeyError, + OSError, + RuntimeError, + ) as exc: + return make_error_response( + "coding_agent_android_unavailable", + ( + str(exc) + if isinstance(exc, (PermissionError, ValueError)) + else "Android operation unavailable; inspect private operator logs." + ), + ) diff --git a/src/backend/services/coding_agent_android_service.py b/src/backend/services/coding_agent_android_service.py new file mode 100644 index 00000000..b2121e18 --- /dev/null +++ b/src/backend/services/coding_agent_android_service.py @@ -0,0 +1,137 @@ +"""Task-authorised Android capability on an enrolled instance's test host. + +The existing operator instance registry binds the route. The authenticated +agent (or its delegator) must still own a live Michael/owner-created assignment. +No caller-supplied actor, organisation, host command or SDK path is accepted. +""" + +from __future__ import annotations + +import base64 +import hashlib +import json +import subprocess + +from . import coding_agent_instance_service as instances +from . import task_management_service as tasks + + +def authorised_slot(instance_id, task_id): + from ..security.access_control import get_effective_organisation_concept_id + from .organisation_membership_governance_service import _trusted_actor + from .organisation_membership_service import resolve_user_organisation_membership + + actor, error = _trusted_actor() + if error: + raise PermissionError(error["error_code"]) + org = get_effective_organisation_concept_id() + slot = instances.load_registry()["instances"].get(instance_id) + if ( + not slot + or slot["organisation_id"] != org + or actor not in {slot["agent_id"], slot["delegator_id"]} + ): + raise PermissionError("No enrolled Android capability in this actor scope") + if not resolve_user_organisation_membership(actor, org): + raise PermissionError("Current organisation membership required") + task = tasks.get_task(task_id) + if ( + not task + or task["organisation_concept_id"] != org + or task["assignee_concept_id"] != slot["agent_id"] + or task["created_by_concept_id"] != slot["delegator_id"] + ): + raise PermissionError("Current delegator-created task assignment required") + return actor, slot, task + + +def coding_agent_android( + *, + instance_id, + task_id, + action="discover", + run_id=None, + x=None, + y=None, + text=None, + fixture=None, +): + actor, slot, task = authorised_slot(instance_id, task_id) + actions = { + "discover", + "start", + "status", + "stop", + "screenshot", + "report", + "tap", + "text", + "open_fixture", + } + if action not in actions: + raise ValueError("Unsupported Android action") + if action not in {"status", "stop"} and task["status"] not in { + "pending", + "in_progress", + }: + raise PermissionError("Task is not executable") + capability = slot.get("android") + if not capability: + return { + "status": "unavailable", + "reason": "No optional Android test host enrolled", + } + request = {"action": action, "task_id": task_id, "run_id": run_id} + for key, value in (("x", x), ("y", y), ("text", text), ("fixture", fixture)): + if value is not None: + request[key] = value + argv = capability["provisioner_command"] + if ( + not isinstance(argv, list) + or not argv + or not all(isinstance(v, str) for v in argv) + ): + raise ValueError("Invalid fixed Android host command") + process = subprocess.run( + argv, input=json.dumps(request), text=True, capture_output=True + ) + if process.returncode: + raise RuntimeError( + "Android host operation failed; inspect private operator log" + ) + result = json.loads(process.stdout) + for key in ("instance_id", "agent_id", "delegator_id", "organisation_id"): + if result.get(key) != slot[key]: + raise RuntimeError("Android receipt identity mismatch") + for key in ("host", "unix_account"): + if result.get(key) != capability[key]: + raise RuntimeError("Android receipt test-host mismatch") + if action != "discover" and ( + result.get("task_id") != task_id or result.get("run_id") != run_id + ): + raise RuntimeError("Android receipt task/run mismatch") + artifact = result.pop("artifact", None) + if artifact: + raw = base64.b64decode(artifact["data_base64"], validate=True) + if hashlib.sha256(raw).hexdigest() != artifact["sha256"]: + raise RuntimeError("Android evidence checksum mismatch") + # Recheck assignment after transport and before publishing private bytes. + authorised_slot(instance_id, task_id) + result["attachment"] = tasks.add_task_attachment_bytes( + task_id, + data=raw, + filename=artifact["filename"], + actor_concept_id=actor, + media_type=artifact["media_type"], + note=json.dumps( + { + "producer": slot["agent_id"], + "run_id": run_id, + "host": result["host"], + "sha256": artifact["sha256"], + "environment": result.get("environment"), + "coverage": "native_android_emulator", + } + ), + ) + return result diff --git a/tests/backend/test_coding_agent_android.py b/tests/backend/test_coding_agent_android.py new file mode 100644 index 00000000..b5ff46b7 --- /dev/null +++ b/tests/backend/test_coding_agent_android.py @@ -0,0 +1,284 @@ +"""Boundary tests for optional task-scoped native Android testing.""" + +import base64 +import hashlib +import json +from types import SimpleNamespace + +import pytest + +from scripts import codex_von_android as host +from src.backend.services import coding_agent_android_service as service + + +@pytest.fixture +def scope(monkeypatch): + from src.backend.security import access_control + from src.backend.services import ( + organisation_membership_governance_service as governance, + ) + from src.backend.services import organisation_membership_service as membership + + slot = dict( + instance_id="test", + agent_id="#V#agent", + delegator_id="#V#owner", + organisation_id="#V#org", + android=dict( + host="mac", unix_account="test", provisioner_command=["/trusted/android"] + ), + ) + task = dict( + organisation_concept_id="#V#org", + assignee_concept_id="#V#agent", + created_by_concept_id="#V#owner", + status="in_progress", + ) + monkeypatch.setattr(governance, "_trusted_actor", lambda: ("#V#agent", None)) + monkeypatch.setattr( + access_control, "get_effective_organisation_concept_id", lambda: "#V#org" + ) + monkeypatch.setattr( + membership, + "resolve_user_organisation_membership", + lambda *_: {"role": "member"}, + ) + monkeypatch.setattr( + service.instances, "load_registry", lambda: {"instances": {"test": slot}} + ) + monkeypatch.setattr(service.tasks, "get_task", lambda *_: task) + return slot, task + + +def test_assigned_agent_route_attaches_verified_bytes(scope, monkeypatch): + slot, _ = scope + data = b"example native evidence" + result = { + **{k: v for k, v in slot.items() if k != "android"}, + "host": "mac", + "unix_account": "test", + "task_id": "#V#task", + "run_id": "run", + "status": "ready", + "artifact": dict( + filename="screen.png", + media_type="image/png", + data_base64=base64.b64encode(data).decode(), + sha256=hashlib.sha256(data).hexdigest(), + ), + } + calls = [] + monkeypatch.setattr( + service.subprocess, + "run", + lambda *a, **kw: SimpleNamespace(returncode=0, stdout=json.dumps(result)), + ) + monkeypatch.setattr( + service.tasks, + "add_task_attachment_bytes", + lambda *a, **kw: calls.append(kw) or {"attachment_id": "verified"}, + ) + response = service.coding_agent_android( + instance_id="test", task_id="#V#task", run_id="run", action="screenshot" + ) + assert response["attachment"]["attachment_id"] == "verified" + assert calls[0]["data"] == data and calls[0]["actor_concept_id"] == "#V#agent" + assert "artifact" not in response + + +@pytest.mark.parametrize( + "field,value", + [ + ("organisation_concept_id", "#V#other"), + ("created_by_concept_id", "#V#other"), + ("assignee_concept_id", "#V#other"), + ], +) +def test_other_org_owner_or_reassigned_task_denied_before_host( + scope, monkeypatch, field, value +): + scope[1][field] = value + monkeypatch.setattr( + service.subprocess, "run", lambda *a, **kw: pytest.fail("host called") + ) + with pytest.raises(PermissionError): + service.coding_agent_android( + instance_id="test", task_id="#V#task", run_id="run", action="start" + ) + + +def test_forged_identity_and_revoked_membership(scope, monkeypatch): + from src.backend.services import ( + organisation_membership_governance_service as governance, + ) + + monkeypatch.setattr( + governance, "_trusted_actor", lambda: (None, {"error_code": "untrusted_actor"}) + ) + with pytest.raises(PermissionError, match="untrusted_actor"): + service.authorised_slot("test", "#V#task") + + +def test_terminal_task_can_clean_up_but_cannot_launch(scope, monkeypatch): + scope[1]["status"] = "completed" + with pytest.raises(PermissionError, match="not executable"): + service.coding_agent_android( + instance_id="test", task_id="#V#task", run_id="run", action="start" + ) + monkeypatch.setattr( + service.subprocess, "run", lambda *a, **kw: SimpleNamespace(returncode=1) + ) + with pytest.raises(RuntimeError, match="host operation"): + service.coding_agent_android( + instance_id="test", task_id="#V#task", run_id="run", action="stop" + ) + + +def test_cross_task_session_and_shell_input_are_denied(): + with pytest.raises(PermissionError): + host.own_session( + {"task_id": "#V#a", "run_id": "r"}, {"task_id": "#V#b", "run_id": "r"} + ) + for value in ["$(touch /tmp/pwned)", "x'; echo secret", "a\nb"]: + with pytest.raises(ValueError): + host.operate({}, {}, {"action": "text", "text": value}) + + +def test_duplicate_launch_and_stopped_attempt_never_respawn(tmp_path, monkeypatch): + spec = {key: key for key in host.IDENTITY} + monkeypatch.setattr(host, "validate", lambda _: ({}, tmp_path)) + monkeypatch.setattr( + host.subprocess, "Popen", lambda *a, **kw: pytest.fail("duplicate launch") + ) + request = {"action": "start", "task_id": "#V#task", "run_id": "same"} + for phase in ["starting", "ready", "stopped"]: + host.write( + tmp_path / "session.json", + dict(task_id="#V#task", run_id="same", status=phase), + ) + assert host.lifecycle(spec, request, "binding")["status"] == phase + with pytest.raises(PermissionError): + host.lifecycle(spec, {**request, "task_id": "#V#another"}, "binding") + + +def test_missing_prerequisites_do_not_launch(tmp_path, monkeypatch): + monkeypatch.setattr(host.platform, "system", lambda: "Darwin") + monkeypatch.setattr(host.platform, "machine", lambda: "arm64") + with pytest.raises(ValueError, match="platform"): + host.validate({"android": {"platform": ["Linux", "x86_64"]}}) + + +def test_worker_request_is_fixed_to_assignment_and_deduplicated(tmp_path, monkeypatch): + from scripts import codex_von_android_control as controller + + root = tmp_path / ".run" / "child" + root.mkdir(parents=True) + receipts = tmp_path / "controller" + config = {"android_instance_id": "test"} + state = {"task_id": "#V#fixed", "attempt": "fixed-run"} + calls = [] + monkeypatch.setattr( + service, + "coding_agent_android", + lambda **kw: calls.append(kw) or {"status": "ready"}, + ) + (root / "request.json").write_text( + json.dumps({"action": "tap", "x": 1, "y": 2, "request_id": "gesture1"}) + ) + controller.poll(config, state, root, receipts) + controller.poll(config, state, root, receipts) + assert ( + len(calls) == 1 + and calls[0]["task_id"] == "#V#fixed" + and calls[0]["run_id"] == "fixed-run" + ) + (root / "request.json").write_text( + json.dumps({"action": "start", "task_id": "#V#another"}) + ) + controller.poll(config, state, root, receipts) + assert len(calls) == 1 + + +def test_worker_lost_effect_ack_is_not_repeated(tmp_path, monkeypatch): + from scripts import codex_von_android_control as controller + + root = tmp_path / ".run" / "child" + root.mkdir(parents=True) + receipts = tmp_path / "receipts" + receipts.mkdir() + raw = json.dumps({"action": "text", "text": "Hello", "request_id": "1"}).encode() + (root / "request.json").write_bytes(raw) + digest = hashlib.sha256(raw).hexdigest() + (receipts / (digest + ".json")).write_text( + json.dumps({"status": "effect_uncertain"}) + ) + monkeypatch.setattr( + service, "coding_agent_android", lambda **kw: pytest.fail("repeated effect") + ) + controller.poll( + {"android_instance_id": "test"}, + {"task_id": "#V#t", "attempt": "r"}, + root, + receipts, + ) + assert ( + json.loads((root / "response.json").read_text())["status"] == "effect_uncertain" + ) + + +def test_stale_socket_recovery_preserves_an_unrelated_process(tmp_path, monkeypatch): + import socket + + spec = {key: key for key in host.IDENTITY} + state = { + "task_id": "#V#task", + "run_id": "r", + "avd_name": "owned-avd", + "status": "ready", + "emulator_pid": 1234, + } + host.write(tmp_path / "session.json", state) + # macOS AF_UNIX addresses are shorter than pytest long temporary paths. + import tempfile + + short = tempfile.TemporaryDirectory(prefix="von-", dir="/tmp") + tmp_path = __import__("pathlib").Path(short.name) + host.write(tmp_path / "session.json", state) + endpoint = tmp_path / "control.sock" + sock = socket.socket(socket.AF_UNIX) + sock.bind(str(endpoint)) + sock.close() + monkeypatch.setattr(host, "validate", lambda _: ({}, tmp_path)) + monkeypatch.setattr( + host.subprocess, + "run", + lambda *a, **kw: SimpleNamespace(returncode=0, stdout="unrelated-editor"), + ) + monkeypatch.setattr( + host.os, "kill", lambda *a: pytest.fail("killed unrelated process") + ) + result = host.lifecycle( + spec, {"action": "stop", "task_id": "#V#task", "run_id": "r"}, "binding" + ) + assert result["status"] == "stopped" and not endpoint.exists() + + +def test_child_mailbox_cannot_redirect_controller_writes(tmp_path, monkeypatch): + from scripts import codex_von_android_control as controller + + root = tmp_path / ".run" / "child" + root.mkdir(parents=True) + secret = tmp_path / "private.txt" + secret.write_text("preserve") + (root / "response.json").symlink_to(secret) + controller.respond(root, {"status": "ready"}) + assert secret.read_text() == "preserve" + (root / "request.json").symlink_to(secret) + with pytest.raises(OSError): + controller.poll({}, {}, root, tmp_path / "receipts") + import shutil + + shutil.rmtree(root) + root.symlink_to(tmp_path, target_is_directory=True) + with pytest.raises(OSError): + controller.respond(root, {"status": "ready"}) From 118e48c3394a919dce12018aad0efa699eec65fb Mon Sep 17 00:00:00 2001 From: witbrock Date: Sun, 13 Sep 2026 18:42:20 -0400 Subject: [PATCH 3/7] Verify native Android controller route and correct ADB port launch --- .../coding_agent_android_testing.md | 74 ++++++++++++------- scripts/codex_von_android.py | 6 +- 2 files changed, 52 insertions(+), 28 deletions(-) diff --git a/docs/engineering/coding_agent_android_testing.md b/docs/engineering/coding_agent_android_testing.md index f5c78b97..445ee95e 100644 --- a/docs/engineering/coding_agent_android_testing.md +++ b/docs/engineering/coding_agent_android_testing.md @@ -81,28 +81,52 @@ a second slot for an emulator on its own host while retaining the first. The intended DGX-to-Mac route uses separate host budgets. Same-host concurrent coding/emulation needs explicit shared reservation support before activation. -## Current delivery decision and remaining acceptance - -Merge decision: **not ready**. This branch is based on unmerged PR #682. -124 targeted tests pass, covering task/org/assignment denial, duplicate launch, command -injection rejection, and the existing shared lock's real contention/inheritance. -The normal canonical service path on DGX successfully discovered the real Mac -SDK through a temporary private SSH socket forward. A subsequent launch was -refused by capacity admission: 10.17 GiB estimated available at the final check, below -8 GiB reserve plus 6 GiB launch headroom. No emulator was launched and the guard -was not reduced to obtain a passing run. - -Still required before delivery: - -- real launch, environment read-back, native interaction and canonical attachment - read-back through the supported controller route; -- real interruption/stop and verified cleanup, including stale supervisor/socket - recovery (targeted coverage does not establish live emulator cleanup); -- demonstrable callable access from the existing coding worker/controller, - rather than only an operator's canonical service invocation; -- core package integration after its own acceptance, without inheriting this - follow-on as a gate for the core package. - -Broader architectures, Chrome152 and physical devices remain disclosed coverage -extensions. No production deployment or additional coding worker is authorised -by this candidate's test setup. +## Current delivery decision and evidence + +The bounded emulator/controller capability passed live acceptance on 13 September +2026. **Publication to main still waits for core PR #682**, on which this separate +follow-on is based. No production deployment is enabled, and this task must not +become a gate for the core package. + +124 targeted tests pass, covering worker/inbox/provisioner regressions, task and +organisation denial, duplicate requests, retained uncertain effects, shared +capacity contention and symlink-safe controller mailboxes. + +The first admission was correctly refused at 10.17 GiB estimated available +versus the unchanged 14 GiB requirement (8 GiB reserve plus 6 GiB launch headroom). +After Michael requested continuation and headroom reached 19.18 GiB, a fresh +bridge attempt ran live acceptance. The ADB build rejected a hostname-qualified +`-L` endpoint; the candidate now uses its supported `-P` private port with default +loopback binding. Per-run host logs survive disposable AVD cleanup. + +Verified route: this coding agent wrote task-bound requests to the existing +controller handler on DGX, under Codex VS Code's canonical task authority; that +handler invoked the enrolled service through a private SSH StreamLocal relay to +the Mac. No second coding worker or model was started. The normal worker polling +hook is covered by tests; it was not deployed to the production DGX worker. + +Live run `native-acceptance-03` verified: + +- Android17 / SDK37 / build CE2A.260420.019, Chrome145.0.7632.218, + Gboard17.2.2.895242737, 3 GiB guest RAM and two cores; +- native Chrome navigation, text entry and a Gboard key tap; repeating the same + key-tap request returned the same receipt and inserted only one character; +- repeating start retained the same emulator PID; competing admission was denied; +- screenshot and report attachments through canonical task services; all three + downloaded byte hashes matched their canonical source hashes; +- another organisation context was denied before access to the host; +- supervisor interruption terminated emulator and ADB, removed mutable AVD data + and the control socket, and released capacity; repeating terminal start did + not relaunch; post-cleanup available memory was 17.43 GiB. + +Canonical task evidence: + +- `#V#computer_file_copy_8daf6dd0ef6c4a59a64db9eddb30819a` (initial screenshot) +- `#V#computer_file_copy_98e796f288344301b52d685109593166` (native keyboard screenshot) +- `#V#computer_file_copy_08b649de8635494fb152e7103190c67d` (environment/action report) + +The report records host operations, not browser DOM input-event telemetry or a +claim that image pasting is fixed. Chrome152 and physical-device coverage remain +outside this acceptance. For release, reconcile this follow-on against the +accepted core package and run applicable CI. Do not repeat the completed native +acceptance merely because a controller output/comment changed. diff --git a/scripts/codex_von_android.py b/scripts/codex_von_android.py index 39668640..1bd90872 100644 --- a/scripts/codex_von_android.py +++ b/scripts/codex_von_android.py @@ -224,7 +224,7 @@ def serve(spec, request): "ANDROID_ADB_SERVER_PORT": str(cfg["adb_port"]), } env.pop("ANDROID_SERIAL", None) - log = (run_root / "host.log").open("ab") + log = (root / (state["avd_name"] + "-host.log")).open("ab") adb = emulator = None def interrupted(*_): @@ -236,8 +236,8 @@ def interrupted(*_): adb = subprocess.Popen( [ str(Path(cfg["sdk_root"]) / "platform-tools/adb"), - "-L", - f"tcp:127.0.0.1:{cfg['adb_port']}", + "-P", + str(cfg["adb_port"]), "nodaemon", "server", ], From 4ff1750cce893b60b00e6d504adba4c30da9b179 Mon Sep 17 00:00:00 2001 From: witbrock Date: Thu, 17 Sep 2026 12:35:16 -0400 Subject: [PATCH 4/7] Handle canonical missing-identity lookup during first provisioning --- src/backend/services/coding_agent_instance_service.py | 5 ++++- tests/backend/test_coding_agent_instances.py | 7 ++++++- 2 files changed, 10 insertions(+), 2 deletions(-) diff --git a/src/backend/services/coding_agent_instance_service.py b/src/backend/services/coding_agent_instance_service.py index 19927702..bcecaba9 100644 --- a/src/backend/services/coding_agent_instance_service.py +++ b/src/backend/services/coding_agent_instance_service.py @@ -107,7 +107,10 @@ def _identity(slot, actor, organisation): from .organisation_membership_service import resolve_user_organisation_membership agent_id = slot["agent_id"] - concept = concept_service.get_concept_by_concept_id(agent_id) + try: + concept = concept_service.get_concept_by_concept_id(agent_id) + except concept_service.ConceptNotFoundError: + concept = None if concept is None: if not slot.get("allow_identity_creation", False): raise PermissionError( diff --git a/tests/backend/test_coding_agent_instances.py b/tests/backend/test_coding_agent_instances.py index b29bd4b9..462bee3c 100644 --- a/tests/backend/test_coding_agent_instances.py +++ b/tests/backend/test_coding_agent_instances.py @@ -274,7 +274,12 @@ def test_tool_reconcile_uses_canonical_identity_and_host_receipts( concepts, members = {}, {"#V#owner"} slot["allow_identity_creation"] = True - monkeypatch.setattr(concept_service, "get_concept_by_concept_id", concepts.get) + def lookup(concept_id): + if concept_id not in concepts: + raise concept_service.ConceptNotFoundError(concept_id) + return concepts[concept_id] + + monkeypatch.setattr(concept_service, "get_concept_by_concept_id", lookup) def create(**kwargs): assert kwargs["organisation_concept_id"] == "#V#org" From 1f1c88008232287938ab82578c3a9cb0a6e78ded Mon Sep 17 00:00:00 2001 From: witbrock Date: Thu, 17 Sep 2026 12:38:31 -0400 Subject: [PATCH 5/7] Use systemd path escaping for enrolled environment files --- scripts/codex_von_instance_schedule.py | 14 +++++++++++++- tests/backend/test_coding_agent_instances.py | 4 ++++ 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/scripts/codex_von_instance_schedule.py b/scripts/codex_von_instance_schedule.py index 33b03e4e..025ad09a 100644 --- a/scripts/codex_von_instance_schedule.py +++ b/scripts/codex_von_instance_schedule.py @@ -16,6 +16,18 @@ def _quote(value, *, command=False): ) +def _environment_path(value): + # EnvironmentFile takes one absolute path, not ExecStart's quoted argv. + # systemd otherwise treats a leading quote as a relative path and ignores it. + value = str(value) + if not Path(value).is_absolute() or any(c in value for c in ("\n", "\r", "\x00")): + raise ValueError("Invalid systemd environment path") + return "".join( + "%%" if c == "%" else "\\x%02x" % ord(c) if c in " \t\\\"'" else c + for c in value + ) + + def install(spec): settings = spec["systemd"] root = Path(spec["instance_root"]) @@ -52,7 +64,7 @@ def install(spec): for path in settings.get("environment_files", []): if not Path(path).is_file(): raise ValueError("Enrolled environment file is not prepared") - service.append("EnvironmentFile=" + _quote(path)) + service.append("EnvironmentFile=" + _environment_path(path)) timer = [ marker, "[Unit]", diff --git a/tests/backend/test_coding_agent_instances.py b/tests/backend/test_coding_agent_instances.py index 462bee3c..802303e3 100644 --- a/tests/backend/test_coding_agent_instances.py +++ b/tests/backend/test_coding_agent_instances.py @@ -337,7 +337,10 @@ def test_systemd_units_are_reused_without_clobbering_other_units( monkeypatch.setattr( schedule.subprocess, "run", lambda argv, **kw: commands.append(argv) ) + env_file = tmp_path / "instance environment" + env_file.write_text("CANARY=1\n") slot["systemd"] = { + "environment_files": [str(env_file)], "unit_directory": str(tmp_path / "units"), "memory_max_bytes": 4 * 1024**3, } @@ -347,6 +350,7 @@ def test_systemd_units_are_reused_without_clobbering_other_units( assert len(files) == 2 service_unit = next(p for p in files if p.suffix == ".service") assert "MemoryMax=4294967296" in service_unit.read_text() + assert "EnvironmentFile=" + str(tmp_path) + "/instance\\x20environment" in service_unit.read_text() assert "releases/current/scripts/codex_von_worker.py" in service_unit.read_text() assert all(command[-1] == "daemon-reload" for command in commands) service_unit.write_text("unrelated unit") From 4c344cd60e2f4c21b84115e92a3d055017d93bfa Mon Sep 17 00:00:00 2001 From: witbrock Date: Thu, 17 Sep 2026 12:42:58 -0400 Subject: [PATCH 6/7] Use format specifier in systemd path escaping --- scripts/codex_von_instance_schedule.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/codex_von_instance_schedule.py b/scripts/codex_von_instance_schedule.py index 025ad09a..0d365751 100644 --- a/scripts/codex_von_instance_schedule.py +++ b/scripts/codex_von_instance_schedule.py @@ -23,7 +23,7 @@ def _environment_path(value): if not Path(value).is_absolute() or any(c in value for c in ("\n", "\r", "\x00")): raise ValueError("Invalid systemd environment path") return "".join( - "%%" if c == "%" else "\\x%02x" % ord(c) if c in " \t\\\"'" else c + "%%" if c == "%" else f"\\x{ord(c):02x}" if c in " \t\\\"'" else c for c in value ) From dcacf5cedaeb3d7a2de97a513f90e2cb9114d63e Mon Sep 17 00:00:00 2001 From: witbrock Date: Thu, 17 Sep 2026 12:45:37 -0400 Subject: [PATCH 7/7] Record bounded live provisioning acceptance and limitations --- docs/engineering/coding_agent_instances.md | 46 ++++++++++++++++------ 1 file changed, 33 insertions(+), 13 deletions(-) diff --git a/docs/engineering/coding_agent_instances.md b/docs/engineering/coding_agent_instances.md index 5a987072..732d9d11 100644 --- a/docs/engineering/coding_agent_instances.md +++ b/docs/engineering/coding_agent_instances.md @@ -143,16 +143,36 @@ Coding children inherit the descriptor so controller death does not release admission while the child still runs. Per-service memory ceilings remain necessary; this lock alone does not constrain non-coding processes. -## Evidence boundary for the initial candidate - -Targeted tests exercise scoped discovery, forged actor rejection, live role -loss, bad settings, canonical identity/membership adapter calls, repeated and -interrupted reconciliation, retained pause, active-run protection, scheduler -ownership and two actual competing lock holders including an inherited child. -Existing worker/inbox tests cover task selection and attachment authority. -Mocked service/scheduler coverage is not a live Von deployment or another Unix -account. Before calling the package ready, run the task's isolated canary through -the normal agent-callable path and retain task, attachment and recipient -read-back. Michael authorised the temporary resource-bounded canary on 17 September 2026; -its current acceptance receipt is retained with the native implementation task. -That authorisation does not request a permanent additional production consumer. +## Acceptance and limitations + +On 17 September 2026, the normal actor-scoped tool provisioned a temporary, +independently attributed instance on the prepared DGX account. Repeated +reconciliation retained one identity and paused unit. Killing its host +provisioner after release preparation and before activation retained an honest +partial state; the same request recovered without another identity or unit. + +The resource-bounded canary used the existing DGX worker's inherited lock as +its host admission inode. Its first scheduled poll reported +`waiting_for_capacity` without selecting a task while that worker was active. +After admission became available, one actual Astra/high execution consumed a +checksum-verified canonical attachment, produced and read back its artefact, +completed the canonical task, and sent its independently attributed result to +the intended recipient. Canonical start/end, partial token usage and unknown +cost were recorded; timing and message replay did not duplicate the result. +An unrelated private file remained unavailable, wrong-assignee task selection +was excluded, and another organisation's provisioning request was denied. +The temporary timer and service were removed afterwards; the identity, task, +artefact and receipts remain audit evidence. + +This acceptance found and repaired two defects missed by initial mocks: +canonical missing-identity lookup raises `ConceptNotFoundError`, and +`EnvironmentFile` requires systemd path escaping rather than command quoting. +The installed unit's environment paths and 8 GiB memory ceiling were read back. +Targeted lifecycle, worker/inbox, capacity, release and timing checks also pass. + +The test used the existing prepared Unix account, not a separate account. +Separate-account authentication and hostile-account isolation remain unverified; +shared-account execution is not a security sandbox between users. No permanent +additional consumer, production registry, or public deployment was activated. +Native task `#V#task_agent_002b293c463ed8e1e52df23cb68f9612` retains the +substantive acceptance receipt and its environment/revision provenance.