Overview
This client handles every event the harness can emit:- Streams assistant tokens in real time
- Shows tool calls and tool results inline
- Prompts the user to answer
ask_user_questioncalls - Prompts the user to allow or deny tool approvals
- Shows in-chat MCP OAuth prompts
- Handles parallel sub-agents
- Continues the conversation across turns
"""Terminal chat client for the TrueFoundry Agent Harness SDK."""
import argparse, json, os, sys
from typing import Optional, Union
from truefoundry_gateway_sdk.agents import (
AgentSession,
AgentSessionClient,
McpAuthRequiredEvent,
McpInitializeEvent,
ModelMessageDeltaEvent,
ModelMessageEvent,
SandboxCreatedEvent,
ThreadCreatedEvent,
ThreadDoneEvent,
ToolApprovalRequiredEvent,
ToolCall,
ToolResponseEvent,
ToolResponseRequiredEvent,
TurnCreatedEvent,
TurnDoneEvent,
TurnEvent,
TurnInputItem,
TurnStateError,
UserToolApprovalEvent,
UserToolResponseEvent,
is_event_delta,
merge_event_delta,
)
from truefoundry_gateway_sdk.types import ApprovalAllow, ApprovalDeny, UserMessage
PendingNextTurnRequest = Union[ToolApprovalRequiredEvent, ToolResponseRequiredEvent, McpAuthRequiredEvent]
class ChatSession:
"""Holds turn-to-turn state for one conversation."""
def __init__(self, session: AgentSession) -> None:
self.session = session
# id-keyed index of assembled events. Base model.message events are stored here
# and their model.message.delta fragments merge into them in place.
self.events: dict[str, TurnEvent] = {}
self.open_subagent_threads: set[str] = set()
self.thread_labels: dict[str, str] = {"main": "main"}
self.pending_next_turn_requests: list[PendingNextTurnRequest] = []
# id of the model.message currently streaming on each thread. A later event on
# the same thread with a different id means that message is done, so flush it.
self._streaming_ids: dict[str, str] = {}
self._main_streaming = False
def label(self, thread_id: str) -> str:
return self.thread_labels.get(thread_id, "main")
def _handle_model_message_delta(self, event: ModelMessageDeltaEvent) -> None:
# Stream main-thread text live; the flush happens later, driven by id change.
if event.thread_id == "main" and event.content:
if not self._main_streaming:
print("\nassistant: ", end="", flush=True)
self._main_streaming = True
print(event.content, end="", flush=True)
def _flush_model_message(self, msg: ModelMessageEvent) -> None:
# The message is complete: finalize the streamed line, print any non-streamed
# content, and emit the now fully assembled tool calls.
thread_id = msg.thread_id
label = self.label(thread_id)
if thread_id == "main":
if self._main_streaming:
print()
self._main_streaming = False
elif msg.content:
print(f"\nassistant: {msg.content}")
elif thread_id in self.open_subagent_threads and msg.content:
print(f"\n[{label}] {msg.content}")
for tool_call in msg.tool_calls or []:
if not tool_call.id or not tool_call.function.name:
continue
if tool_call.tool_info.name == "ask_user_question":
continue
self._print_tool_call(thread_id, tool_call)
def _print_tool_call(self, thread_id: str, tool_call: ToolCall) -> None:
label = self.label(thread_id)
try:
args = json.loads(tool_call.function.arguments or "{}")
except json.JSONDecodeError:
args = tool_call.function.arguments
args_preview = json.dumps(args)[:160] if isinstance(args, dict) else str(args)[:160]
prefix = "main" if thread_id == "main" else label
print(f"\n[{prefix} -> {tool_call.tool_info.name}] {args_preview}")
def _handle_tool_response(self, event: ToolResponseEvent) -> None:
label = self.label(event.thread_id)
prefix = "main" if event.thread_id == "main" else label
preview = event.content[:200].replace("\n", " ")
print(f"\n[{prefix} tool result] {preview}")
def _prompt_client_side_response(
self,
pending_event: ToolResponseRequiredEvent,
tool_call: ToolCall,
) -> UserToolResponseEvent:
if tool_call.tool_info.name == "ask_user_question":
try:
args = json.loads(tool_call.function.arguments or "{}")
except json.JSONDecodeError:
args = {}
question = args.get("question", "Answer:")
options = args.get("options") or []
print(f"\nassistant asks: {question}")
for idx, option in enumerate(options, 1):
print(f" {idx}. {option}")
if options:
print(" Or type a free-form answer.")
answer = input("you: ").strip()
if options and answer.isdigit() and 1 <= int(answer) <= len(options):
content = options[int(answer) - 1]
else:
content = answer
else:
print(f"\nclient-side tool response required: {tool_call.tool_info.name}")
print(f" arguments: {tool_call.function.arguments}")
content = input("you: ").strip()
return UserToolResponseEvent(
thread_id=pending_event.thread_id,
tool_call_id=tool_call.id,
content=content,
)
def _prompt_approval(
self,
pending_event: ToolApprovalRequiredEvent,
tool_call: ToolCall,
) -> UserToolApprovalEvent:
try:
arguments = json.loads(tool_call.function.arguments or "{}")
except json.JSONDecodeError:
arguments = tool_call.function.arguments
print(f"\napproval needed: {tool_call.tool_info.name}")
print(f" arguments: {json.dumps(arguments, indent=2)}")
choice = input(" allow / deny / reason for denial: ").strip().lower()
if choice in {"a", "allow", "yes", "y"}:
approval = ApprovalAllow()
elif choice in {"d", "deny", "no", "n"}:
approval = ApprovalDeny()
else:
approval = ApprovalDeny(reason=choice)
return UserToolApprovalEvent(
thread_id=pending_event.thread_id,
tool_call_id=tool_call.id,
approval=approval,
)
def _prompt_mcp_auth_required(self, event: McpAuthRequiredEvent) -> None:
print("\n[auth] MCP authorization required. Complete OAuth for each server:")
for server in event.mcp_servers:
print(f" {server.name}: {server.auth_url}")
while True:
answer = input(" Type 'yes' once you've authenticated: ").strip().lower()
if answer in {"y", "yes"}:
break
def build_next_turn_input(self) -> list[TurnInputItem]:
next_turn: list[UserToolApprovalEvent | UserToolResponseEvent] = []
pending = self.pending_next_turn_requests[:]
self.pending_next_turn_requests.clear()
for pending_event in pending:
if isinstance(pending_event, McpAuthRequiredEvent):
self._prompt_mcp_auth_required(pending_event)
continue
for tool_ref in pending_event.tool_calls:
# tool_ref.source_event_id points at the model.message that emitted this tool
# call; find the matching tool call inside it by id.
model_message = self.events[tool_ref.source_event_id]
tool_call = next(
tc for tc in (model_message.tool_calls or []) if tc.id == tool_ref.id
)
if isinstance(pending_event, ToolResponseRequiredEvent):
next_turn.append(
self._prompt_client_side_response(pending_event, tool_call)
)
else:
next_turn.append(
self._prompt_approval(pending_event, tool_call)
)
return next_turn
def _handle_event(self, event: object) -> Optional[str]:
# Keep the id-keyed index current: store base events, merge deltas into them.
if is_event_delta(event):
merge_event_delta(self.events[event.id], event)
elif isinstance(event, ModelMessageEvent):
self.events[event.id] = event
# A later event on a thread with a different id means that thread's streaming
# model.message is complete - flush it before handling the new event.
thread_id = getattr(event, "thread_id", None)
if thread_id is not None:
streaming_id = self._streaming_ids.get(thread_id)
if streaming_id is not None and streaming_id != event.id:
self._flush_model_message(self.events[streaming_id])
del self._streaming_ids[thread_id]
if isinstance(event, TurnCreatedEvent):
return None
if isinstance(event, ThreadCreatedEvent):
if event.thread_id != "main":
self.open_subagent_threads.add(event.thread_id)
name = event.agent_info.name
self.thread_labels[event.thread_id] = name
print(f"\n[subagent {name} started]")
return None
if isinstance(event, ThreadDoneEvent):
if event.thread_id != "main":
self.open_subagent_threads.discard(event.thread_id)
label = self.label(event.thread_id)
if event.state.status == "error":
print(f"\n[{label} error] {event.state.error}")
else:
print(f"[subagent {label} finished]")
return None
if isinstance(event, McpInitializeEvent):
for server in event.mcp_servers:
print(f"[mcp] connected to {server.name}")
return None
if isinstance(event, McpAuthRequiredEvent):
self.pending_next_turn_requests.append(event)
return None
if isinstance(event, SandboxCreatedEvent):
print(f"[sandbox] provisioned ({event.sandbox_id})")
return None
if isinstance(event, ModelMessageEvent):
# Base model.message: a new message starts streaming on this thread. Its
# content and tool_calls fill in as the matching deltas arrive.
self._streaming_ids[event.thread_id] = event.id
return None
if isinstance(event, ModelMessageDeltaEvent):
self._handle_model_message_delta(event)
return None
if isinstance(event, ToolResponseEvent):
self._handle_tool_response(event)
return None
if isinstance(event, ToolResponseRequiredEvent):
self.pending_next_turn_requests.append(event)
return None
if isinstance(event, ToolApprovalRequiredEvent):
self.pending_next_turn_requests.append(event)
return None
if isinstance(event, TurnDoneEvent):
# Turn-level event (thread_id is null), so flush any messages still
# streaming - the last message on each thread has no later event to
# trigger its flush.
for streaming_id in self._streaming_ids.values():
self._flush_model_message(self.events[streaming_id])
self._streaming_ids.clear()
# event.state is the terminal TurnState (same as get_turn().state).
if isinstance(event.state, TurnStateError):
print(f"\n[error] {event.state.message}")
return event.state.status
print(f"\n[unknown event] {event}")
return None
def run_turn(self, input_items: list[TurnInputItem] | None = None) -> str:
last_status = "done"
turn = self.session.prepare_turn(input=input_items)
for data in turn.execute(stream=True):
status = self._handle_event(data.event)
if status is not None:
last_status = status
return last_status
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(
description="Terminal chat client for the TrueFoundry Agent Harness.",
)
parser.add_argument(
"agent_name",
nargs="?",
default=os.environ.get("AGENT_NAME"),
help="Registered agent name (default: AGENT_NAME environment variable)",
)
args = parser.parse_args()
if not args.agent_name:
parser.error("agent_name is required (pass as argument or set AGENT_NAME)")
return args
def main() -> None:
args = parse_args()
base_url = os.environ.get("TFY_GATEWAY_URL")
api_key = os.environ.get("TFY_API_KEY")
if not base_url or not api_key:
print("Set TFY_GATEWAY_URL and TFY_API_KEY environment variables.")
sys.exit(1)
client = AgentSessionClient(
base_url=base_url,
api_key=api_key,
)
session = client.create_session(agent_name=args.agent_name)
chat = ChatSession(session)
print("Type a message, or Ctrl-D to exit.")
while True:
while chat.pending_next_turn_requests:
chat.run_turn(chat.build_next_turn_input())
try:
text = input("\nyou: ").strip()
except EOFError:
print()
break
if not text:
continue
chat.run_turn([UserMessage(content=text)])
if __name__ == "__main__":
main()
/** Terminal chat client for the TrueFoundry Agent Harness SDK. */
import { createInterface, type Interface } from "node:readline/promises";
import { argv, env, exit, stdin, stdout } from "node:process";
import {
AgentSessionClient,
isEventDelta,
mergeEventDelta,
type AgentSession,
type McpAuthRequiredEvent,
type ModelMessageEvent,
type ModelMessageDeltaEvent,
type ToolApprovalRequiredEvent,
type ToolCall,
type ToolResponseEvent,
type ToolResponseRequiredEvent,
type TurnEvent,
type TurnInputItem,
type TurnStreamingEvent,
type UserToolApprovalEvent,
type UserToolResponseEvent,
} from "truefoundry-gateway-sdk/agents";
type TurnTerminalStatus = "done" | "cancelled" | "error";
type PendingNextTurnRequest =
| ToolApprovalRequiredEvent
| ToolResponseRequiredEvent
| McpAuthRequiredEvent;
/** Holds turn-to-turn state for one conversation. */
class ChatSession {
readonly pendingNextTurnRequests: PendingNextTurnRequest[] = [];
// id-keyed index of assembled events. Base model.message events are stored here and
// their model.message.delta fragments merge into them in place.
private readonly events = new Map<string, TurnEvent>();
private readonly openSubagentThreads = new Set<string>();
private readonly threadLabels = new Map<string, string>([["main", "main"]]);
// id of the model.message currently streaming on each thread. A later event on the
// same thread with a different id means that message is done, so flush it.
private readonly streamingIds = new Map<string, string>();
private mainStreaming = false;
constructor(
private readonly session: AgentSession,
private readonly rl: Interface,
) {}
private label(threadId: string): string {
return this.threadLabels.get(threadId) ?? "main";
}
private handleModelMessageDelta(event: ModelMessageDeltaEvent): void {
// Stream main-thread text live; the flush happens later, driven by id change.
if (event.threadId === "main" && event.content) {
if (!this.mainStreaming) {
stdout.write("\nassistant: ");
this.mainStreaming = true;
}
stdout.write(event.content);
}
}
private flushModelMessage(msg: ModelMessageEvent): void {
// The message is complete: finalize the streamed line, print any non-streamed
// content, and emit the now fully assembled tool calls.
const threadId = msg.threadId;
const label = this.label(threadId);
if (threadId === "main") {
if (this.mainStreaming) {
stdout.write("\n");
this.mainStreaming = false;
} else if (msg.content) {
console.log(`\nassistant: ${msg.content}`);
}
} else if (this.openSubagentThreads.has(threadId) && msg.content) {
console.log(`\n[${label}] ${msg.content}`);
}
for (const toolCall of msg.toolCalls) {
if (!toolCall.id || !toolCall.function.name) {
continue;
}
if (toolCall.toolInfo?.name === "ask_user_question") {
continue;
}
this.printToolCall(threadId, toolCall);
}
}
private printToolCall(threadId: string, toolCall: ToolCall): void {
const label = this.label(threadId);
let argsPreview: string;
try {
argsPreview = JSON.stringify(JSON.parse(toolCall.function.arguments || "{}")).slice(0, 160);
} catch {
argsPreview = (toolCall.function.arguments ?? "").slice(0, 160);
}
const prefix = threadId === "main" ? "main" : label;
console.log(`\n[${prefix} -> ${toolCall.toolInfo?.name}] ${argsPreview}`);
}
private handleToolResponse(event: ToolResponseEvent): void {
const label = this.label(event.threadId);
const prefix = event.threadId === "main" ? "main" : label;
const preview = event.content.slice(0, 200).replace(/\n/g, " ");
console.log(`\n[${prefix} tool result] ${preview}`);
}
private async promptClientSideResponse(
pendingEvent: ToolResponseRequiredEvent,
toolCall: ToolCall,
): Promise<UserToolResponseEvent> {
let content: string;
if (toolCall.toolInfo?.name === "ask_user_question") {
let args: Record<string, unknown> = {};
try {
args = JSON.parse(toolCall.function.arguments || "{}");
} catch {
args = {};
}
const question = (args.question as string) ?? "Answer:";
const options = (args.options as string[]) ?? [];
console.log(`\nassistant asks: ${question}`);
options.forEach((option, idx) => console.log(` ${idx + 1}. ${option}`));
if (options.length) {
console.log(" Or type a free-form answer.");
}
const answer = (await this.rl.question("you: ")).trim();
const choice = Number(answer);
content =
options.length && Number.isInteger(choice) && choice >= 1 && choice <= options.length
? options[choice - 1]
: answer;
} else {
console.log(`\nclient-side tool response required: ${toolCall.toolInfo?.name}`);
console.log(` arguments: ${toolCall.function.arguments}`);
content = (await this.rl.question("you: ")).trim();
}
return {
type: "user.tool_response",
threadId: pendingEvent.threadId,
toolCallId: toolCall.id,
content,
};
}
private async promptApproval(
pendingEvent: ToolApprovalRequiredEvent,
toolCall: ToolCall,
): Promise<UserToolApprovalEvent> {
let parsedArgs: unknown;
try {
parsedArgs = JSON.parse(toolCall.function.arguments || "{}");
} catch {
parsedArgs = toolCall.function.arguments;
}
console.log(`\napproval needed: ${toolCall.toolInfo?.name}`);
console.log(` arguments: ${JSON.stringify(parsedArgs, null, 2)}`);
const choice = (await this.rl.question(" allow / deny / reason for denial: ")).trim().toLowerCase();
let approval: UserToolApprovalEvent["approval"];
if (["a", "allow", "yes", "y"].includes(choice)) {
approval = { status: "allow" };
} else if (["d", "deny", "no", "n"].includes(choice)) {
approval = { status: "deny" };
} else {
approval = { status: "deny", reason: choice };
}
return {
type: "user.tool_approval",
threadId: pendingEvent.threadId,
toolCallId: toolCall.id,
approval,
};
}
private async promptMcpAuthRequired(event: McpAuthRequiredEvent): Promise<void> {
console.log("\n[auth] MCP authorization required. Complete OAuth for each server:");
for (const server of event.mcpServers) {
console.log(` ${server.name}: ${server.authUrl}`);
}
for (;;) {
const answer = (await this.rl.question(" Type 'yes' once you've authenticated: "))
.trim()
.toLowerCase();
if (answer === "y" || answer === "yes") {
break;
}
}
}
async buildNextTurnInput(): Promise<TurnInputItem[]> {
const nextTurn: (UserToolApprovalEvent | UserToolResponseEvent)[] = [];
const pending = [...this.pendingNextTurnRequests];
this.pendingNextTurnRequests.length = 0;
for (const pendingEvent of pending) {
if (pendingEvent.type === "mcp.auth_required") {
await this.promptMcpAuthRequired(pendingEvent);
continue;
}
for (const toolRef of pendingEvent.toolCalls) {
// toolRef.sourceEventId points at the model.message that emitted this tool call;
// find the matching tool call inside it by id.
const modelMessage = this.events.get(toolRef.sourceEventId) as ModelMessageEvent;
const toolCall = modelMessage.toolCalls!.find((tc) => tc.id === toolRef.id)!;
if (pendingEvent.type === "tool.response_required") {
nextTurn.push(await this.promptClientSideResponse(pendingEvent, toolCall));
} else {
nextTurn.push(await this.promptApproval(pendingEvent, toolCall));
}
}
}
return nextTurn;
}
private handleEvent(event: TurnStreamingEvent): TurnTerminalStatus | null {
// Keep the id-keyed index current: store base events, merge deltas into them.
if (isEventDelta(event)) {
mergeEventDelta(this.events.get(event.id) as ModelMessageEvent, event);
} else if (event.type === "model.message") {
this.events.set(event.id, event);
}
// A later event on a thread with a different id means that thread's streaming
// model.message is complete - flush it before handling the new event.
const threadId = event.threadId;
if (threadId != null) {
const streamingId = this.streamingIds.get(threadId);
if (streamingId != null && streamingId !== event.id) {
this.flushModelMessage(this.events.get(streamingId) as ModelMessageEvent);
this.streamingIds.delete(threadId);
}
}
switch (event.type) {
case "turn.created":
return null;
case "thread.created":
if (event.threadId !== "main") {
this.openSubagentThreads.add(event.threadId);
const name = event.agentInfo.name;
this.threadLabels.set(event.threadId, name);
console.log(`\n[subagent ${name} started]`);
}
return null;
case "thread.done":
if (event.threadId !== "main") {
this.openSubagentThreads.delete(event.threadId);
const label = this.label(event.threadId);
if (event.state.status === "error") {
console.log(`\n[${label} error] ${event.state.error}`);
} else {
console.log(`[subagent ${label} finished]`);
}
}
return null;
case "mcp.initialize":
for (const server of event.mcpServers) {
console.log(`[mcp] connected to ${server.name}`);
}
return null;
case "mcp.auth_required":
this.pendingNextTurnRequests.push(event);
return null;
case "sandbox.created":
console.log(`[sandbox] provisioned (${event.sandboxId})`);
return null;
case "model.message":
// Base model.message: a new message starts streaming on this thread. Its
// content and tool_calls fill in as the matching deltas arrive.
this.streamingIds.set(event.threadId, event.id);
return null;
case "model.message.delta":
this.handleModelMessageDelta(event);
return null;
case "tool.response":
this.handleToolResponse(event);
return null;
case "tool.response_required":
this.pendingNextTurnRequests.push(event);
return null;
case "tool.approval_required":
this.pendingNextTurnRequests.push(event);
return null;
case "turn.done":
// Turn-level event (threadId is null), so flush any messages still
// streaming - the last message on each thread has no later event to
// trigger its flush.
for (const streamingId of this.streamingIds.values()) {
this.flushModelMessage(this.events.get(streamingId) as ModelMessageEvent);
}
this.streamingIds.clear();
// event.state is the terminal TurnState (same as execute({ stream: false })).
if (event.state.status === "error") {
console.log(`\n[error] ${event.state.message}`);
}
return event.state.status;
default:
console.log(`\n[unknown event] ${JSON.stringify(event)}`);
return null;
}
}
async runTurn(inputItems?: TurnInputItem[]): Promise<TurnTerminalStatus> {
let lastStatus: TurnTerminalStatus = "done";
const turn = this.session.prepareTurn({ input: inputItems });
for await (const data of turn.execute({ stream: true })) {
const status = this.handleEvent(data.event);
if (status != null) {
lastStatus = status;
}
}
return lastStatus;
}
}
async function main(): Promise<void> {
const agentName = argv[2] ?? env.AGENT_NAME;
if (!agentName) {
console.error("agent name is required (pass as argument or set AGENT_NAME)");
exit(1);
}
const baseUrl = env.TFY_GATEWAY_URL;
const apiKey = env.TFY_API_KEY;
if (!baseUrl || !apiKey) {
console.error("Set TFY_GATEWAY_URL and TFY_API_KEY environment variables.");
exit(1);
}
const client = new AgentSessionClient({ apiKey, baseUrl });
const session = await client.createSession({ agentName });
const rl = createInterface({ input: stdin, output: stdout });
const chat = new ChatSession(session, rl);
console.log("Type a message, or Ctrl-D to exit.");
for (;;) {
while (chat.pendingNextTurnRequests.length) {
await chat.runTurn(await chat.buildNextTurnInput());
}
// rl.question rejects when the input stream closes (Ctrl-D); treat that as exit.
const text = await rl.question("\nyou: ").then((t) => t.trim(), () => null);
if (text == null) {
console.log();
break;
}
if (!text) {
continue;
}
await chat.runTurn([{ type: "user.message", content: text }]);
}
rl.close();
}
main();
Sample run
$ python example.py support-bot
Type a message, or Ctrl-D to exit.
you: Refund invoice INV-2031 for tenant acme.
[mcp] connected to truefoundry-mcp
[sandbox] provisioned (my-tenant.a1b2c3d4)
assistant asks: Which refund reason should I record?
1. Customer churn
2. Billing error
3. Goodwill credit
Or type a free-form answer.
you: 2
assistant: Found invoice INV-2031 for $1,240.00 on 2026-04-12.
approval needed: process_refund
arguments: {
"invoice_id": "INV-2031",
"amount": 1240.0,
"reason": "Billing error"
}
allow / deny / reason for denial: allow
[main tool result] {"refund_id":"RF-9914","status":"completed"}
assistant: Refund RF-9914 of $1,240.00 was processed for INV-2031.
you: ^D
Adapting the flow
A few common variations on top of the same skeleton:- Persisted history. Save
session.idafter the first turn. On the next process start, callclient.getSession({ sessionId })and uselistTurns()withlistEvents()to rebuild UI state, orturn.stream()if the latest turn is still running. - JSON output for piping. Replace the
printcalls inside_handle_eventwithjson.dumps(event.model_dump())to emit one event per line for downstream tools. - Generative UI. When you detect a fenced
```openuiblock in assembledmodel.messagecontent, hand the block to the OpenUI React renderer instead of printing it. Everything else stays the same. - Browser / web UI. The same event shape is delivered over Server-Sent Events on
POST /sessions/{session_id}/turns. Replaceturn.execute({ stream: true })with an SSE consumer; the handlers above are unchanged.