Skip to content

Commit f853785

Browse files
committed
feat(harness): add a Gemini CLI harness (tap, turn, and init templates)
Adds Gemini CLI as a framework harness alongside Claude Code and Codex: - convert_gemini_cli_to_agentex_events: maps the CLI's stream-json events (init, message deltas, tool_use, tool_result, error, result; schema per packages/core/src/output/types.ts in google-gemini/gemini-cli) onto the canonical StreamTaskMessage* stream. Assistant deltas open one text slot that closes on the next tool event, the result, or end of stream, so every Start has a Done; tool requests and results pair by tool_id. - GeminiCliTurn: HarnessTurn wrapper exposing session_id and model from the init event and normalising result.stats into TurnUsage. - Both exported from agentex.lib.adk. - agentex init templates sync-gemini-cli, default-gemini-cli and temporal-gemini-cli (registered in TemplateType, file map and menus), cloned from the Claude Code templates: prompt passed via -p with stdin closed (the CLI reads stdin to EOF in headless mode), optional GEMINI_MODEL, GEMINI_API_KEY credential, npm install -g @google/gemini-cli in the Dockerfile. Turns are independent prompts: the CLI's --resume takes latest/index, not a session id. - Tests: tap (text deltas, whole messages, tools, errors, callbacks, source close on cancel), turn (usage mapping, protocol), harness end to end through UnifiedEmitter with span derivation; template suite covers the three new templates. Offline tests only; a live smoke run needs a Gemini API key. Claude-Session: https://claude.ai/code/session_01HCVKnA7LeJZ44nxZz1uzF3
1 parent 0db6037 commit f853785

43 files changed

Lines changed: 3326 additions & 0 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎src/agentex/lib/adk/__init__.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@
2222
)
2323
from agentex.lib.adk._modules._codex_sync import convert_codex_to_agentex_events
2424
from agentex.lib.adk._modules._codex_turn import CodexTurn, codex_usage_to_turn_usage
25+
from agentex.lib.adk._modules._gemini_cli_sync import convert_gemini_cli_to_agentex_events
26+
from agentex.lib.adk._modules._gemini_cli_turn import GeminiCliTurn, gemini_cli_usage_to_turn_usage
2527
from agentex.lib.adk._modules.events import EventsModule
2628
from agentex.lib.adk._modules.messages import MessagesModule
2729
from agentex.lib.adk._modules.state import StateModule
@@ -101,6 +103,10 @@
101103
"convert_codex_to_agentex_events",
102104
"CodexTurn",
103105
"codex_usage_to_turn_usage",
106+
# Gemini CLI
107+
"convert_gemini_cli_to_agentex_events",
108+
"GeminiCliTurn",
109+
"gemini_cli_usage_to_turn_usage",
104110
# Unified harness surface (AGX1-375)
105111
"UnifiedEmitter",
106112
"SpanTracer",
Lines changed: 276 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,276 @@
1+
"""Gemini CLI stream-json parser tap for the unified harness surface.
2+
3+
Converts the newline-delimited JSON events emitted by
4+
``gemini -p <prompt> --output-format stream-json`` into the canonical
5+
``StreamTaskMessage*`` stream consumed by the Agentex harness.
6+
7+
Event → canonical mapping
8+
-------------------------
9+
init
10+
Fires ``on_init`` with the raw event (``session_id``, ``model``). Nothing
11+
is emitted: session metadata is a provider concern.
12+
13+
message (role=user)
14+
Ignored. The CLI echoes the prompt back as the first message.
15+
16+
message (role=assistant)
17+
The CLI streams the answer as ``delta: true`` chunks. The first chunk
18+
opens a text slot (Start(TextContent)); every chunk is a Delta(TextDelta).
19+
The slot is closed (Done) when a ``tool_use``, ``tool_result`` or
20+
``result`` event arrives, or when the stream ends. A non-delta assistant
21+
message whose content matches the open slot closes it; otherwise it is
22+
delivered as Start + Delta + Done.
23+
24+
tool_use
25+
Start(ToolRequestContent) + Done. ``tool_id`` → ``tool_call_id``,
26+
``tool_name`` → ``name``, ``parameters`` → ``arguments``.
27+
28+
tool_result
29+
Full(ToolResponseContent) keyed by ``tool_id``. ``output`` (or the error
30+
message when ``status == "error"``) becomes ``content["result"]``;
31+
``is_error`` is set for error results.
32+
33+
error
34+
Logged (``severity`` + ``message``). Nothing is emitted.
35+
36+
result
37+
Closes any open text slot, then fires ``on_result`` with the raw event so
38+
the caller can read ``stats`` (tokens, duration, tool calls).
39+
40+
Reference: ``packages/core/src/output/types.ts`` in google-gemini/gemini-cli.
41+
"""
42+
43+
from __future__ import annotations
44+
45+
import json
46+
from typing import Any, Callable, Awaitable, AsyncIterator
47+
48+
from agentex.lib.utils.logging import make_logger
49+
from agentex.types.text_content import TextContent
50+
from agentex.types.task_message_delta import TextDelta
51+
from agentex.types.task_message_update import (
52+
StreamTaskMessageDone,
53+
StreamTaskMessageFull,
54+
StreamTaskMessageDelta,
55+
StreamTaskMessageStart,
56+
)
57+
from agentex.types.tool_request_content import ToolRequestContent
58+
from agentex.types.tool_response_content import ToolResponseContent
59+
60+
logger = make_logger(__name__)
61+
62+
_MAX_RESULT_LENGTH = 4000
63+
64+
65+
def _truncate(text: str) -> str:
66+
return str(text)[:_MAX_RESULT_LENGTH]
67+
68+
69+
async def convert_gemini_cli_to_agentex_events(
70+
lines: AsyncIterator[str | dict[str, Any]],
71+
on_result: Callable[[dict[str, Any]], Awaitable[None]] | None = None,
72+
on_init: Callable[[dict[str, Any]], Awaitable[None]] | None = None,
73+
) -> AsyncIterator[StreamTaskMessageStart | StreamTaskMessageDelta | StreamTaskMessageFull | StreamTaskMessageDone]:
74+
"""Public tap: convert a Gemini CLI ``stream-json`` line stream to events.
75+
76+
Thin wrapper over :func:`_convert_gemini_cli_impl` that owns the
77+
cancellation backstop: a ``finally`` closes the underlying ``lines``
78+
iterator (when it exposes ``aclose``) whenever this generator is closed,
79+
including on the ``GeneratorExit``/``CancelledError`` raised when the
80+
consuming task is cancelled mid-turn, so the CLI stdout handle and
81+
subprocess are not leaked.
82+
"""
83+
inner = _convert_gemini_cli_impl(lines, on_result=on_result, on_init=on_init)
84+
try:
85+
async for event in inner:
86+
yield event
87+
finally:
88+
inner_aclose = getattr(inner, "aclose", None)
89+
if inner_aclose is not None:
90+
await inner_aclose()
91+
aclose = getattr(lines, "aclose", None)
92+
if aclose is not None:
93+
await aclose()
94+
95+
96+
async def _convert_gemini_cli_impl(
97+
lines: AsyncIterator[str | dict[str, Any]],
98+
on_result: Callable[[dict[str, Any]], Awaitable[None]] | None = None,
99+
on_init: Callable[[dict[str, Any]], Awaitable[None]] | None = None,
100+
) -> AsyncIterator[StreamTaskMessageStart | StreamTaskMessageDelta | StreamTaskMessageFull | StreamTaskMessageDone]:
101+
"""Convert a Gemini CLI ``stream-json`` line stream into ``StreamTaskMessage*`` events.
102+
103+
Each item in ``lines`` is either a raw JSON string (as read from the CLI's
104+
stdout) or an already-parsed dict. Empty strings are skipped; unparseable
105+
JSON is logged and skipped. The event → canonical mapping is documented in
106+
this module's docstring.
107+
"""
108+
next_index = 0
109+
tool_call_count = 0
110+
111+
# One open assistant text slot at a time: the CLI streams the answer as
112+
# ``delta: true`` message chunks with no explicit start/stop markers.
113+
text_open = False
114+
text_index: int | None = None
115+
text_buf = ""
116+
117+
def _close_text() -> StreamTaskMessageDone | None:
118+
nonlocal text_open, text_index, text_buf
119+
if not text_open or text_index is None:
120+
return None
121+
done = StreamTaskMessageDone(type="done", index=text_index)
122+
text_open = False
123+
text_index = None
124+
text_buf = ""
125+
return done
126+
127+
async for raw in lines:
128+
if not raw:
129+
continue
130+
131+
if isinstance(raw, dict):
132+
evt = raw
133+
else:
134+
line = raw.strip()
135+
if not line:
136+
continue
137+
try:
138+
evt = json.loads(line)
139+
except json.JSONDecodeError:
140+
logger.debug("gemini-cli: skipping non-JSON line: %r", line[:120])
141+
continue
142+
143+
if not isinstance(evt, dict):
144+
continue
145+
evt_type = evt.get("type", "")
146+
147+
if evt_type == "message":
148+
if evt.get("role") != "assistant":
149+
continue # the CLI echoes the user prompt; nothing to emit
150+
content = evt.get("content", "")
151+
if not isinstance(content, str) or not content:
152+
continue
153+
154+
if evt.get("delta"):
155+
if not text_open:
156+
text_open = True
157+
text_index = next_index
158+
next_index += 1
159+
text_buf = ""
160+
yield StreamTaskMessageStart(
161+
type="start",
162+
index=text_index,
163+
content=TextContent(type="text", author="agent", content=""),
164+
)
165+
text_buf += content
166+
assert text_index is not None
167+
yield StreamTaskMessageDelta(
168+
type="delta",
169+
index=text_index,
170+
delta=TextDelta(type="text", text_delta=content),
171+
)
172+
continue
173+
174+
# A complete (non-delta) assistant message. If it materialises the
175+
# slot we are already streaming, just close the slot; otherwise
176+
# deliver it as its own Start + Delta + Done.
177+
if text_open and text_buf and content.startswith(text_buf):
178+
done = _close_text()
179+
if done is not None:
180+
yield done
181+
continue
182+
done = _close_text()
183+
if done is not None:
184+
yield done
185+
msg_index = next_index
186+
next_index += 1
187+
yield StreamTaskMessageStart(
188+
type="start",
189+
index=msg_index,
190+
content=TextContent(type="text", author="agent", content=""),
191+
)
192+
yield StreamTaskMessageDelta(
193+
type="delta",
194+
index=msg_index,
195+
delta=TextDelta(type="text", text_delta=content),
196+
)
197+
yield StreamTaskMessageDone(type="done", index=msg_index)
198+
199+
elif evt_type == "tool_use":
200+
done = _close_text()
201+
if done is not None:
202+
yield done
203+
tool_call_count += 1
204+
tool_id = evt.get("tool_id") or f"tool_{tool_call_count}"
205+
name = evt.get("tool_name") or "unknown"
206+
arguments = evt.get("parameters")
207+
if not isinstance(arguments, dict):
208+
arguments = {}
209+
msg_index = next_index
210+
next_index += 1
211+
yield StreamTaskMessageStart(
212+
type="start",
213+
index=msg_index,
214+
content=ToolRequestContent(
215+
type="tool_request",
216+
author="agent",
217+
tool_call_id=str(tool_id),
218+
name=str(name),
219+
arguments=arguments,
220+
),
221+
)
222+
yield StreamTaskMessageDone(type="done", index=msg_index)
223+
224+
elif evt_type == "tool_result":
225+
done = _close_text()
226+
if done is not None:
227+
yield done
228+
tool_id = str(evt.get("tool_id") or "")
229+
is_error = evt.get("status") == "error"
230+
output = evt.get("output")
231+
if output is None:
232+
error = evt.get("error") or {}
233+
output = error.get("message", "") if isinstance(error, dict) else str(error)
234+
result_content: dict[str, Any] = {"result": _truncate(str(output))}
235+
if is_error:
236+
result_content["is_error"] = True
237+
msg_index = next_index
238+
next_index += 1
239+
yield StreamTaskMessageFull(
240+
type="full",
241+
index=msg_index,
242+
content=ToolResponseContent(
243+
type="tool_response",
244+
author="agent",
245+
tool_call_id=tool_id,
246+
name="",
247+
content=result_content,
248+
),
249+
)
250+
251+
elif evt_type == "init":
252+
if on_init is not None:
253+
await on_init(evt)
254+
255+
elif evt_type == "error":
256+
logger.warning(
257+
"gemini-cli: %s: %s",
258+
evt.get("severity", "error"),
259+
str(evt.get("message", ""))[:300],
260+
)
261+
262+
elif evt_type == "result":
263+
done = _close_text()
264+
if done is not None:
265+
yield done
266+
if on_result is not None:
267+
await on_result(evt)
268+
269+
else:
270+
logger.debug("gemini-cli: unhandled event type %r", evt_type)
271+
272+
# Stream ended without a result event (truncated / interrupted): close the
273+
# slot so every Start has a matching Done.
274+
done = _close_text()
275+
if done is not None:
276+
yield done

0 commit comments

Comments
 (0)