diff --git a/docs/integrations/ego-source-reader.md b/docs/integrations/ego-source-reader.md index 2c044aa79c..43839b478f 100644 --- a/docs/integrations/ego-source-reader.md +++ b/docs/integrations/ego-source-reader.md @@ -1,13 +1,18 @@ -# Rendered public-source reading with an existing Ego Page +# Rendered public-source reading with Ego The optional `loopx.extensions.ego_source_reader` stdio MCP adapter gives an existing Codex host narrow `read_public_url` and `read_public_image` tools when its model shell cannot -reach Ego's local bootstrap. It reuses an installed Ego browser and an existing -TaskSpace/Page. It does not replace the Chat Session owner or create a browser, +reach Ego's local bootstrap. It reuses an installed Ego browser with a reserved +Page. It does not replace the Chat Session owner or create a browser, model thread, background service, material catalog or permission authority. -Operator setup is explicit. Reserve an existing Page for this MCP process and -allow only the public-source origins needed for the task. Do not reserve a Page +Operator setup is explicit. Choose `auto` to lazily create one TaskSpace and its +initial `p1` per MCP process. Ego's named factory can reuse an existing space, +so the adapter generates a unique host nonce once and keeps that name throughout +its lifecycle and confirmed-missing recovery. This avoids expired fixed ids and +keeps concurrent hosts on separate Pages. Alternatively, reserve an existing numeric TaskSpace +and Page for this process. Allow only the public-source origins needed for the +task. Do not reserve a Page shared with another process or grant a private account/admin origin. Navigation may use the existing browser's session; origin configuration does not prove that every page on that origin is public. Follow the browser's installed skill for @@ -25,18 +30,32 @@ tool_timeout_sec = 40 [mcp_servers.loopx_ego_source_read.env] LOOPX_EGO_READ_BIN = "/absolute/path/to/installed/ego-browser" -LOOPX_EGO_READ_TASK_SPACE = "7" -LOOPX_EGO_READ_PAGE = "p2" +LOOPX_EGO_READ_TASK_SPACE = "auto" +LOOPX_EGO_READ_PAGE = "p1" LOOPX_EGO_READ_ORIGINS = "https://example.com,https://www.example.org" ``` Use a supported LoopX installation containing this module. Restart an idle host through its existing service path and resume the original Session. Do not change its sandbox, approval policy, workspace grants or authentication to make the -tool work. No new dependency is needed beyond LoopX's existing MCP dependency. +tool work. The Python environment needs the optional `loopx[ego-source-reader]` +extra (`mcp==1.28.1`); the base LoopX CLI has no Python runtime dependencies. +Keep this environment outside PATH when another installation owns `loopx`. Disable by removing only this MCP entry and restarting that idle host. Keep a private configuration backup and the installed/source revision for rollback. +In `auto` mode, URL/configuration/origin validation runs before creation. +The process reuses its space for text and image calls, and replaces it once only +when Ego explicitly reports `task space not found`. Other browser errors, +verification walls and user-control stops do not create replacements. An +ambiguous creation receipt fails closed until the operator inspects/restarts the +host. On normal MCP shutdown or SIGTERM, it finishes only its own created, +still-agent-owned space. SIGTERM cleanup can complete while the stdio server +still waits for its host to close stdin; callers should also close the pipe when +stopping the process. Configured numeric spaces are never finished by the +adapter. Shutdown failures may require operator cleanup; a killed process cannot +guarantee cleanup. No login/profile selection or browser-control tool is exposed. + The tool accepts an HTTPS URL, checks the configured origin before navigation, and uses WHATWG URL normalization for the target before checking the exact resulting URL before DOM extraction. Equivalent dot segments and query escaping diff --git a/loopx/extensions/ego_source_reader.py b/loopx/extensions/ego_source_reader.py index 350cc0a0c1..9eef04cb86 100644 --- a/loopx/extensions/ego_source_reader.py +++ b/loopx/extensions/ego_source_reader.py @@ -1,4 +1,4 @@ -"""Opt-in rendered-source MCP adapter for an existing, reserved Ego Page. +"""Opt-in rendered-source MCP adapter for a reserved Ego Page. This transport owns no Session, grant, material store or model runner. Operator configuration selects the browser endpoint and origins; tool input selects only @@ -11,11 +11,14 @@ import json import os import re +import signal import subprocess import threading +import time import tempfile import struct -from dataclasses import dataclass +import uuid +from dataclasses import dataclass, replace from pathlib import Path from urllib.parse import urlsplit @@ -27,6 +30,7 @@ MAX_IMAGE_ITEMS = 128 TIMEOUT_SECONDS = 30 MARKER = "LOOPX_PUBLIC_SOURCE:" +SPACE_MARKER = "LOOPX_READER_SPACE:" _READ_LOCK = threading.Lock() @@ -50,7 +54,7 @@ def _url(value: str) -> tuple[str, str]: @dataclass(frozen=True) class ReaderConfig: executable: str - task_space: int + task_space: int | None page: str origins: frozenset[str] @@ -59,10 +63,12 @@ def from_environment(cls) -> ReaderConfig: executable = Path(os.environ["LOOPX_EGO_READ_BIN"]).expanduser() if not executable.is_absolute() or not executable.is_file(): raise ValueError("configure an installed executable") - space = int(os.environ["LOOPX_EGO_READ_TASK_SPACE"]) - page = os.environ["LOOPX_EGO_READ_PAGE"] - if space <= 0 or not re.fullmatch(r"p[1-9][0-9]*", page): - raise ValueError("configure an existing TaskSpace and Page label") + setting = os.environ["LOOPX_EGO_READ_TASK_SPACE"] + space = None if setting == "auto" else int(setting) + page = os.environ.get("LOOPX_EGO_READ_PAGE", "p1") + if ((space is not None and space <= 0) or not re.fullmatch(r"p[1-9][0-9]*", page) + or (space is None and page != "p1")): + raise ValueError("configure an existing Page or auto with p1") origins = set() for entry in os.environ["LOOPX_EGO_READ_ORIGINS"].split(","): canonical, origin = _url(entry.strip()) @@ -72,6 +78,74 @@ def from_environment(cls) -> ReaderConfig: return cls(str(executable.resolve(strict=True)), space, page, frozenset(origins)) +class _OwnedSpace: + """One lazily created space per MCP process; never owns configured spaces.""" + + def __init__(self) -> None: + # Ego's named factory reuses existing agent-owned spaces. A stable + # nonce belongs to this MCP host, not to all hosts of this provider. + self.name = f"LoopX public-source reader {uuid.uuid4().hex}" + self.space: int | None = None + self.executable: str | None = None + self.creation_attempted = False + + def resolve(self, config: ReaderConfig, *, deadline: float | None = None) -> ReaderConfig: + if config.task_space is not None: + return config + if self.executable is not None and self.executable != config.executable: + raise ValueError("reader executable changed") + if self.space is None: + # A lost creation receipt is ambiguous: don't create another space + # on the next tool call. An operator must inspect/restart the host. + if self.creation_attempted: + raise ValueError("reader space creation outcome unknown") + self.creation_attempted = True + self.executable = config.executable + script = (f"const t=await taskSpace({json.dumps(self.name)});" + f"console.log({json.dumps(SPACE_MARKER)}+JSON.stringify({{id:t.spaceId}}));") + result = _run(config.executable, script, deadline=deadline) + values = [line[len(SPACE_MARKER):] for line in (result.stdout + "\n" + result.stderr).splitlines() + if line.startswith(SPACE_MARKER)] + if result.returncode or len(values) != 1: + raise ValueError("reader space creation failed") + value = json.loads(values[0]) + if not isinstance(value, dict) or type(value.get("id")) is not int or value["id"] <= 0: + raise ValueError("invalid reader space receipt") + self.space = value["id"] + return replace(config, task_space=self.space) + + def forget_closed(self) -> None: + self.space = None + self.creation_attempted = False + + def close(self) -> None: + if self.space is None or self.executable is None: + return + space, self.space = self.space, None + # Only this process's created space, and only while still agent-owned. + # Never claim/take over a space after the user or another owner stops it. + script = (f"const t=await taskSpace({space});" + "if(t.ownership==='agent')await t.finish({keep:[]});") + try: + _run(self.executable, script) + except (OSError, UnicodeError, subprocess.TimeoutExpired): + pass + + +_OWNED_SPACE = _OwnedSpace() + + +def _run(executable: str, script: str, *, deadline: float | None = None) -> subprocess.CompletedProcess[str]: + timeout = TIMEOUT_SECONDS if deadline is None else deadline - time.monotonic() + if timeout <= 0: + raise subprocess.TimeoutExpired(executable, TIMEOUT_SECONDS) + return subprocess.run( + [executable, "nodejs", "-e", script], + stdin=subprocess.DEVNULL, capture_output=True, text=True, encoding="utf-8", + timeout=timeout, check=False, + ) + + def _navigation(config: ReaderConfig, url: str) -> str: # Use the browser's URL rules before navigation, including dot segments and # query escaping. The fixed operator-owned script, not page data, supplies @@ -80,7 +154,10 @@ def _navigation(config: ReaderConfig, url: str) -> str: f"const requestedUrl={json.dumps(url)};const target=new URL(requestedUrl);target.hash='';" f"const origins={json.dumps(sorted(config.origins))}.map(o=>new URL(o).origin);" "if(!origins.includes(target.origin))throw new Error('source_origin_not_authorized');" - f"const t=await taskSpace({config.task_space});const p=t.page({json.dumps(config.page)});" + f"let t;try{{t=await taskSpace({config.task_space});}}catch(e){{" + "if(/task space not found/i.test(String(e?.message)))" + f"console.log({json.dumps(SPACE_MARKER)}+JSON.stringify({{closed:true}}));throw e;}}" + f"const p=t.page({json.dumps(config.page)});" "await p.goto(target.href);" ) @@ -253,17 +330,21 @@ def _read(url: str, image_index: int | None = None, screenshot_path: str = "") - if origin not in config.origins: return {"ok": False, "error": "source_origin_not_authorized"} # Concurrent calls within this MCP process do not navigate the reserved Page - # over one another. Separate processes must reserve separate existing Pages. + # over one another. Auto mode gives each process a distinct owned space. if not _READ_LOCK.acquire(blocking=False): return {"ok": False, "error": "source_reader_busy"} try: - script = (_script(config, canonical) if image_index is None else - _image_script(config, canonical, image_index, screenshot_path)) - result = subprocess.run( - [config.executable, "nodejs", "-e", script], - stdin=subprocess.DEVNULL, capture_output=True, text=True, encoding="utf-8", - timeout=TIMEOUT_SECONDS, check=False, - ) + deadline = time.monotonic() + TIMEOUT_SECONDS + for attempt in range(2): + resolved = _OWNED_SPACE.resolve(config, deadline=deadline) + script = (_script(resolved, canonical) if image_index is None else + _image_script(resolved, canonical, image_index, screenshot_path)) + result = _run(config.executable, script, deadline=deadline) + if (attempt == 0 and config.task_space is None and result.returncode + and SPACE_MARKER + '{"closed":true}' in (result.stdout + "\n" + result.stderr).splitlines()): + _OWNED_SPACE.forget_closed() + continue + break if result.returncode: return {"ok": False, "error": "browser_read_failed", "exit_code": result.returncode} @@ -274,6 +355,8 @@ def _read(url: str, image_index: int | None = None, screenshot_path: str = "") - return {"ok": False, "error": "browser_read_timeout"} except (OSError, UnicodeError): return {"ok": False, "error": "browser_read_unavailable"} + except (TypeError, ValueError): + return {"ok": False, "error": "source_reader_space_unavailable"} finally: _READ_LOCK.release() @@ -319,7 +402,15 @@ def main() -> None: readOnlyHint=True, destructiveHint=False, idempotentHint=True, openWorldHint=True, ))(read_public_image) - server.run(transport="stdio") + previous = signal.getsignal(signal.SIGTERM) + def terminate(_signum: int, _frame: object) -> None: + raise SystemExit(0) + signal.signal(signal.SIGTERM, terminate) + try: + server.run(transport="stdio") + finally: + _OWNED_SPACE.close() + signal.signal(signal.SIGTERM, previous) if __name__ == "__main__": diff --git a/pyproject.toml b/pyproject.toml index 2f1d485b54..7427b6e474 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -23,6 +23,9 @@ Issues = "https://github.com/loopx-project/loopx/issues" Changelog = "https://github.com/loopx-project/loopx/releases" [project.optional-dependencies] +ego-source-reader = [ + "mcp==1.28.1", +] deepseek-harness = [ "deepseek-harness-sdk==0.1.5rc1", ] diff --git a/tests/test_ego_source_reader.py b/tests/test_ego_source_reader.py index bf466361aa..496de4e345 100644 --- a/tests/test_ego_source_reader.py +++ b/tests/test_ego_source_reader.py @@ -10,6 +10,11 @@ URL = "https://example.com/article?q=1" +@pytest.fixture(autouse=True) +def isolated_reader_space(monkeypatch): + monkeypatch.setattr(reader, "_OWNED_SPACE", reader._OwnedSpace()) + + @pytest.fixture def configured(monkeypatch, tmp_path): executable = tmp_path / "ego-browser" @@ -49,7 +54,7 @@ def run(args, **kwargs): assert result["images_read"] is False and result["image_count"] == 2 args, kwargs = calls[0] assert args[:3] == [str(configured), "nodejs", "-e"] - assert kwargs["stdin"] == subprocess.DEVNULL and kwargs["timeout"] == 30 + assert kwargs["stdin"] == subprocess.DEVNULL and 0 < kwargs["timeout"] <= 30 assert "shell" not in kwargs assert "does not prove article" in result["limitations"] @@ -108,6 +113,167 @@ def test_bad_operator_configuration_fails_closed(configured, monkeypatch, key, v assert reader.read_public_url(URL)["error"] == "source_reader_not_configured" +def auto_config(monkeypatch): + monkeypatch.setenv("LOOPX_EGO_READ_TASK_SPACE", "auto") + monkeypatch.setenv("LOOPX_EGO_READ_PAGE", "p1") + + +def space_response(space, *, stderr=False): + output = reader.SPACE_MARKER + json.dumps({"id": space}) + return subprocess.CompletedProcess([], 0, "" if stderr else output, output if stderr else "") + + +@pytest.mark.parametrize("stderr", [False, True]) +def test_auto_space_is_lazy_and_reused_for_text_and_image(configured, monkeypatch, tmp_path, stderr): + auto_config(monkeypatch) + calls = [] + def run(_executable, script, **_kwargs): + calls.append(script) + if 'taskSpace("LoopX public-source reader ' in script: + return space_response(19, stderr=stderr) + if "screenshot(" in script: + path = tmp_path / "image.png" + path.write_bytes(png()) + return response({"url": URL, "index": 0, "alt": "figure"}) + return response(extraction()) + monkeypatch.setattr(reader, "_run", run) + assert reader.read_public_url("https://other.test/private")["error"] == "source_origin_not_authorized" + assert not calls + assert reader.read_public_url(URL)["ok"] + assert reader.read_public_url(URL)["ok"] + assert reader._read(URL, 0, str(tmp_path / "image.png"))["ok"] + assert len(calls) == 4 + assert all("taskSpace(19)" in s for s in calls[1:]) + + +def test_separate_reader_process_state_owns_separate_spaces(configured, monkeypatch): + auto_config(monkeypatch) + config = reader.ReaderConfig.from_environment() + # Model the real factory's name-based reuse, not predetermined distinct ids. + spaces = {} + def run(_executable, script, **_kwargs): + name = json.loads(script.split("taskSpace(", 1)[1].split(");", 1)[0]) + return space_response(spaces.setdefault(name, 19 + len(spaces))) + monkeypatch.setattr(reader, "_run", run) + first, second = reader._OwnedSpace(), reader._OwnedSpace() + assert first.resolve(config).task_space == 19 + assert second.resolve(config).task_space == 20 + assert first.resolve(config).task_space == 19 + assert first.name != second.name and len(spaces) == 2 + first.forget_closed() + first.resolve(config) + assert len(spaces) == 2 # The owner identity remains stable during recovery. + + +@pytest.mark.parametrize("failure", ["user_control", "timeout", "ambiguous_creation"]) +def test_errors_do_not_create_replacement_spaces(configured, monkeypatch, failure): + auto_config(monkeypatch) + calls = [] + def run(executable, script, **_kwargs): + calls.append(script) + if 'taskSpace("LoopX public-source reader ' in script: + if failure == "ambiguous_creation": + raise subprocess.TimeoutExpired(executable, 30) + return space_response(19) + if failure == "timeout": + raise subprocess.TimeoutExpired(executable, 30) + return subprocess.CompletedProcess([], 1, "", "user took control") + monkeypatch.setattr(reader, "_run", run) + assert not reader.read_public_url(URL)["ok"] + assert not reader.read_public_url(URL)["ok"] + assert sum('taskSpace("LoopX public-source reader ' in s for s in calls) == 1 + + +@pytest.mark.parametrize("stderr", [False, True]) +def test_only_confirmed_missing_owned_space_is_replaced_once(configured, monkeypatch, stderr): + auto_config(monkeypatch) + calls = [] + def run(_executable, script, **_kwargs): + calls.append(script) + if 'taskSpace("LoopX public-source reader ' in script: + return space_response(19 if len(calls) == 1 else 20) + if "taskSpace(19)" in script: + output = reader.SPACE_MARKER + '{"closed":true}' + return subprocess.CompletedProcess([], 1, "" if stderr else output, output if stderr else "") + return response(extraction()) + monkeypatch.setattr(reader, "_run", run) + assert reader.read_public_url(URL)["ok"] + assert reader._OWNED_SPACE.space == 20 and len(calls) == 4 + + +@pytest.mark.parametrize("auto", [True, False]) +def test_shutdown_only_finishes_its_created_space_once(configured, monkeypatch, auto): + if auto: + auto_config(monkeypatch) + calls = [] + def run(_executable, script, **_kwargs): + calls.append(script) + return space_response(19) + monkeypatch.setattr(reader, "_run", run) + reader._OWNED_SPACE.resolve(reader.ReaderConfig.from_environment()) + reader._OWNED_SPACE.close() + reader._OWNED_SPACE.close() + if auto: + assert len(calls) == 2 + assert "taskSpace(19)" in calls[-1] and "ownership==='agent'" in calls[-1] + assert "finish({keep:[]})" in calls[-1] + else: + assert not calls + + +def test_creation_and_read_share_one_timeout_budget(configured, monkeypatch): + auto_config(monkeypatch) + clock = iter([0, 5, 29]) + monkeypatch.setattr(reader.time, "monotonic", lambda: next(clock)) + timeouts = [] + def run(_args, **kwargs): + timeouts.append(kwargs["timeout"]) + return space_response(19) if len(timeouts) == 1 else response(extraction()) + monkeypatch.setattr(reader.subprocess, "run", run) + assert reader.read_public_url(URL)["ok"] + assert timeouts == [25, 1] + + +def test_stdio_hosts_create_isolated_spaces_and_clean_up_on_eof(configured, monkeypatch, tmp_path): + import asyncio + import os + import sys + from mcp import ClientSession, StdioServerParameters + from mcp.client.stdio import stdio_client + auto_config(monkeypatch) + log = tmp_path / "calls.jsonl" + configured.write_text("#!" + sys.executable + "\n" + + "import json,os,sys\nfrom pathlib import Path\n" + + "script=sys.argv[-1]\n" + + f"with Path({str(log)!r}).open('a') as f: f.write(json.dumps([os.getppid(),script])+'\\n')\n" + + "if 'taskSpace(\"LoopX public-source reader ' in script:\n" + + " print('LOOPX_READER_SPACE:'+json.dumps({'id':os.getppid()}),file=sys.stderr)\n" + + "elif 'finish({keep:[]})' not in script:\n" + + " print('LOOPX_PUBLIC_SOURCE:'+" + repr(response(extraction()).stdout.split(":", 1)[1]) + ",file=sys.stderr)\n") + configured.chmod(0o700) + async def host(): + params = StdioServerParameters(command=sys.executable, + args=["-m", "loopx.extensions.ego_source_reader"], env=dict(os.environ)) + async with stdio_client(params) as (incoming, outgoing): + async with ClientSession(incoming, outgoing) as session: + await session.initialize() + for _ in range(2): + result = await session.call_tool("read_public_url", {"url": URL}) + assert json.loads(result.content[0].text)["ok"] + async def journey(): + await asyncio.gather(host(), host()) + asyncio.run(journey()) + calls = [json.loads(line) for line in log.read_text().splitlines()] + hosts = {pid for pid, _ in calls} + assert len(hosts) == 2 + for pid in hosts: + scripts = [script for owner, script in calls if owner == pid] + assert len(scripts) == 4 + assert sum('taskSpace("LoopX public-source reader ' in s for s in scripts) == 1 + assert all(f"taskSpace({pid})" in s for s in scripts[1:]) + assert scripts[-1].endswith("if(t.ownership==='agent')await t.finish({keep:[]});") + + @pytest.mark.parametrize("value", [ extraction(url="https://other.test/private"), extraction(url="https://example.com/another-article"),