Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -252,6 +252,45 @@ host after a lost submission response before repeating the host message. The
receiver's `manager-inbox read`, decision and `report` receipts remain distinct
from host submission and from each other.

### Real-host qualification

`examples/peer-handoff-live-smoke.py` is an explicit opt-in qualification of the
existing Codex app-server adapter and the same request/return CLI. It creates
two synthetic host threads in the selected authenticated home, a disposable
Goal/registry/runtime, and a bounded artifact pinned to the checkout head and
SHA-256. It does not resume a user thread or create replacement child workers.

From the source checkout:

```bash
uv run --extra test python examples/peer-handoff-live-smoke.py --execute-real-host
```

On Windows select `--codex-bin codex.cmd` if the installed launcher needs it.
The command consumes model quota and leaves the host's own test-thread records
in that home; no authentication or session records are copied between homes.
The caller authorizes these test submissions separately from route resolution.
An unreachable historical binding keeps automatic selection ambiguous; an
explicit exact link pins the existing reviewer. The real receiver reads and
adopts the request, independently checks the artifact, and returns its exact
head/digest. A fresh requester process restores its own thread, reads the
result and acknowledges consumption. Retry recovers one request. Output contains
compact assertions, without thread links, local paths or raw conversations.

The qualification exposed two Windows blockers in this journey: private request
hashes and lock/claim suffixes exceed `MAX_PATH`, and POSIX-only input flags
prevent artifact readback. Private store/lock I/O now addresses the same physical
files using Win32 extended paths; identities, lock exclusion and storage layout
remain unchanged. Regular input files use the platform's binary/nonblocking
flags. Focused regression checks cover mutual exclusion, release, artifact
readback and one request across repeated delivery/consumption.

This qualifies this local owned-host request/adopt/return slice. It does not
qualify remote hosts, grant message permission to an arbitrary App task, transfer
a lease, or close the overall R2/R3 collaboration acceptance. Route previews
continue to report `host_delivery: not_attempted`; the smoke's explicit host
submission and receiver receipts are separate evidence.

## Target Identity Pinning

A bounded wait, or the readback that a delivery produced a turn, must be pinned
Expand Down
238 changes: 238 additions & 0 deletions examples/peer-handoff-live-smoke.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,238 @@
#!/usr/bin/env python3
"""Opt-in real Codex peer handoff using only disposable LoopX state.

Run from the source checkout with an authenticated Codex home. This creates two
synthetic host threads and consumes model quota; it never resumes a user thread.
Only compact assertions are printed. Host transcripts stay in the selected home.
"""
from __future__ import annotations

import argparse
import hashlib
import json
import os
import subprocess
import sys
import tempfile
from pathlib import Path

from loopx.chat_agent import CodexChatAgentSession
from loopx.control_plane.effect_runtime import restart_effect_runtime


def tool(name, description, properties=None):
fields = properties or {}
return {"name": name, "description": description, "inputSchema": {
"type": "object", "properties": fields, "required": list(fields),
"additionalProperties": False,
}}


def qualify(workspace: Path, *, codex_bin: str, codex_home: Path) -> dict:
"""The host owns thread submission; LoopX owns request/read/return receipts."""
runtime = workspace / "runtime"
registry = workspace / "registry.json"
head = subprocess.check_output(
["git", "rev-parse", "HEAD"], text=True
).strip()
artifact = workspace / "review-packet.json"
artifact.write_text(json.dumps({"head": head, "demand": 15,
"allocated": 13, "reserved": 2}), encoding="utf-8")
digest = hashlib.sha256(artifact.read_bytes()).hexdigest()
brief = workspace / "brief.json"
brief.write_text(json.dumps({
"schema_version": "collaboration_brief_v0",
"purpose": "Independently review the synthetic allocation artifact",
"context": f"Review exactly checkout head {head}; do not substitute another peer.",
"constraints": ["No child Agents", "No shell commands", "No external writes"],
"inputs": [{"ref": artifact.name, "description": "Synthetic review packet",
"sha256": digest}],
"acceptance": ["Allocated plus reserved equals demand", "Head and digest match"],
"return_requirement": "Return the verdict, exact head and artifact digest to requester",
}), encoding="utf-8")
env = {**os.environ, "LOOPX_CODEX_HOMES": str(codex_home)}

def cli(agent, action, *args, ok=True):
process = subprocess.run([sys.executable, "-m", "loopx.cli",
"--registry", str(registry), "--runtime-root", str(runtime),
"--format", "json", "manager-inbox", action,
"--goal-id", "peer-qualification", "--agent-id", agent, *args],
env=env, cwd=workspace, text=True, capture_output=True, timeout=45)
value = json.loads(process.stdout)
assert process.returncode == (0 if ok else 1), value
return value

events, calls = [], []
sessions = []

def observe(method, value):
if method == "item/completed":
item = value.get("item") or {}
events.append(item.get("type"))

def start(agent, tools, *, resume=None):
session = CodexChatAgentSession.start(
codex_bin=codex_bin, codex_home=codex_home, work_dir=workspace,
goal_id="peer-qualification", objective=f"Synthetic {agent} qualification",
execution_mode=True, sandbox="read-only", resume_thread_id=resume,
dynamic_tools=tools, isolate_process_tree=True, hard_timeout_sec=180,
host_config={"features.multi_agent": False},
)
sessions.append(session)
return session

review_tools = [
tool("peer_review_read", "Read your scoped Inbox request and its exact artifact"),
tool("peer_review_adopt", "Adopt the request after reading it"),
tool("peer_review_report", "Return the independent review to its requester", {
"head": {"type": "string"}, "sha256": {"type": "string"},
"verdict": {"type": "string", "enum": ["accept", "reject"]},
}),
]
requester_tools = [
tool("peer_return_read", "Read the review returned to this requester"),
tool("peer_return_consume", "Acknowledge the returned result after reading it"),
]
try:
reviewer = start("reviewer", review_tools)
requester = start("requester", requester_tools)
reviewer_id, requester_id = reviewer.thread_id, requester.thread_id
# These are new test sessions. Finishing a real Turn makes the production
# local-store observer readable without fabricating its SQLite records.
for session in (reviewer, requester):
session.send("Reply fixture-ready. Do not use tools or create child Agents.",
on_event=observe)
reviewer.close()
requester.close()
registry.write_text(json.dumps({"goals": [{"id": "peer-qualification",
"repo": str(workspace), "coordination": {
"registered_agents": ["requester", "reviewer"],
"thread_agent_bindings": [
{"agent_id": "reviewer", "host_surface": "codex-app",
"thread_id": reviewer_id},
{"agent_id": "reviewer", "host_surface": "codex-app",
"thread_id": "historical-unreachable-fixture"},
{"agent_id": "requester", "host_surface": "codex-app",
"thread_id": requester_id},
],
}}]}), encoding="utf-8")
common = ("--peer-agent-id", "reviewer", "--operation-id", "exact-head-review",
"--brief-file", str(brief), "--require-host-route")
refused = cli("requester", "request", *common, ok=False)
assert "ambiguous" in refused["error"]
selected = (*common, "--peer-thread-link", f"codex://threads/{reviewer_id}")
sent = cli("requester", "request", *selected)
rid = sent["request_id"]
assert sent["host_delivery"]["status"] == "not_attempted"
assert sent["host_delivery"]["thread_id"] == reviewer_id
assert cli("requester", "request", *selected)["replayed"]
# Resume only the pinned test thread through the owning host. This is an
# explicitly authorized test submission, not authority inferred from a route.
reviewer = start("reviewer", review_tools, resume=reviewer_id)

def review_handler(name, arguments, identity):
assert identity["thread_id"] == reviewer_id
calls.append(name)
if name == "peer_review_read":
inbox = cli("reviewer", "read")
row = next(item for item in inbox["items"] if item["request_id"] == rid)
assert row["brief"]["inputs"][0]["sha256"] == digest
assert hashlib.sha256(artifact.read_bytes()).hexdigest() == digest
return {"ok": True, "request": row, "artifact": json.loads(artifact.read_text()),
"sha256": digest}
if name == "peer_review_adopt":
assert "peer_review_read" in calls
return cli("reviewer", "acknowledge", "--request-id", rid,
"--decision", "adopt", "--reason", "Independent exact-head review")
if name == "peer_review_report":
assert "peer_review_adopt" in calls
assert arguments == {"head": head, "sha256": digest, "verdict": "accept"}
return cli("reviewer", "report", "--request-id", rid,
"--phase", "conclusion", "--reply-text", f"ACCEPT head={head} sha256={digest}")
raise ValueError("unsupported qualification tool")

reviewer.bound_tool_handler = review_handler
reviewer.send(sent["host_delivery"]["message"] +
" Use peer_review_read, inspect the arithmetic independently, then "
"peer_review_adopt and peer_review_report. Use only these three tools.",
on_event=observe)
reviewer.close()
# A fresh process restores the requester's original identity for result
# consumption; neither retry nor restart creates another peer request.
replay = cli("requester", "request", *selected)
assert replay["replayed"] and replay["request_id"] == rid
requester = start("requester", requester_tools, resume=requester_id)

def return_handler(name, arguments, identity):
assert identity["thread_id"] == requester_id
calls.append(name)
if name == "peer_return_read":
result = cli("requester", "read")
item = next(item for item in result["peer_returns"]["items"]
if item["request_id"] == rid)
assert item["text"] == f"ACCEPT head={head} sha256={digest}"
return {"ok": True, "result": item}
if name == "peer_return_consume":
assert "peer_return_read" in calls
return cli("requester", "acknowledge-return", "--request-id", rid)
raise ValueError("unsupported qualification tool")

requester.bound_tool_handler = return_handler
requester.send("The existing peer has returned its review. Use peer_return_read "
"then peer_return_consume; summarize its exact head and digest. "
"Use only these two tools.", on_event=observe)
required = {entry["name"] for entry in review_tools + requester_tools}
assert required <= set(calls), "real peers did not complete the exchange"
assert not cli("requester", "read").get("peer_returns", {}).get("items", [])
assert not {"collabAgentToolCall", "commandExecution"} & set(events)
status = cli("reviewer", "status", "--request-id", rid)
return {"ok": True, "host": "codex_app_server", "peers": 2,
"exact_thread_resumed": True, "request_replayed": True,
"receiver_adopted": True, "result_returned_and_consumed": True,
"replacement_workers": 0, "host_delivery_preview": "not_attempted",
"tracked_requests": len(status["rows"])}
finally:
for session in sessions:
session.close()


def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--execute-real-host", action="store_true",
help="Authorize two synthetic Codex sessions and model usage")
parser.add_argument("--codex-bin", default="codex")
parser.add_argument("--codex-home", type=Path,
default=Path(os.environ.get("CODEX_HOME") or "~/.codex").expanduser())
args = parser.parse_args()
if not args.execute_real_host:
parser.error("--execute-real-host is required; this smoke uses a real authenticated host")
with tempfile.TemporaryDirectory(prefix="lxp-") as folder:
# CLI reads start a reusable typed runtime whose Windows cwd keeps the
# workspace open. Give this smoke its own locator, then stop that owner
# before deleting the workspace; never stop a shared user's runtime.
runtime_temp = Path(folder) / "tmp"
runtime_temp.mkdir()
previous_temp = tempfile.tempdir
previous_env = {key: os.environ.get(key) for key in ("TMPDIR", "TEMP", "TMP")}
try:
tempfile.tempdir = str(runtime_temp)
os.environ.update({key: str(runtime_temp) for key in previous_env})
try:
result = qualify(Path(folder), codex_bin=args.codex_bin,
codex_home=args.codex_home.resolve())
finally:
restart = restart_effect_runtime()
if restart["status"] == "shutdown_pending":
raise RuntimeError("isolated typed runtime shutdown did not complete")
finally:
tempfile.tempdir = previous_temp
for key, value in previous_env.items():
if value is None:
os.environ.pop(key, None)
else:
os.environ[key] = value
print(json.dumps(result, sort_keys=True))


if __name__ == "__main__":
main()
5 changes: 4 additions & 1 deletion loopx/control_plane/collaboration/inbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
from pathlib import Path
from typing import TYPE_CHECKING, Any
from ...file_lock import exclusive_file_lock
from ..runtime.file_paths import windows_extended_path
from ..content_digest import BARE_SHA256_PATTERN, ENVELOPED_SHA256_PATTERN
from ..todos.contract import TODO_ID_PATTERN

Expand All @@ -40,7 +41,9 @@ def _hash(value: Any) -> str:

def _root(runtime_root: Path) -> Path:
"""Retain the shipped storage address; Agent topology is not encoded in it."""
return runtime_root / ".local" / "manager-context"
# Full request hashes plus lock sidecars can exceed MAX_PATH even in an
# ordinary workspace. Keep extended syntax inside the private store.
return windows_extended_path(runtime_root / ".local" / "manager-context")


def _write(path: Path, value: dict) -> None:
Expand Down
2 changes: 1 addition & 1 deletion loopx/control_plane/collaboration/peers.py
Original file line number Diff line number Diff line change
Expand Up @@ -544,7 +544,7 @@ def input_readiness(
# Nonblocking open plus fstat prevents a FIFO/device reference
# from hanging the worker's entire Inbox read.
with os.fdopen(
os.open(path, os.O_RDONLY | os.O_NONBLOCK), "rb"
os.open(path, os.O_RDONLY | getattr(os, "O_NONBLOCK", 0) | getattr(os, "O_BINARY", 0)), "rb"
) as stream:
if not stat.S_ISREG(os.fstat(stream.fileno()).st_mode):
raise OSError("input is not a regular file")
Expand Down
18 changes: 18 additions & 0 deletions loopx/control_plane/runtime/file_paths.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
"""Native filesystem addresses without Goal discovery or routing dependencies."""

from __future__ import annotations

import os
from pathlib import Path


def windows_extended_path(path: Path) -> Path:
"""Address the same Windows file beyond MAX_PATH; leave other hosts alone."""
if os.name != "nt":
return path
address = os.path.abspath(path)
if address.startswith("\\\\?\\"):
return Path(address)
if address.startswith("\\\\"):
return Path("\\\\?\\UNC\\" + address[2:])
return Path("\\\\?\\" + address)
16 changes: 12 additions & 4 deletions loopx/file_lock.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
from typing import Any, Iterator, TextIO
from uuid import uuid4

from .control_plane.runtime.file_paths import windows_extended_path

try: # pragma: no cover - exercised on POSIX hosts in integration smokes.
fcntl: Any = importlib.import_module("fcntl")
except ImportError: # pragma: no cover
Expand Down Expand Up @@ -121,7 +123,7 @@ def _policy(value: LockAcquisitionPolicy | str) -> LockAcquisitionPolicy:


def _lock_path(path: Path) -> Path:
return path.with_name(f"{path.name}.lock")
return windows_extended_path(path.with_name(f"{path.name}.lock"))


def _open_lock_descriptor(path: Path, *, flags: int) -> int:
Expand Down Expand Up @@ -159,8 +161,14 @@ def lock_incident_path(path: Path) -> Path:


def _lock_id(path: Path) -> str:
resolved = str(path.expanduser().resolve(strict=False)).encode("utf-8")
return hashlib.sha256(resolved).hexdigest()[:16]
address = str(path.expanduser().resolve(strict=False))
if os.name == "nt":
# The address syntax must not create a second diagnostic lock identity.
if address.startswith("\\\\?\\UNC\\"):
address = "\\\\" + address[8:]
elif address.startswith("\\\\?\\"):
address = address[4:]
return hashlib.sha256(address.encode("utf-8")).hexdigest()[:16]


def _identity(
Expand Down Expand Up @@ -602,7 +610,7 @@ def try_exclusive_file_lock(


def _effect_mutation_lock_path(path: Path) -> Path:
return Path(f"{path}{EFFECT_MUTATION_LOCK_SUFFIX}")
return windows_extended_path(Path(f"{path}{EFFECT_MUTATION_LOCK_SUFFIX}"))


def _effect_mutation_claim_path(path: Path, token: str) -> Path:
Expand Down
Loading
Loading