跳转到内容
搜索文档

可调用方法

最后更新 查看 MarkdownAgent 设置

可调用方法(callable methods)允许客户端通过 WebSocket 使用 RPC(远程过程调用)调用 Agent 方法。为方法添加 @callable() 装饰器,即可向浏览器、移动应用或其他服务等外部客户端暴露它们。

概览

import { Agent, callable } from "agents";

export class MyAgent extends Agent {
	@callable()
	async greet(name) {
		return `Hello, ${name}!`;
	}
}
import { Agent, callable } from "agents";

export class MyAgent extends Agent {
	@callable()
	async greet(name: string): Promise<string> {
		return `Hello, ${name}!`;
	}
}
// Client
const result = await agent.stub.greet("World");
console.log(result); // "Hello, World!"
// Client
const result = await agent.stub.greet("World");
console.log(result); // "Hello, World!"

工作原理

sequenceDiagram
    participant Client
    participant Agent
    Client->>Agent: agent.stub.greet("World")
    Note right of Agent: 检查 @callable<br/>执行方法
    Agent-->>Client: "Hello, World!"

何时使用 @callable()

场景 用法
浏览器/移动端调用 Agent @callable()
外部服务调用 Agent @callable()
同一 Worker 内调用 Agent Durable Object RPC(无需装饰器)
Agent 调用另一 Agent 通过 getAgentByName() 的 Durable Object RPC

@callable() 装饰器专用于来自外部客户端的 WebSocket RPC。在同一 Worker 内或 Agent 间调用时,直接使用标准 Durable Object RPC

基本用法

定义可调用方法

为要暴露的方法添加 @callable() 装饰器:

import { Agent, callable } from "agents";

export class CounterAgent extends Agent {
	initialState = { count: 0, items: [] };

	@callable()
	increment() {
		this.setState({ ...this.state, count: this.state.count + 1 });
		return this.state.count;
	}

	@callable()
	decrement() {
		this.setState({ ...this.state, count: this.state.count - 1 });
		return this.state.count;
	}

	@callable()
	async addItem(item) {
		this.setState({ ...this.state, items: [...this.state.items, item] });
		return this.state.items;
	}

	@callable()
	getStats() {
		return {
			count: this.state.count,
			itemCount: this.state.items.length,
		};
	}
}
import { Agent, callable } from "agents";

export type CounterState = {
	count: number;
	items: string[];
};

export class CounterAgent extends Agent<Env, CounterState> {
	initialState: CounterState = { count: 0, items: [] };

	@callable()
	increment(): number {
		this.setState({ ...this.state, count: this.state.count + 1 });
		return this.state.count;
	}

	@callable()
	decrement(): number {
		this.setState({ ...this.state, count: this.state.count - 1 });
		return this.state.count;
	}

	@callable()
	async addItem(item: string): Promise<string[]> {
		this.setState({ ...this.state, items: [...this.state.items, item] });
		return this.state.items;
	}

	@callable()
	getStats(): { count: number; itemCount: number } {
		return {
			count: this.state.count,
			itemCount: this.state.items.length,
		};
	}
}

从客户端调用

有两种方式从客户端调用方法:

使用 agent.stub(推荐):

// Clean, typed syntax
const count = await agent.stub.increment();
const items = await agent.stub.addItem("new item");
const stats = await agent.stub.getStats();
// Clean, typed syntax
const count = await agent.stub.increment();
const items = await agent.stub.addItem("new item");
const stats = await agent.stub.getStats();

使用 agent.call()

// Explicit method name as string
const count = await agent.call("increment");
const items = await agent.call("addItem", ["new item"]);
const stats = await agent.call("getStats");
// Explicit method name as string
const count = await agent.call("increment");
const items = await agent.call("addItem", ["new item"]);
const stats = await agent.call("getStats");

stub 代理提供更好的 ergonomics 与 TypeScript 支持。

方法签名

可序列化类型

参数与返回值须为 JSON 可序列化:

// Valid - primitives and plain objects
class MyAgent extends Agent {
	@callable()
	processData(input) {
		return { result: true };
	}
}

// Valid - arrays
class MyAgent extends Agent {
	@callable()
	processItems(items) {
		return items.map((item) => item.length);
	}
}

// Invalid - non-serializable types
// Functions, Dates, Maps, Sets, etc. cannot be serialized
// Valid - primitives and plain objects
class MyAgent extends Agent {
	@callable()
	processData(input: { name: string; count: number }): { result: boolean } {
		return { result: true };
	}
}

// Valid - arrays
class MyAgent extends Agent {
	@callable()
	processItems(items: string[]): number[] {
		return items.map((item) => item.length);
	}
}

// Invalid - non-serializable types
// Functions, Dates, Maps, Sets, etc. cannot be serialized

异步方法

同步与异步方法均可:

// Sync method
class MyAgent extends Agent {
	@callable()
	add(a, b) {
		return a + b;
	}
}

// Async method
class MyAgent extends Agent {
	@callable()
	async fetchUser(id) {
		const user = await this.sql`SELECT * FROM users WHERE id = ${id}`;
		return user[0];
	}
}
// Sync method
class MyAgent extends Agent {
	@callable()
	add(a: number, b: number): number {
		return a + b;
	}
}

// Async method
class MyAgent extends Agent {
	@callable()
	async fetchUser(id: string): Promise<User> {
		const user = await this.sql`SELECT * FROM users WHERE id = ${id}`;
		return user[0];
	}
}

无返回值方法

不返回值的方法:

class MyAgent extends Agent {
	@callable()
	async logEvent(event) {
		await this.sql`INSERT INTO events (name) VALUES (${event})`;
	}
}
class MyAgent extends Agent {
	@callable()
	async logEvent(event: string): Promise<void> {
		await this.sql`INSERT INTO events (name) VALUES (${event})`;
	}
}

客户端上这些方法仍返回 Promise,在方法完成时 resolve:

await agent.stub.logEvent("user-clicked");
// Resolves when the server confirms execution
await agent.stub.logEvent("user-clicked");
// Resolves when the server confirms execution

流式响应

对随时间产生数据的方法(如 AI 文本生成),使用流式:

定义流式方法

import { Agent, callable } from "agents";

export class AIAgent extends Agent {
	@callable({ streaming: true })
	async generateText(stream, prompt) {
		// First parameter is always StreamingResponse for streaming methods

		for await (const chunk of this.llm.stream(prompt)) {
			stream.send(chunk); // Send each chunk to the client
		}

		stream.end(); // Signal completion
	}

	@callable({ streaming: true })
	async streamNumbers(stream, count) {
		for (let i = 0; i < count; i++) {
			stream.send(i);
			await new Promise((resolve) => setTimeout(resolve, 100));
		}
		stream.end(count); // Optional final value
	}
}
import { Agent, callable, type StreamingResponse } from "agents";

export class AIAgent extends Agent {
	@callable({ streaming: true })
	async generateText(stream: StreamingResponse, prompt: string) {
		// First parameter is always StreamingResponse for streaming methods

		for await (const chunk of this.llm.stream(prompt)) {
			stream.send(chunk); // Send each chunk to the client
		}

		stream.end(); // Signal completion
	}

	@callable({ streaming: true })
	async streamNumbers(stream: StreamingResponse, count: number) {
		for (let i = 0; i < count; i++) {
			stream.send(i);
			await new Promise((resolve) => setTimeout(resolve, 100));
		}
		stream.end(count); // Optional final value
	}
}

在客户端消费流

// Preferred format (supports timeout and other options)
await agent.call("generateText", [prompt], {
	stream: {
		onChunk: (chunk) => {
			// Called for each chunk
			appendToOutput(chunk);
		},
		onDone: (finalValue) => {
			// Called when stream ends
			console.log("Stream complete", finalValue);
		},
		onError: (error) => {
			// Called if an error occurs
			console.error("Stream error:", error);
		},
	},
});

// Legacy format (still supported for backward compatibility)
await agent.call("generateText", [prompt], {
	onChunk: (chunk) => appendToOutput(chunk),
	onDone: (finalValue) => console.log("Done", finalValue),
	onError: (error) => console.error("Error:", error),
});
// Preferred format (supports timeout and other options)
await agent.call("generateText", [prompt], {
	stream: {
		onChunk: (chunk) => {
			// Called for each chunk
			appendToOutput(chunk);
		},
		onDone: (finalValue) => {
			// Called when stream ends
			console.log("Stream complete", finalValue);
		},
		onError: (error) => {
			// Called if an error occurs
			console.error("Stream error:", error);
		},
	},
});

// Legacy format (still supported for backward compatibility)
await agent.call("generateText", [prompt], {
	onChunk: (chunk) => appendToOutput(chunk),
	onDone: (finalValue) => console.log("Done", finalValue),
	onError: (error) => console.error("Error:", error),
});

StreamingResponse API

方法 描述
send(chunk) 向客户端发送 chunk
end(finalChunk?) 结束流,可选最终值
error(message) 向客户端发送错误并关闭流
class MyAgent extends Agent {
	@callable({ streaming: true })
	async processWithProgress(stream, items) {
		for (let i = 0; i < items.length; i++) {
			await this.process(items[i]);
			stream.send({ progress: (i + 1) / items.length, item: items[i] });
		}
		stream.end({ completed: true, total: items.length });
	}
}
class MyAgent extends Agent {
	@callable({ streaming: true })
	async processWithProgress(stream: StreamingResponse, items: string[]) {
		for (let i = 0; i < items.length; i++) {
			await this.process(items[i]);
			stream.send({ progress: (i + 1) / items.length, item: items[i] });
		}
		stream.end({ completed: true, total: items.length });
	}
}

TypeScript 集成

类型化客户端调用

传入 Agent 类作为类型参数以获得完整类型安全:

import { useAgent } from "agents/react";

function App() {
	const agent = useAgent({
		agent: "MyAgent",
		name: "default",
	});

	async function handleGreet() {
		// TypeScript knows the method signature
		const result = await agent.stub.greet("World");
		// ^? string
	}

	// TypeScript catches errors
	// await agent.stub.greet(123); // Error: Argument of type 'number' is not assignable
	// await agent.stub.nonExistent(); // Error: Property 'nonExistent' does not exist
}
import { useAgent } from "agents/react";
import type { MyAgent } from "./server";

function App() {
	const agent = useAgent<MyAgent>({
		agent: "MyAgent",
		name: "default",
	});

	async function handleGreet() {
		// TypeScript knows the method signature
		const result = await agent.stub.greet("World");
		// ^? string
	}

	// TypeScript catches errors
	// await agent.stub.greet(123); // Error: Argument of type 'number' is not assignable
	// await agent.stub.nonExistent(); // Error: Property 'nonExistent' does not exist
}

排除不可调用方法

若有未用 @callable() 装饰的方法,可从类型中排除:

class MyAgent extends Agent {
	@callable()
	publicMethod() {
		return "public";
	}

	// Not callable from clients
	internalMethod() {
		// internal logic
	}
}

// Exclude internal methods from the client type
const agent = useAgent({
	agent: "MyAgent",
});

agent.stub.publicMethod(); // Works
// agent.stub.internalMethod(); // TypeScript error
class MyAgent extends Agent {
	@callable()
	publicMethod(): string {
		return "public";
	}

	// Not callable from clients
	internalMethod(): void {
		// internal logic
	}
}

// Exclude internal methods from the client type
const agent = useAgent<Omit<MyAgent, "internalMethod">>({
	agent: "MyAgent",
});

agent.stub.publicMethod(); // Works
// agent.stub.internalMethod(); // TypeScript error

错误处理

在可调用方法中抛出错误

可调用方法中抛出的错误会传播到客户端:

class MyAgent extends Agent {
	@callable()
	async riskyOperation(data) {
		if (!isValid(data)) {
			throw new Error("Invalid data format");
		}

		try {
			await this.processData(data);
		} catch (e) {
			throw new Error("Processing failed: " + e.message);
		}
	}
}
class MyAgent extends Agent {
	@callable()
	async riskyOperation(data: unknown): Promise<void> {
		if (!isValid(data)) {
			throw new Error("Invalid data format");
		}

		try {
			await this.processData(data);
		} catch (e) {
			throw new Error("Processing failed: " + e.message);
		}
	}
}

客户端错误处理

try {
	const result = await agent.stub.riskyOperation(data);
} catch (error) {
	// Error thrown by the agent method
	console.error("RPC failed:", error.message);
}
try {
	const result = await agent.stub.riskyOperation(data);
} catch (error) {
	// Error thrown by the agent method
	console.error("RPC failed:", error.message);
}

流式错误处理

流式方法使用 onError 回调:

await agent.call("streamData", [input], {
	stream: {
		onChunk: (chunk) => handleChunk(chunk),
		onError: (errorMessage) => {
			console.error("Stream error:", errorMessage);
			showErrorUI(errorMessage);
		},
		onDone: (result) => handleComplete(result),
	},
});
await agent.call("streamData", [input], {
	stream: {
		onChunk: (chunk) => handleChunk(chunk),
		onError: (errorMessage) => {
			console.error("Stream error:", errorMessage);
			showErrorUI(errorMessage);
		},
		onDone: (result) => handleComplete(result),
	},
});

服务端可用 stream.error() 在流中途优雅发送错误:

class MyAgent extends Agent {
	@callable({ streaming: true })
	async processItems(stream, items) {
		for (const item of items) {
			try {
				const result = await this.process(item);
				stream.send(result);
			} catch (e) {
				stream.error(`Failed to process ${item}: ${e.message}`);
				return; // Stream is now closed
			}
		}
		stream.end();
	}
}
class MyAgent extends Agent {
	@callable({ streaming: true })
	async processItems(stream: StreamingResponse, items: string[]) {
		for (const item of items) {
			try {
				const result = await this.process(item);
				stream.send(result);
			} catch (e) {
				stream.error(`Failed to process ${item}: ${e.message}`);
				return; // Stream is now closed
			}
		}
		stream.end();
	}
}

连接错误

若 WebSocket 连接在 RPC 调用 pending 时关闭,它们会自动以「Connection closed」错误 reject:

try {
	const result = await agent.call("longRunningMethod", []);
} catch (error) {
	if (error.message === "Connection closed") {
		// Handle disconnection
		console.log("Lost connection to agent");
	}
}
try {
	const result = await agent.call("longRunningMethod", []);
} catch (error) {
	if (error.message === "Connection closed") {
		// Handle disconnection
		console.log("Lost connection to agent");
	}
}

重连后重试

客户端断开后会自动重连。要在重连后重试失败的调用,重试前 await agent.ready

async function callWithRetry(agent, method, args = []) {
	try {
		return await agent.call(method, args);
	} catch (error) {
		if (error.message === "Connection closed") {
			await agent.ready; // Wait for reconnection
			return await agent.call(method, args); // Retry once
		}
		throw error;
	}
}

// Usage
const result = await callWithRetry(agent, "processData", [data]);
async function callWithRetry<T>(
	agent: AgentClient,
	method: string,
	args: unknown[] = [],
): Promise<T> {
	try {
		return await agent.call(method, args);
	} catch (error) {
		if (error.message === "Connection closed") {
			await agent.ready; // Wait for reconnection
			return await agent.call(method, args); // Retry once
		}
		throw error;
	}
}

// Usage
const result = await callWithRetry(agent, "processData", [data]);

何时不使用 @callable

Worker 到 Agent 调用

从同一 Worker 调用 Agent(例如在 fetch handler 中)时,直接使用 Durable Object RPC:

import { getAgentByName } from "agents";

export default {
	async fetch(request, env) {
		// Get the agent stub
		const agent = await getAgentByName(env.MyAgent, "instance-name");

		// Call methods directly - no @callable needed
		const result = await agent.processData(data);

		return Response.json(result);
	},
};
import { getAgentByName } from "agents";

export default {
	async fetch(request: Request, env: Env) {
		// Get the agent stub
		const agent = await getAgentByName(env.MyAgent, "instance-name");

		// Call methods directly - no @callable needed
		const result = await agent.processData(data);

		return Response.json(result);
	},
} satisfies ExportedHandler<Env>;

Agent 到 Agent 调用

一个 Agent 需要调用另一个时:

class OrchestratorAgent extends Agent {
	async delegateWork(taskId) {
		// Get another agent
		const worker = await getAgentByName(this.env.WorkerAgent, taskId);

		// Call its methods directly
		const result = await worker.doWork();

		return result;
	}
}
class OrchestratorAgent extends Agent {
	async delegateWork(taskId: string) {
		// Get another agent
		const worker = await getAgentByName(this.env.WorkerAgent, taskId);

		// Call its methods directly
		const result = await worker.doWork();

		return result;
	}
}

为何区分?

RPC 类型 传输 用例
@callable WebSocket 外部客户端(浏览器、应用)
Durable Object RPC 内部 Worker 到 Agent、Agent 到 Agent

Durable Object RPC 对内部调用更高效,无需经过 WebSocket 序列化。@callable 装饰器为外部客户端添加必要的 WebSocket RPC 处理。

API 参考

@callable(metadata?) 装饰器

将方法标记为可供外部客户端调用。

import { callable } from "agents";

class MyAgent extends Agent {
	@callable()
	method() {}

	@callable({ streaming: true })
	streamingMethod(stream) {}

	@callable({ description: "Fetches user data" })
	getUser(id) {}
}
import { callable } from "agents";

class MyAgent extends Agent {
	@callable()
	method(): void {}

	@callable({ streaming: true })
	streamingMethod(stream: StreamingResponse): void {}

	@callable({ description: "Fetches user data" })
	getUser(id: string): User {}
}

CallableMetadata 类型

type CallableMetadata = {
	/** Optional description of what the method does */
	description?: string;
	/** Whether the method supports streaming responses */
	streaming?: boolean;
};

StreamingResponse 类

在流式可调用方法中用于向客户端发送数据。

import {} from "agents";

class MyAgent extends Agent {
	@callable({ streaming: true })
	async streamData(stream, input) {
		stream.send("chunk 1");
		stream.send("chunk 2");
		stream.end("final");
	}
}
import { type StreamingResponse } from "agents";

class MyAgent extends Agent {
	@callable({ streaming: true })
	async streamData(stream: StreamingResponse, input: string) {
		stream.send("chunk 1");
		stream.send("chunk 2");
		stream.end("final");
	}
}
方法 Signature 描述
send (chunk: unknown) => void 向客户端发送 chunk
end (finalChunk?: unknown) => void 结束流
error (message: string) => void 发送错误并关闭流

客户端方法

方法 Signature 描述
agent.call (method, args?, options?) => Promise 按名称调用方法
agent.stub Proxy 类型化方法调用
// Using call()
await agent.call("methodName", [arg1, arg2]);
await agent.call("streamMethod", [arg], {
	stream: { onChunk, onDone, onError },
});

// With timeout (rejects if call does not complete in time)
await agent.call("slowMethod", [], { timeout: 5000 });

// Using stub
await agent.stub.methodName(arg1, arg2);
// Using call()
await agent.call("methodName", [arg1, arg2]);
await agent.call("streamMethod", [arg], {
	stream: { onChunk, onDone, onError },
});

// With timeout (rejects if call does not complete in time)
await agent.call("slowMethod", [], { timeout: 5000 });

// Using stub
await agent.stub.methodName(arg1, arg2);

CallOptions 类型

type CallOptions = {
	/** Timeout in milliseconds. Rejects if call does not complete in time. */
	timeout?: number;
	/** Streaming options */
	stream?: {
		onChunk?: (chunk: unknown) => void;
		onDone?: (finalChunk: unknown) => void;
		onError?: (error: string) => void;
	};
};

getCallableMethods() 方法

返回 Agent 上所有可调用方法及其元数据的映射。适用于 introspection 与自动文档。

const methods = agent.getCallableMethods();
// Map<string, CallableMetadata>

for (const [name, meta] of methods) {
	console.log(`${name}: ${meta.description || "(no description)"}`);
	if (meta.streaming) console.log("  (streaming)");
}
const methods = agent.getCallableMethods();
// Map<string, CallableMetadata>

for (const [name, meta] of methods) {
	console.log(`${name}: ${meta.description || "(no description)"}`);
	if (meta.streaming) console.log("  (streaming)");
}

故障排除

SyntaxError: Invalid or unexpected token

使用 @callable() 时若 dev server 因 SyntaxError: Invalid or unexpected token 失败,需要两件事:

1. 添加 agents/vite 插件 — Vite 8 使用 Oxc 转译,尚不支持 TC39 装饰器。该插件添加所需 transform:

vite.config.tsts
import agents from "agents/vite";

export default defineConfig({
	plugins: [agents(), react(), cloudflare()],
});

2. 扩展 agents/tsconfig — 设置 "target": "ES2021" 及所有其他推荐编译器选项:

tsconfig.jsonjson
{
	"extends": "agents/tsconfig"
}

若无法扩展共享配置,请在 tsconfig.json 中手动设置 "target": "ES2021"

后续步骤

这篇文档对您有帮助吗?