Skip to content
Merged
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
10 changes: 6 additions & 4 deletions loopx/control_plane/collaboration/peers.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
from ...agent_registry import registered_agent_ids_for_goal
from ...thread_agent_binding import resolve_thread_agent_binding
from ..projects.registry_codec import load_project_registry
from ..runtime.file_paths import windows_extended_path
from ..content_digest import BARE_SHA256_PATTERN

PEER_INSTRUCTION = (
Expand Down Expand Up @@ -535,7 +536,9 @@ def _input_readiness_for_goal(
and Path(alias["canonical_project"]).resolve() == goal_workspace
):
selected = Path(workspace).resolve()
workspace = selected
# Resolve native addresses before checking workspace confinement, including
# deep Windows junctions that cross the workspace boundary.
workspace = windows_extended_path(selected).resolve()
result = []
for item in brief.get("inputs", []):
path = (workspace / item["ref"]).resolve()
Expand All @@ -546,9 +549,8 @@ def _input_readiness_for_goal(
try:
# 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 | getattr(os, "O_NONBLOCK", 0) | getattr(os, "O_BINARY", 0)), "rb"
) as stream:
flags = os.O_RDONLY | getattr(os, "O_NONBLOCK", 0) | getattr(os, "O_BINARY", 0)
with os.fdopen(os.open(path, flags), "rb") as stream:
if not stat.S_ISREG(os.fstat(stream.fileno()).st_mode):
raise OSError("input is not a regular file")
content = stream.read(4 * 1024 * 1024 + 1)
Expand Down
2 changes: 1 addition & 1 deletion loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -1031,7 +1031,7 @@
},
{
"site": "loopx/control_plane/collaboration/peers.py::<module>._goal::codec_read:load_project_registry#1",
"line": 52,
"line": 53,
"column": 22,
"kind": "codec_read",
"api": "load_project_registry",
Expand Down
58 changes: 58 additions & 0 deletions tests/test_peer_collaboration.py
Original file line number Diff line number Diff line change
Expand Up @@ -789,6 +789,64 @@ def test_peer_binary_artifact_preserves_crlf_and_ctrl_z_digest(scenario):
assert readiness["content_supplied"] is False


@pytest.mark.skipif(sys.platform != "win32", reason="Win32 workspace input path regression")
def test_peer_read_qualifies_long_workspace_input(scenario):
import shutil

from loopx.control_plane.runtime.file_paths import windows_extended_path

root, registry, brief, *_ = scenario
long_ref = "inputs/" + "/".join(("a" * 65, "b" * 65, "c" * 65))
long_directory = windows_extended_path(root / long_ref)
long_directory.mkdir(parents=True)
outside = root.parent / "outside-inputs"
try:
content = b"before\r\n\x1aafter\r\n\x00\xff"
artifact = long_directory / "packet.bin"
changed = long_directory / "changed.bin"
artifact.write_bytes(content)
changed.write_bytes(content)
expected = hashlib.sha256(content).hexdigest()
outside.mkdir()
(outside / "secret.bin").write_bytes(b"outside-secret")
subprocess.run(
["cmd", "/c", "mklink", "/J", str(long_directory / "junction"), str(outside)],
check=True, capture_output=True, text=True,
)
assert len(str(root / long_ref / "packet.bin")) > 260
inputs = [
{"ref": long_ref + "/packet.bin", "description": "Long input", "sha256": expected},
{"ref": long_ref + "/changed.bin", "description": "Changed input", "sha256": "0" * 64},
{"ref": long_ref + "/missing.bin", "description": "Missing input", "sha256": expected},
{"ref": long_ref + "/junction/secret.bin", "description": "Outside input",
"sha256": hashlib.sha256(b"outside-secret").hexdigest()},
]
packet = root / "long-input-brief.json"
packet.write_text(json.dumps({**brief, "inputs": inputs}), encoding="utf-8")
sent = cli(root, registry, "builder", "request", "--peer-agent-id", "reviewer",
"--operation-id", "long-workspace-input", "--brief-file", str(packet))
item = cli(root, registry, "reviewer", "read")["items"][0]
assert item["request_id"] == sent["request_id"]
readiness = item["input_readiness"]
assert [row["status"] for row in readiness] == [
"available", "changed", "unavailable", "outside_workspace"
]
assert readiness[0]["observed_sha256"] == readiness[0]["expected_sha256"] == expected
assert readiness[1]["observed_sha256"] == expected
assert readiness[3]["observed_sha256"] is None
assert all(row["content_supplied"] is False for row in readiness)
assert input_readiness(registry, "delivery", {"inputs": [{"ref": "../outside.txt"}]})[0][
"status"
] == "outside_workspace"
finally:
long_tree = windows_extended_path(root / "inputs" / ("a" * 65))
assert long_tree.resolve().is_relative_to(windows_extended_path(root).resolve())
shutil.rmtree(long_tree)
if outside.exists():
assert outside.resolve().parent == root.parent.resolve()
shutil.rmtree(outside)


def test_default_local_forwarding_keeps_revocations_after_original_delivery(external_scenario):
root, registry, brief, _store, session, _turn, parent = external_scenario
policy_path = _root(root) / "policy.json"
Expand Down
Loading