diff --git a/src/agent_env/a2a_agent/object_transfer.py b/src/agent_env/a2a_agent/object_transfer.py index 63f76abe..e9127ea7 100644 --- a/src/agent_env/a2a_agent/object_transfer.py +++ b/src/agent_env/a2a_agent/object_transfer.py @@ -7,9 +7,9 @@ import math import re from collections.abc import AsyncIterator, Callable, Collection, Mapping -from contextlib import asynccontextmanager +from contextlib import asynccontextmanager, nullcontext from dataclasses import dataclass -from typing import Any, Literal, TypeVar +from typing import TYPE_CHECKING, Any, Literal, TypeVar import httpx from agentenv_protocol.a2a_agent import ( @@ -37,6 +37,7 @@ from agentenv_protocol.transfers import ReadObject, WriteNamespaceGrant, WriteObject from pydantic import BaseModel, ValidationError +from agent_env.a2a_agent import protocol from agent_env.a2a_agent.protocol import raise_for_extension_status from agent_env.a2a_agent.staging import StagedObjectStore, staged_store from agent_env.config import get_config @@ -45,6 +46,9 @@ from agent_env.store.object_store.object_store import readable_url from agent_env.store.object_store.local.grant_server import unreachable_hint +if TYPE_CHECKING: + from agent_env.task_step.context import DeployedAgent + logger = logging.getLogger(__name__) # "objects": the call carries grants and the agent moves the bytes. "legacy": the inline forms, @@ -308,11 +312,13 @@ async def readable_parts( card: Mapping[str, Any] | None, sandbox_type: str | None, expires_in: int, + shareable: Collection[str] | None = None, ) -> AsyncIterator[list[dict]]: """``parts`` as the agent at ``a2a_url`` can read them: a file part naming an object a configured store owns names an HTTPS URL for it instead, one ``readable_url`` gives for at least ``expires_in`` seconds, or else a copy staged on the agent for the length of the block. Other parts, and file parts naming - anything else, are sent as they are. Raises when an owned object can be given no URL the agent can read.""" + anything else, are sent as they are; so is an owned object ``shareable`` (when given) does not name. + Raises when an owned object can be given no URL the agent can read.""" config = get_config() readable = list(parts) staged: dict[int, StagedObjectStore] = {} # by identity: a store need not be hashable @@ -322,7 +328,7 @@ async def readable_parts( if not isinstance(uri, str): continue store = config.get_object_store_at(uri) - if not store.owns(uri): + if not store.owns(uri) or (shareable is not None and uri not in shareable): continue url = await asyncio.to_thread(readable_url, store, uri, sandbox_type=sandbox_type, expires_in=expires_in) if url is None: @@ -346,6 +352,42 @@ async def readable_parts( await staging.release() +async def send_and_wait( + a2a_url: str, + parts: list[dict], + *, + agent: DeployedAgent | None, + message_id: str, + context_id: str | None, + timeout_seconds: int, + poll_interval_seconds: int, + before_send: Callable[[], None] | None = None, + shareable: Collection[str] | None = None, +) -> tuple[str, dict]: + """Send ``parts`` to the A2A peer at ``a2a_url`` and wait for its task to end: the task's id and its + final state. A peer that is an ``agent`` agent-env deployed is sent each file part a configured store + owns, of those ``shareable`` names when given, as an HTTPS URL it can read (``readable_parts``); any + other peer, such as a human's hub, reads the store itself and is sent the parts as they are. + ``before_send`` runs once the parts are ready, just before the message goes out. An ``agent``'s sandbox is + watched while it works, so one that dies is given up on (``poll_a2a_task``).""" + sending = ( + readable_parts( + parts, a2a_url=a2a_url, card=agent.a2a_card, sandbox_type=agent.sandbox_type, + expires_in=timeout_seconds, shareable=shareable, + ) + if agent is not None + else nullcontext(parts) + ) + async with sending as sent: + if before_send is not None: + before_send() + task_id, _ = await protocol.send_a2a_message(a2a_url, sent, message_id, context_id, timeout_seconds) + return task_id, await protocol.poll_a2a_task( + a2a_url, task_id, timeout_seconds, poll_interval_seconds, + sandbox_id=agent.sandbox_id if agent is not None else None, + ) + + def skill_bundle_request( store: ObjectStore, *, name: str, description: str, object_url: str ) -> BundleSkillRequest: diff --git a/src/agent_env/task_step/task_steps/prompt_agent.py b/src/agent_env/task_step/task_steps/prompt_agent.py index af14884b..6b34136b 100644 --- a/src/agent_env/task_step/task_steps/prompt_agent.py +++ b/src/agent_env/task_step/task_steps/prompt_agent.py @@ -17,7 +17,7 @@ from agent_env.a2a_agent.object_transfer import ( TrajectoryUpload, fetch_trajectory, - readable_parts, + send_and_wait, trajectory_mode, ) from agent_env.a2a_agent.staging import draining, staged_changelogs, transfer_store @@ -91,6 +91,11 @@ def _duplicates_prompt_text(parts: list[dict], prompt_text: Optional[str]) -> bo return parts == [{"kind": "text", "text": prompt_text}] +def _file_uris(parts: list[dict]) -> list[str]: + files = (part.get("file") for part in parts if part.get("kind") == "file") + return [file["uri"] for file in files if isinstance(file, dict) and isinstance(file.get("uri"), str)] + + _DEFAULT_USER_SIM_OUTPUT_FORMAT: dict[str, Any] = { "type": "json_schema", "schema": { @@ -474,18 +479,20 @@ async def _execute_conversation( traj_ext_cached = A2AAgent.find_extension(card, A2AAgent.EXT_TRAJECTORY) final_state: str = TaskState.completed.value trajectory_s3_uri: Optional[str] = None + # What the target has been sent: of the files its replies name, the only ones a user-sim is made able + # to read, so a reply naming any other object a store owns can't read it out through the user-sim. + sent_to_target: set[str] = set() for turn in range(self.max_conversation_turns): # `target_a2a_task_id` is the client A2A message id sent to the target # agent and recorded on the conversation as `a2a_task_id`. # This id is also used as part of the key name for the trajectory S3 object. target_a2a_task_id = uuid.uuid4().hex - # Only the sent copy names readable URLs; the turn records the objects' own URLs, and only - # once the copy is ready, so an object the agent can't be sent leaves no turn waiting. - async with readable_parts( - current_user_parts, a2a_url=target_url, card=card, - sandbox_type=agent.sandbox_type, expires_in=self.timeout_seconds, - ) as sent_parts: + + # The turn records the objects' own URLs, and only once the parts are ready to send, so an + # object the agent can't be sent leaves no turn waiting. + def record_turn() -> None: + sent_to_target.update(_file_uris(current_user_parts)) conversation_store.add_a2a_task( conversation_id=conversation_id, parts=current_user_parts, @@ -502,14 +509,11 @@ async def _execute_conversation( else list(current_user_parts) ) - sent_task_id, _ = await protocol.send_a2a_message( - target_url, sent_parts, target_a2a_task_id, - solver_context_id, self.timeout_seconds, - ) - result = await protocol.poll_a2a_task( - target_url, sent_task_id, self.timeout_seconds, self.poll_interval_seconds, - sandbox_id=agent.sandbox_id, - ) + sent_task_id, result = await send_and_wait( + target_url, current_user_parts, agent=agent, message_id=target_a2a_task_id, + context_id=solver_context_id, timeout_seconds=self.timeout_seconds, + poll_interval_seconds=self.poll_interval_seconds, before_send=record_turn, + ) target_state = result["status"]["state"] status_msg = (result.get("status") or {}).get("message") or {} final_terminal = protocol.TerminalResponse.from_message(status_msg) @@ -573,14 +577,14 @@ async def _execute_conversation( user_a2a_task_id = uuid.uuid4().hex try: - sent_user_task_id, _ = await protocol.send_a2a_message( - user_url, agent_response_parts, user_a2a_task_id, - conversation_id, self.user_agent_timeout_seconds, - ) - user_result = await protocol.poll_a2a_task( - user_url, sent_user_task_id, - self.user_agent_timeout_seconds, self.poll_interval_seconds, - sandbox_id=user_sim.sandbox_id if is_user_sim else None, + # A user-sim runs in a sandbox agent-env deployed; a human peer, registered or named by + # user_a2a_url, has none and reads the store itself. + _, user_result = await send_and_wait( + user_url, agent_response_parts, + agent=user_sim if is_user_sim and user_sim.sandbox_id else None, shareable=sent_to_target, + message_id=user_a2a_task_id, context_id=conversation_id, + timeout_seconds=self.user_agent_timeout_seconds, + poll_interval_seconds=self.poll_interval_seconds, ) except TimeoutError as e: logger.warning(f"user_a2a_url timeout for conversation {conversation_id} ({e}); marking abandoned") diff --git a/tst/data/a2a_agent/agent.py b/tst/data/a2a_agent/agent.py index 6fc085c4..35e2d253 100644 --- a/tst/data/a2a_agent/agent.py +++ b/tst/data/a2a_agent/agent.py @@ -3,10 +3,17 @@ Deterministic: it echoes the prompt and records a native trajectory. Agent-config, mcp-config, trajectory and triggers come from the agentenv-protocol framework. ``model_params`` is a write-only config field so the agent-config negotiation and its redaction on read-back can be observed. + +It also moves files both ways. After the echo it adds a ``read over : `` line for +each file part it is sent, fetching only inline bytes and HTTP(S) URLs (``could not read ...`` otherwise), +and it answers each ``send-file `` line of its prompt with a file part naming that URI. """ +import base64 from typing import Any +from urllib.parse import urlsplit +import httpx from agentenv_protocol.a2a_agent import ( MCP_CONFIG_V1, TRAJECTORY_V1, @@ -14,6 +21,7 @@ AgentConfig, AgentEnvAgent, AgentIdentity, + FilePart, TaskRequest, TaskResult, TextPart, @@ -22,6 +30,9 @@ a2a_agent, ) +SEND_FILE = "send-file " +READ_TIMEOUT_SECONDS = 120 + class EchoAgentConfig(AgentConfig): model: str | None = None @@ -30,6 +41,22 @@ class EchoAgentConfig(AgentConfig): model_params: WriteOnly[dict[str, Any] | None] = None +async def _read(part: FilePart) -> str: + """What a file part holds, inline or at an HTTP(S) URL.""" + scheme = "bytes" if part.bytes is not None else urlsplit(part.uri).scheme + try: + if part.bytes is not None: + data = base64.b64decode(part.bytes) + else: + async with httpx.AsyncClient(timeout=READ_TIMEOUT_SECONDS) as client: + response = await client.get(part.uri) + response.raise_for_status() + data = response.content + except Exception as exc: # noqa: BLE001 -- said in the reply, so a test sees why + return f"could not read {part.name} over {scheme}: {type(exc).__name__}" + return f"read {part.name} over {scheme}: {data.decode(errors='replace').strip()}" + + @a2a_agent( identity=AgentIdentity( name="agentenv-echo-agent", @@ -43,10 +70,13 @@ class EchoAgent(AgentEnvAgent): async def run(self, request: TaskRequest[EchoAgentConfig]) -> TaskResult: prompt = "\n".join(part.text for part in request.parts if isinstance(part, TextPart)) reply = f"Echo: {prompt}" + reply = "\n".join([reply, *[await _read(part) for part in request.parts if isinstance(part, FilePart)]]) + sends = [line.removeprefix(SEND_FILE).strip() for line in prompt.splitlines() if line.startswith(SEND_FILE)] + files = [FilePart(uri=uri, name=uri.rsplit("/", 1)[-1], mime_type="text/plain") for uri in sends] return ( TaskResult.builder() .succeeded() - .add_text(reply) + .parts([TextPart(text=reply), *files]) .usage(Usage(tool_call_count=0, input_tokens=len(prompt), output_tokens=len(reply))) .native_trajectory( format="agentenv-echo-agent/v1", diff --git a/tst/integration/task_step/test_prompt_agent_file_parts_local.py b/tst/integration/task_step/test_prompt_agent_file_parts_local.py new file mode 100644 index 00000000..c8bdc57f --- /dev/null +++ b/tst/integration/task_step/test_prompt_agent_file_parts_local.py @@ -0,0 +1,198 @@ +"""prompt_agent hands store files to the agents it talks to, on the local defaults. A deployed agent is sent a URL it +can fetch for each file in its prompt, and a user-sim one for each file the solver passes back of those it was sent; +a file the solver names but was never sent stays as it is, so the user-sim can't read it. A human peer, which reads +the store itself, is sent the store's own URL. The run records the store's own URLs throughout. The agents are the +echo agent (tst/data/a2a_agent), which reports what each file part held. Needs Docker and a throwaway local registry.""" + +import importlib.util +import json +import re +import shutil +import socket +import subprocess +import threading +import time +import uuid + +import httpx +import pytest +import uvicorn + +from agent_env.a2a_agent import conversation_store +from agent_env.artifact.store import reset_artifact_store +from agent_env.config import configure, get_config, reset_config, set_image_store +from agent_env.store.image_store import LocalRegistryImageStore +from agent_env.store.object_store.local.grant_server import grant_server +from agent_env.task import Task +from agent_env.task_step.task_steps.deploy_agent import DeployAgentTaskStep +from agent_env.task_step.task_steps.deploy_human_agent import DeployHumanAgentTaskStep +from agent_env.task_step.task_steps.prompt_agent import PromptAgentTaskStep +from tst.util.a2a_test_agent import AGENT_DIR, put_test_agent + +pytestmark = pytest.mark.integration + +REGISTRY_READY_SECONDS = 45 +PEER_READY_SECONDS = 15 + + +def _docker(*args): + return subprocess.run(["docker", *args], capture_output=True, text=True) + + +def _free_port() -> int: + with socket.socket() as s: + s.bind(("127.0.0.1", 0)) + return s.getsockname()[1] + + +@pytest.fixture +def local_registry(monkeypatch, tmp_path): + sandboxes = tmp_path / "sandboxes" + monkeypatch.chdir(tmp_path) + monkeypatch.setenv("XDG_STATE_HOME", str(tmp_path / "state")) + monkeypatch.setenv("AGENT_ENV_DOCUMENT_STORE", "local") + monkeypatch.setenv("AGENT_ENV_OBJECT_STORE", "local") + monkeypatch.setenv("AGENT_ENV_IMAGE_STORE", "local") + monkeypatch.setenv("AGENT_ENV_LOCAL_SANDBOX_DIR", str(sandboxes)) + port = _free_port() + name = f"file-parts-registry-{uuid.uuid4().hex[:8]}" + started = _docker("run", "-d", "--rm", "--name", name, "-p", f"127.0.0.1:{port}:5000", "registry:2") + assert started.returncode == 0, started.stderr + host = f"localhost:{port}" + deadline = time.time() + REGISTRY_READY_SECONDS + while time.time() < deadline: + try: + if httpx.get(f"http://{host}/v2/", timeout=2).status_code in (200, 401): + break + except httpx.HTTPError: + pass + time.sleep(0.5) + configure() + set_image_store(LocalRegistryImageStore(host)) + reset_artifact_store() + try: + yield tmp_path + finally: + # Only this test's containers: every local sandbox it deployed has a work dir in its own sandbox root, and its + # container is named after that sandbox. Images are shared, so matching by image would reach other runs. + for work_dir in sandboxes.glob("agent-env-*"): + if m := re.match(r"agent-env-(local-[0-9a-f]+)-", work_dir.name): + _docker("rm", "-f", f"agent-{m.group(1)}") + _docker("rm", "-f", name) + shutil.rmtree(sandboxes, ignore_errors=True) + # The agents fetched through this process's grant server, whose certificate is from this test's state root; + # a later test's agents trust another root's CA, so it must start afresh. + grant_server(None, None).close() + reset_artifact_store() + reset_config() + + +@pytest.fixture +def human_peer(): + """The echo agent served in this process, standing in for a human's hub: its URL, and the file URIs it is sent.""" + spec = importlib.util.spec_from_file_location("echo_agent_here", AGENT_DIR / "agent.py") + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + agent, received = module.EchoAgent(), [] + run = agent.run + + async def recording(request: module.TaskRequest[module.EchoAgentConfig]): + received.extend(part.uri for part in request.parts if isinstance(part, module.FilePart) and part.uri) + return await run(request) + + agent.run = recording + port = _free_port() + server = uvicorn.Server(uvicorn.Config(agent.create_app(), host="127.0.0.1", port=port, log_level="warning")) + thread = threading.Thread(target=server.run, daemon=True) + thread.start() + deadline = time.time() + PEER_READY_SECONDS + while not server.started and time.time() < deadline: + time.sleep(0.1) + try: + yield f"http://127.0.0.1:{port}", received + finally: + server.should_exit = True + thread.join(timeout=10) + + +def _files(suffix: str) -> dict[str, str]: + """Two text files in the store, by name: the object URL of each, whose content is ``-``.""" + store = get_config().get_object_store() + return { + name: store.put(f"file-parts-{suffix}/{name}.txt", f"{name}-{suffix}".encode(), content_type="text/plain") + for name in ("brief", "report") + } + + +def _deploy(agent, name: str) -> DeployAgentTaskStep: + return DeployAgentTaskStep( + id=f"deploy-{name}", version=None, env_ids=[], a2a_agent_id=agent.id, a2a_agent_version=agent.version, + agent_name=name, sandbox_type="local", + ) + + +def _conversation(ctx) -> list[dict]: + """The conversation's messages: the solver's prompt, its reply, the peer's reply, the solver's second reply.""" + assert not ctx.metadata.get("failed_steps"), ctx.metadata.get("failed_steps") + return conversation_store.get_conversation(ctx.metadata["a2a_conversations"]["ask"])["messages"] + + +def _text(message: dict) -> str: + return "\n".join(part["text"] for part in message["parts"] if part["kind"] == "text") + + +def _uris(message: dict) -> list[str]: + return [part["file"]["uri"] for part in message["parts"] if part["kind"] == "file"] + + +@pytest.mark.asyncio +async def test_the_solver_and_a_user_sim_fetch_the_files_the_conversation_shares_and_no_other(local_registry): + suffix = uuid.uuid4().hex[:8] + urls = _files(suffix) + agent = put_test_agent(f"file-parts-agent-{suffix}") + task = Task.put(id=f"file-parts-{suffix}", steps=[ + _deploy(agent, "solver"), + _deploy(agent, "user"), + PromptAgentTaskStep( + id="ask", version=None, prompt_id="ask", agent_name="solver", user_agent_name="user", + max_conversation_turns=2, poll_interval_seconds=1, + parts=[ + {"kind": "text", "text": f"send-file {urls['brief']}\nsend-file {urls['report']}"}, + {"kind": "file", "file": {"uri": urls["brief"], "name": "brief.txt", "mimeType": "text/plain"}}, + ], + ), + ]) + + prompt, reply, user_reply, _ = _conversation(await task.run()) + + assert f"read brief.txt over https: brief-{suffix}" in _text(reply) + # The user-sim's own lines follow its echo of the reply, one per file it was sent. + shared, never_sent = _text(user_reply).splitlines()[-2:] + assert shared == f"read brief.txt over https: brief-{suffix}" + assert never_sent.startswith("could not read report.txt over file:") + assert (_uris(prompt), _uris(reply)) == ([urls["brief"]], [urls["brief"], urls["report"]]) + assert "https://" not in json.dumps([prompt, reply, user_reply]) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("registered", [False, True], ids=["user_a2a_url", "registered"]) +async def test_a_human_peer_is_sent_the_stores_own_url(local_registry, human_peer, registered): + suffix = uuid.uuid4().hex[:8] + urls = _files(suffix) + agent = put_test_agent(f"file-parts-agent-{suffix}") + peer_url, received = human_peer + peer = {"user_agent_name": "human"} if registered else {"user_a2a_url": peer_url} + task = Task.put(id=f"file-parts-human-{suffix}", steps=[ + _deploy(agent, "solver"), + *([DeployHumanAgentTaskStep(id="register-human", version=None, agent_name="human", a2a_url=peer_url)] + if registered else []), + PromptAgentTaskStep( + id="ask", version=None, prompt_id="ask", agent_name="solver", max_conversation_turns=2, + poll_interval_seconds=1, prompt=f"send-file {urls['report']}", **peer, + ), + ]) + + _, reply, _, _ = _conversation(await task.run()) + + assert _uris(reply) == [urls["report"]] + assert received == [urls["report"]] diff --git a/tst/unit/a2a_agent/readable_parts_test.py b/tst/unit/a2a_agent/readable_parts_test.py index f5fe8344..5f486683 100644 --- a/tst/unit/a2a_agent/readable_parts_test.py +++ b/tst/unit/a2a_agent/readable_parts_test.py @@ -52,6 +52,17 @@ async def test_parts_that_name_no_owned_object_are_sent_as_they_are(store): assert store.granted == [] +@pytest.mark.asyncio +async def test_only_the_owned_objects_named_shareable_are_made_readable(store): + shared, other = store.put("seeds/x.png", b"png"), store.put("answers/y.png", b"png") + parts = [_file(shared), _file(other)] + + async with readable_parts(parts, a2a_url=URL, card=None, sandbox_type="local", expires_in=600, shareable={shared}) as sent: + assert sent == [_file(f"{GRANT_ORIGIN}/seeds/x.png?sig=read"), _file(other)] + + assert store.granted == [shared] + + @pytest.mark.asyncio async def test_an_owned_object_the_agent_can_be_given_no_url_for_is_not_sent(store): url = store.put("seeds/s1/x.png", b"png") diff --git a/tst/unit/task_step/prompt_agent_file_parts_test.py b/tst/unit/task_step/prompt_agent_file_parts_test.py index b11eb7a1..cb1d196c 100644 --- a/tst/unit/task_step/prompt_agent_file_parts_test.py +++ b/tst/unit/task_step/prompt_agent_file_parts_test.py @@ -1,7 +1,9 @@ """prompt_agent sends the agent an HTTPS URL for each file part naming an object the store owns, on every turn, while the run records the object's own URL.""" +import httpx import pytest +from agentenv_protocol.a2a_agent import STAGING_V1_URI from agent_env.config import configure from agent_env.task_step.context import DeployedAgent, TaskStepContext @@ -13,6 +15,7 @@ AGENT_URL = "http://agent.test" USER_URL = "http://user.test" DONE = [{"kind": "text", "text": "done"}] +USER_DONE = [{"kind": "text", "text": '{"message": "thanks", "done": true}'}] @pytest.fixture @@ -96,3 +99,117 @@ async def test_an_object_the_agent_could_not_read_is_never_sent(monkeypatch, sto assert agents.sent[AGENT_URL] == [] assert recorded == [] # no turn left waiting for a reply + + +def _replying_with(url): + return [{"kind": "text", "text": "here it is"}, _file(url)] + + +def _two_turns(*files): + return PromptAgentTaskStep( + id="solve", version=None, agent_name="solver", poll_interval_seconds=0, max_conversation_turns=2, + parts=[{"kind": "text", "text": "hi"}, *(_file(url) for url in files)], + ) + + +def _user_sim(url, sandbox_type="local", card=None): + return DeployedAgent( + agent_name="human_agent", api_url=url, a2a_url=url, sandbox_id="sb-user", sandbox_type=sandbox_type, a2a_card=card, + ) + + +@pytest.mark.asyncio +async def test_a_user_sim_agent_is_sent_a_file_the_target_passes_on_as_a_grant(monkeypatch, store, recorded): + url = store.put("inputs/z.png", b"png") + agents = FakeA2AAgents({AGENT_URL: _replying_with(url), USER_URL: USER_DONE}).serve(monkeypatch) + context = _context() + context.deployed_agents.append(_user_sim(USER_URL)) + + await _two_turns(url).execute(context) + + assert agents.sent[USER_URL] == [ + [{"kind": "text", "text": "here it is"}, _file(f"{GRANT_ORIGIN}/inputs/z.png?sig=read")] + ] + + +@pytest.mark.asyncio +async def test_a_user_sim_is_not_made_able_to_read_a_file_the_target_was_never_sent(monkeypatch, store, recorded): + sent = store.put("inputs/z.png", b"png") + other = store.put("someone-elses/answer.json", b"{}") + agents = FakeA2AAgents({AGENT_URL: _replying_with(other), USER_URL: USER_DONE}).serve(monkeypatch) + context = _context() + context.deployed_agents.append(_user_sim(USER_URL)) + + await _two_turns(sent).execute(context) + + assert agents.sent[USER_URL] == [[{"kind": "text", "text": "here it is"}, _file(other)]] + assert store.granted == [sent] # the target's own file, and nothing for the user-sim + + +@pytest.mark.asyncio +async def test_a_user_sim_the_stores_grants_cannot_reach_is_sent_a_copy_staged_on_it(monkeypatch, store, recorded): + url = store.put("inputs/z.png", b"png") + user_url = "https://user.test" + staging = [] # each staging request, with how many messages the user-sim had been sent by then + + def staging_route(request): + staging.append((request.method, len(agents.sent[user_url]))) + return httpx.Response(201 if request.method == "PUT" else 204) + + agents = FakeA2AAgents({AGENT_URL: _replying_with(url), user_url: USER_DONE}, other=staging_route).serve(monkeypatch) + context = _context() + card = {"capabilities": {"extensions": [{"uri": STAGING_V1_URI, "params": {"endpoint": "/ext/staging"}}]}} + context.deployed_agents.append(_user_sim(user_url, "modal", card)) + + await _two_turns(url).execute(context) + + ((text, file),) = agents.sent[user_url] + assert file["file"]["uri"].startswith(f"{user_url}/ext/staging/") + assert staging == [("PUT", 0), ("DELETE", 1)] # staged before the message went, cleared once it was answered + + +@pytest.mark.asyncio +async def test_a_user_sim_that_can_be_sent_no_readable_url_fails_the_step_before_it_is_sent_anything( + monkeypatch, store, recorded +): + url = store.put("inputs/z.png", b"png") + agents = FakeA2AAgents({AGENT_URL: _replying_with(url), USER_URL: USER_DONE}).serve(monkeypatch) + context = _context() + context.deployed_agents.append(_user_sim(USER_URL, "modal")) + + with pytest.raises(RuntimeError, match="cannot be sent to the agent"): + await _two_turns(url).execute(context) + + assert agents.sent[USER_URL] == [] + assert recorded == [[{"kind": "text", "text": "hi"}, _file(url)]] # the solver's turn, which ran + + +@pytest.mark.asyncio +async def test_a_registered_human_peer_is_sent_the_objects_own_url(monkeypatch, store, recorded): + url = store.put("outputs/z.png", b"png") + agents = FakeA2AAgents({AGENT_URL: _replying_with(url), USER_URL: DONE}).serve(monkeypatch) + context = _context() + context.deployed_agents.append(DeployedAgent(agent_name="human_agent", api_url=USER_URL, a2a_url=USER_URL)) + step = PromptAgentTaskStep( + id="solve", version=None, agent_name="solver", prompt="hi", poll_interval_seconds=0, max_conversation_turns=2, + ) + + await step.execute(context) + + assert agents.sent[USER_URL] == [[{"kind": "text", "text": "here it is"}, _file(url)]] + assert store.granted == [] + + +@pytest.mark.asyncio +async def test_a_humans_hub_is_sent_the_objects_own_url(monkeypatch, store, recorded): + url = store.put("outputs/z.png", b"png") + agents = FakeA2AAgents({AGENT_URL: _replying_with(url), USER_URL: DONE}).serve(monkeypatch) + step = PromptAgentTaskStep( + id="solve", version=None, agent_name="solver", prompt="hi", poll_interval_seconds=0, + max_conversation_turns=2, user_a2a_url=USER_URL, + ) + + await step.execute(_context()) + + assert agents.sent[USER_URL] == [[{"kind": "text", "text": "here it is"}, _file(url)]] + assert store.granted == []