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
88 changes: 79 additions & 9 deletions fixtures/mock-acp-v2.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@
state_lock = threading.Lock()
next_session = 1
next_injection = 1
# state_lock owns IDs and pending contents; delivery removes before emitting.
pending_injections = {}
supports_fork = "--no-fork" not in sys.argv
supports_steer = "--steer" in sys.argv
supports_models = "--models" in sys.argv
Expand Down Expand Up @@ -38,6 +40,8 @@ def argument_value(name):
prompt_release_text = option("--prompt-release-text")
inject_release = option("--inject-release")
inject_ack_release = option("--inject-ack-release")
mutation_ack_release = option("--mutation-ack-release")
mutation_committed = option("--mutation-committed")
close_release = option("--close-release")
close_release_session = option("--close-release-session")
fail_close_session = option("--fail-close-session")
Expand Down Expand Up @@ -89,6 +93,8 @@ def log_request(request):
entry["sessionId"] = params["sessionId"]
if request.get("method") in ("session/new", "session/fork"):
entry["cwd"] = params["cwd"]
if "additionalDirectories" in params:
entry["additionalDirectories"] = params["additionalDirectories"]
if request.get("method") == "session/prompt":
entry["text"] = params["prompt"][0]["text"]
if request.get("method") == "session/inject":
Expand Down Expand Up @@ -275,6 +281,12 @@ def prompt(request):
sys.stderr.write("\x01kit-runtime\x01" + json.dumps({"event": "child_finished", "call": "call-1:compose:shell", "tool": "shell", "ok": True, "summary": "done", "millis": 2}) + "\nmock diagnostic\n")
sys.stderr.flush()
text = "rich done"
if "MOCK_SHARED_MESSAGE_IDS" in text:
for role, value in [("user_message", "User text"), ("agent_thought_chunk", "Thought text"), ("agent_message_chunk", "Assistant text")]:
content = {"type": "text", "text": value}
if role == "user_message":
content = [content]
send({"jsonrpc": "2.0", "method": "session/update", "params": {"sessionId": session_id, "update": {"sessionUpdate": role, "messageId": "shared-id", "content": content}}})
send({
"jsonrpc": "2.0",
"method": "session/update",
Expand Down Expand Up @@ -309,17 +321,49 @@ def inject(request):
with state_lock:
message_id = f"injected-{next_injection}"
next_injection += 1
pending_injections[message_id] = (params["sessionId"], params["content"])
while inject_ack_release is not None and not os.path.exists(inject_ack_release):
time.sleep(0.01)
respond(request["id"], {"messageId": message_id})
if inject_ack_release is not None:
send({"jsonrpc": "2.0", "method": "session/update", "params": {"sessionId": params["sessionId"], "update": {"sessionUpdate": "notice", "title": "Injection acknowledged"}}})
while inject_release is not None and not os.path.exists(inject_release):
time.sleep(0.01)
send({
"jsonrpc": "2.0", "method": "session/update",
"params": {"sessionId": params["sessionId"], "update": {
"sessionUpdate": "user_message", "messageId": message_id, "content": params["content"],
}},
})
with state_lock:
pending = pending_injections.pop(message_id, None)
if pending is not None:
send({
"jsonrpc": "2.0", "method": "session/update",
"params": {"sessionId": pending[0], "update": {
"sessionUpdate": "user_message", "messageId": message_id, "content": pending[1],
}},
})


def mutate_pending_injection(request):
params = request["params"]
replacement = params.get("content")
rejected = replacement is not None and any("MOCK_REJECT_REPLACE" in block.get("text", "") for block in replacement)
with state_lock:
pending = pending_injections.get(params["messageId"])
accepted = not rejected and pending is not None and pending[0] == params["sessionId"]
if accepted:
if replacement is None:
del pending_injections[params["messageId"]]
else:
pending_injections[params["messageId"]] = (pending[0], replacement)
if mutation_committed is not None:
with open(mutation_committed, "w", encoding="utf-8") as marker:
marker.write("committed" if accepted else "rejected")
# No I/O or callbacks while state_lock is held. Delayed replies exercise
# delivery racing an already committed mutation without resurrecting rows.
while mutation_ack_release is not None and not os.path.exists(mutation_ack_release):
time.sleep(0.01)
if accepted:
respond(request["id"], {})
else:
send({"jsonrpc": "2.0", "id": request["id"], "error": {"code": -32602, "message": "pending injection unavailable or rejected"}})
send({"jsonrpc": "2.0", "method": "session/update", "params": {"sessionId": params["sessionId"], "update": {"sessionUpdate": "notice", "title": "Mutation acknowledged"}}})


def close(request):
Expand Down Expand Up @@ -361,7 +405,8 @@ def close(request):
"authMethods": [{"type": "agent", "methodId": "browser-login", "name": "secret-name", "description": "secret-description"}],
"capabilities": {"session": {
"prompt": {name: {} for name in prompt_capability_names if name},
**({"inject": {"modes": ["steer"], "steerInStream": ["finish"]}} if supports_steer else {}),
**({"additionalDirectories": {}} if "--additional-directories" in sys.argv else {}),
**({"inject": {"modes": ["steer"], "steerInStream": ["finish"], "pending": {"replace": "--pending-replace" in sys.argv}}} if supports_steer else {}),
}},
})
elif method == "session/new":
Expand Down Expand Up @@ -437,10 +482,17 @@ def close(request):
sys.stderr.write("\x01kit-runtime\x01" + json.dumps({"event": "session_started", "session_id": session_id}) + "\n")
sys.stderr.flush()
selected_models[session_id] = model_ids[0]
for update in [
replay_updates = [
{"sessionUpdate": "user_message_chunk", "messageId": "user-1", "content": {"type": "text", "text": "replayed user"}},
{"sessionUpdate": "agent_message_chunk", "messageId": "agent-1", "content": {"type": "text", "text": "replayed assistant"}},
]:
]
if "--replay-shared-ids" in sys.argv:
replay_updates = [
{"sessionUpdate": kind, "messageId": "shared", "content": {"type": "text", "text": "text"}}
for kind in ["user_message_chunk", "agent_thought_chunk", "agent_message_chunk"]
] + [{"sessionUpdate": "agent_message", "messageId": "shared",
"content": [{"type": "text", "text": "answer"}]}]
for update in replay_updates:
send({"jsonrpc": "2.0", "method": "session/update", "params": {"sessionId": session_id, "update": update}})
result = {}
if supports_models:
Expand All @@ -455,10 +507,28 @@ def close(request):
}})
elif method == "session/fork":
threading.Thread(target=fork, args=(request,), daemon=True).start()
elif method == "kit/subagent/transcript/read":
params = request["params"]
if "session_id" not in params or params.get("generation") != 1 or params.get("id") != "fixture-child":
send({"jsonrpc": "2.0", "id": request["id"], "error": {"code": -32602, "message": "invalid child identity"}})
else:
cursor = params["cursor"]
updates = [] if cursor else [
{"sessionUpdate": "agent_message_chunk", "messageId": "child-message", "content": {"type": "text", "text": "Child-only fixture output"}},
]
respond(request["id"], {"updates": updates, "next_cursor": 128, "generation": 1, "caught_up": True})
elif method == "kit/subagent/steer":
params = request["params"]
if "session_id" not in params or params.get("generation") != 1 or params.get("id") != "fixture-child" or not isinstance(params.get("prompt"), str):
send({"jsonrpc": "2.0", "id": request["id"], "error": {"code": -32602, "message": "invalid child steer"}})
else:
respond(request["id"], {"receipt": {"accepted": True}})
elif method == "session/prompt":
threading.Thread(target=prompt, args=(request,), daemon=True).start()
elif method == "session/inject":
threading.Thread(target=inject, args=(request,), daemon=True).start()
elif method in ("session/replace_inject", "session/revoke_inject"):
threading.Thread(target=mutate_pending_injection, args=(request,), daemon=True).start()
elif method == "session/set_config_option":
params = request["params"]
value = params["value"]
Expand Down
51 changes: 51 additions & 0 deletions fixtures/mock-openrouter-shell.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
#!/usr/bin/env python3
"""Local provider boundary for the desktop's real ACP/compose/shell regression.

Prints one ephemeral loopback port, then serves OpenRouter-compatible SSE.
No model service, credentials, or shell execution occurs in this fixture.
"""
import json
from http.server import BaseHTTPRequestHandler, HTTPServer


class Provider(BaseHTTPRequestHandler):
def log_message(self, *args):
pass

def do_GET(self):
self.respond(json.dumps({"data": []}).encode(), "application/json")

def do_POST(self):
body = json.loads(self.rfile.read(int(self.headers["Content-Length"])))
results = [m for m in body["messages"] if m.get("role") == "tool"]
if results:
# Do not claim success merely because a tool returned an error.
result = json.loads(results[-1]["content"])
success = result.get("success") is True and result.get("exit_code") == 0
delta = {"role": "assistant", "content": "REAL_SHELL_COMPLETE" if success else "REAL_SHELL_FAILED"}
reason = "stop"
else:
name = next(t["function"]["name"] for t in body["tools"]
if t["function"]["name"].split(".")[-1] == "compose")
delta = {"role": "assistant", "tool_calls": [{
"index": 0, "id": "real-shell-pwd", "type": "function",
"function": {"name": name, "arguments": json.dumps({
"script": 'return shell({command: "pwd", timeout_seconds: 3})'
})}
}]}
reason = "tool_calls"
chunk = {"id": "local-shell", "choices": [{"index": 0, "delta": delta, "finish_reason": reason}]}
self.respond(("data: " + json.dumps(chunk) + "\n\ndata: [DONE]\n\n").encode(), "text/event-stream")

def respond(self, body, content_type):
self.send_response(200)
self.send_header("Content-Type", content_type)
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)


if __name__ == "__main__":
with HTTPServer(("127.0.0.1", 0), Provider) as server:
print(server.server_port, flush=True)
server.serve_forever()
Loading