Skip to content
Merged
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
85 changes: 70 additions & 15 deletions services/core/tests/official_workspace_images_native.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import importlib.metadata
import json
from pathlib import Path
import re
import secrets
import time

Expand All @@ -17,12 +18,17 @@ def verify_workspace_images(client, foreign, http, agent_options, session_option
distribution = importlib.metadata.distribution("openai")
assert distribution.version == pin["sdk_version"]
assert json.loads(distribution.read_text("direct_url.json"))["vcs_info"]["commit_id"] == pin["commit"]
# This qualification selects the two declarations with message/result images.
# claudesdk.Configuration rejects failed images; codex.Configuration admits them.
harness = agent_options["x_agents_core"]["harness"]
assert harness in {"codex", "claude_sdk"}, "Images qualification requires Codex or Claude Code"
sessions = client.beta.agents.sessions
workspace = session_options["environment"].get("workspace_directory", "/workspace").rstrip("/")
root = str(client.base_url).rstrip("/") + "/agents"
headers = {"Authorization": "Bearer " + client.api_key, "OpenAI-Beta": "agents=v1"}
foreign_headers = {**headers, "Authorization": "Bearer " + foreign.api_key}
proof = {"engine": agent_options["x_agents_core"]["harness"], "checks": [], "runs": [], "calls": []}
proof = {"engine": harness, "checks": [], "runs": [], "calls": [],
"failed_result_images": "unverified" if harness == "codex" else "declared_unsupported"}
expected_messages = []
sid = None

Expand Down Expand Up @@ -71,7 +77,9 @@ def history():
assert [i["content"] for i in stored if i["type"] == "message" and i.get("role") == "user"] == [m["content"] for m in expected_messages]
outputs = {i["call_id"]: i for i in stored if i["type"] == "function_call_output"}
for call in proof["calls"]:
assert outputs[call["call_id"]]["output"] == call["output"]
result = outputs[call["call_id"]]
assert result["output"] == call["output"]
assert result["error"] == call.get("error")
return stored

def post(events, key=None, foreign_request=False):
Expand All @@ -80,13 +88,16 @@ def post(events, key=None, foreign_request=False):
hdr["Idempotency-Key"] = key
return http.post(root + "/sessions/" + sid + "/events", headers=hdr, json={"events": events})

def verify_file(path, expected, turn):
def verify_file(path, expected, turn, exact_answer=False):
assert path.startswith(workspace + "/")
public_path = "/workspace" + path[len(workspace):]
answers = [i for i in items() if i["type"] == "message" and i.get("role") == "assistant"]
answer = " ".join(p["text"] for p in answers[-1]["content"] if p["type"] == "output_text").lower()
positions = [answer.find(color) for color in expected]
assert all(p >= 0 for p in positions) and positions == sorted(positions), answer
if exact_answer:
assert re.sub(r"\s+", "", answer) == ",".join(expected), answer
else:
positions = [answer.find(color) for color in expected]
assert all(p >= 0 for p in positions) and positions == sorted(positions), answer
artifacts = list(sessions.artifacts.list(sid, limit=100))
matches = [a for a in artifacts if a.path == public_path and a.turn_id == turn]
assert len(matches) == 1, [a.to_dict() for a in artifacts]
Expand All @@ -100,6 +111,7 @@ def verify_file(path, expected, turn):
def run(value, handler=None):
observed = []
handled = False
prior_calls = len(proof["calls"])
with sessions.events.stream(sid, timeout=300) as stream:
sessions.events.create(sid, events=[event(value)])
expected_messages.extend(value)
Expand Down Expand Up @@ -129,20 +141,33 @@ def run(value, handler=None):
turn = sessions.turns.list(sid, order="desc").data[0]
assert terminal[0] == "agent.session.turn." + turn.status
history()
results = [e["item"] for e in observed if e["type"] == "agent.session.turn.item.added" and e["item"]["type"] == "function_call_output"]
for call in proof["calls"][prior_calls:]:
matching = [item for item in results if item["call_id"] == call["call_id"]]
assert len(matching) == 1 and matching[0]["output"] == call["output"]
assert matching[0]["error"] == call.get("error")
return turn

def submit(action, output, validate=False):
def submit(action, output, validate=False, success=True):
result = {"type": "agent.session.input.tool_result", "turn_id": action.turn_id,
"call_id": action.call_id, "success": True, "output": output}
"call_id": action.call_id, "success": success, "output": output}
if not success:
result["error"] = "Visual lookup failed; the attached diagnostics remain available."
key = "result-" + action.call_id
before = items()
if validate:
if agent_options["x_agents_core"]["harness"] == "claude_sdk":
if harness == "claude_sdk":
invalid = [{**result, "success": False}, {**result, "output": [{"type": "input_image", "image_url": "https://example.test/image.png"}]}]
pending = [a.to_dict() for a in sessions.retrieve(sid).required_actions]
for bad in invalid:
assert post([bad], key).status_code == 400
response = post([bad], key)
assert response.status_code == 400
error = response.json()["error"]
assert error["type"] == "invalid_request_error" and error["code"] == "unsupported_or_invalid_configuration"
assert items() == before
assert sessions.retrieve(sid).required_actions[0].call_id == action.call_id
assert [a.to_dict() for a in sessions.retrieve(sid).required_actions] == pending
# Reuse the same receipt key for the accepted result below.
proof["failed_result_images"] = "unsupported_rejection_verified"
assert post([result], foreign_request=True).status_code == 404
assert post([{**result, "call_id": "unknown-call"}]).status_code in {400, 404, 409}
assert items() == before
Expand Down Expand Up @@ -176,11 +201,12 @@ def initial_done():
history()
check("initial_png_native_workspace_files_artifact")

jpeg_colors = ["yellow", "blue", "red", "green"]
jpeg = "data:image/jpeg;base64," + base64.b64encode((Path(__file__).parent / "testdata/function-bands.jpg").read_bytes()).decode()
idle = messages([{"type": "input_image", "image_url": jpeg}]) + messages([text(
f"Write this image's four lowercase band colors from left to right, comma-separated with no newline, to {workspace}/outputs/idle.txt using native tools. Reply with those colors. Do not call get_visual.")])
turn = run(idle)
verify_file(f"{workspace}/outputs/idle.txt", ["yellow", "blue", "red", "green"], turn.id)
verify_file(f"{workspace}/outputs/idle.txt", jpeg_colors, turn.id)
check("prepared_image_only_jpeg_and_message_boundary")

active_colors = colors[1:] + colors[:1]
Expand All @@ -197,12 +223,41 @@ def active(action):
check("active_image_input_retry_and_native_tools")

result_colors = colors[2:] + colors[:2]
if result_colors == jpeg_colors:
result_colors = result_colors[1:] + result_colors[:1]
output = [text("Inspect this visual."), {"type": "input_image", "image_url": picture(result_colors, 15)}, text("Remember these band colors.")]
turn = run(messages([text(f"Call get_visual exactly once. Read its returned image and use native tools to write its four lowercase band colors in left-to-right order, comma-separated with no newline, to {workspace}/outputs/result.txt. Reply with the colors. Do not call get_visual again.")]),
lambda action: submit(action, output, validate=True))
verify_file(f"{workspace}/outputs/result.txt", result_colors, turn.id)
check("large_function_image_result_receipt_retry_isolation_and_artifact")

# Keep the large single PNG above and add two distinguishable images in
# one ordered result, reusing the JPEG already used for message input.
variants = [("multiple", True, result_colors)]
if harness == "codex":
failed_colors = result_colors[1:] + result_colors[:1]
if failed_colors == jpeg_colors:
failed_colors = failed_colors[1:] + failed_colors[:1]
variants.append(("failed", False, failed_colors))
for name, success, png_colors in variants:
multiple = [text("First image:"), {"type": "input_image", "image_url": picture(png_colors)},
text("Second image:"), {"type": "input_image", "image_url": jpeg},
text("Read each image left to right, in the order supplied.")]
ordered_colors = png_colors + jpeg_colors
path = f"{workspace}/outputs/{name}.txt"
prompt = ("Call get_visual exactly once. Even if it reports a failure, inspect its diagnostic images. "
"Read the first image's four band colors, then the second image's four band colors. "
f"Use native tools to write all eight lowercase names in that order, comma-separated with no newline, to {path}. "
"Reply with exactly the same eight-name comma-separated sequence. Do not repeat get_visual.")
turn = run(messages([text(prompt)]), lambda action: submit(action, multiple, success=success))
assert turn.status == "completed"
verify_file(path, ordered_colors, turn.id, exact_answer=True)
check(name + "_function_result_png_jpeg_order_sse_history_and_artifact")
if not success:
proof["failed_result_images"] = "passed"
save()

turns_before_continuation = len(sessions.turns.list(sid, limit=100).data)
before = history()
invalid = [event(messages([text("must not persist")])), event(messages([{"type": "input_image", "image_url": "https://example.test/image.png"}]))]
assert post(invalid).status_code == 400
Expand Down Expand Up @@ -235,9 +290,9 @@ def active(action):
if restart is not None:
restart()
assert history() == before
turn = run(messages([text(f"Recall the most recent image returned by get_visual from conversation history, without reading any file or calling get_visual. Write its four lowercase band colors in order, comma-separated with no newline, to {workspace}/outputs/resumed.txt using native tools, then reply with them.")]))
verify_file(f"{workspace}/outputs/resumed.txt", result_colors, turn.id)
assert len(sessions.turns.list(sid, limit=100).data) == 5
turn = run(messages([text(f"Recall both images in the most recent get_visual result from conversation history, without reading any file or calling get_visual. Write their eight lowercase band colors in image order, comma-separated with no newline, to {workspace}/outputs/resumed.txt using native tools, then reply with exactly that comma-separated sequence.")]))
verify_file(f"{workspace}/outputs/resumed.txt", ordered_colors, turn.id, exact_answer=True)
assert len(sessions.turns.list(sid, limit=100).data) == turns_before_continuation + 1
check(("cold_agent_host" if restart is not None else "warm") + "_history_continuation_without_replay")

pending = []
Expand All @@ -249,7 +304,7 @@ def cancel(action):
assert turn.status == "cancelled" and not sessions.retrieve(sid).required_actions
assert not any(i["type"] == "function_call_output" and i["call_id"] == pending[0] for i in items())
run(messages([text("Reply only PLAIN_OK. Do not call any tools.")]))
assert len(sessions.turns.list(sid, limit=100).data) == 7
assert len(sessions.turns.list(sid, limit=100).data) == turns_before_continuation + 3
assert any(i["type"] == "message" and i.get("role") == "assistant" and any("PLAIN_OK" in p.get("text", "") for p in i["content"]) for i in items())
check("pending_cancel_retry_and_text_continuation")
proof["passed"] = True
Expand Down