diff --git a/docs/feature/agent-tools.md b/docs/feature/agent-tools.md
new file mode 100644
index 0000000000..e04e7ec4ac
--- /dev/null
+++ b/docs/feature/agent-tools.md
@@ -0,0 +1,110 @@
+# Agent Tools & Event Triggers
+
+For users who want an EventMesh-hosted agent to **act on** the mesh — deliver
+LLM decisions through connectors, poll external systems mid-conversation, or
+run agents fully event-driven without a human in the loop. Covers the
+`eventmesh-agent` extension points (LLM client, conversation memory, tools,
+event triggers) and how they map onto the connector plugin ecosystem.
+*(Experimental)*
+
+---
+
+## The extension surface
+
+`StreamingAgent` is assembled from constructor-injected interfaces — every
+piece below has a shipped default and can be replaced without forking:
+
+| Extension point | Interface | Default | Replace it to… |
+| --- | --- | --- | --- |
+| Chat backend | `agent.llm.LlmClient` | `OpenAiLlmClient` (OpenAI-compatible SSE) | call Anthropic native, Bedrock, or an internal inference service |
+| Conversation history | `agent.ConversationMemory` | `ConversationStore` (in-memory window) | persist sessions to Redis / RocksDB / a database |
+| Callable tools | `agent.tool.AgentTool` | none registered (plain streaming) | expose any capability to the model via function calling |
+| Event triggers | `StreamingAgent.onEvent(...)` | off | drive the agent from topic subscriptions instead of user prompts |
+
+## Connector plugins as agent tools
+
+`ConnectorToolAdapter` turns the shipped connector SPI into `AgentTool`s, so
+every connector plugin doubles as a tool the LLM may call:
+
+- **Sink tool (write)** — the model's arguments object is wrapped as one
+ CloudEvent and pushed through `SinkConnector.put()`. Example: a DingTalk
+ sink connector becomes a `notify` tool the model can call to alert a
+ channel.
+- **Source tool (read)** — each call polls one batch via
+ `SourceConnector.poll()` and returns it as a JSON array. Example: a canal
+ source connector becomes a `query-changes` tool for mid-conversation data
+ lookups.
+
+Wire-up in `AgentApplication` (system properties):
+
+```properties
+# a write tool backed by any sink connector on the classpath
+agent.tools.sink.notify=org.apache.eventmesh.connector.dingtalk.sink.DingTalkSinkConnector
+agent.tools.props.notify.connector.sink.webhook=https://oapi.dingtalk.com/robot/send?...
+
+# a read tool backed by any source connector
+agent.tools.source.query-changes=org.apache.eventmesh.connector.canal.source.CanalSourceConnector
+agent.tools.props.query-changes.connector.batchSize=10
+```
+
+With tools registered the agent switches from plain token streaming to a
+**function-calling loop**: the model may request tools, receives their
+results as messages, and continues until it produces a final answer (bounded
+at 5 tool iterations). Without tools, behavior is unchanged
+token-by-token streaming.
+
+## Event-driven agents (no user in the loop)
+
+Set a subscription list and an output topic:
+
+```properties
+agent.subscribe.topics=orders.changed,risk.alerts
+agent.trigger.output.topic=agent.decisions
+```
+
+Each consumed CloudEvent becomes a prompt, is answered (with tools if
+registered), and the answer is published onto the output topic — where any
+**sink connector** can deliver it (DingTalk, Kafka, Spring, ...). Trigger
+conversations are keyed `trigger:` and never bleed into user
+sessions. This is the agent-side twin of the connector pipeline:
+
+```
+source connector ──> topic ──> agent (trigger) ──> LLM (+ tools)
+ │
+ └── output topic ──> sink connector ──> external system
+```
+
+## Using the interfaces directly (embedder API)
+
+```java
+LlmClient llm = new OpenAiLlmClient(baseUrl, apiKey, model); // or your own impl
+ConversationMemory memory = new ConversationStore(20); // or a persistent impl
+ToolRegistry tools = new ToolRegistry()
+ .register(ConnectorToolAdapter.sinkTool("notify", "Send a DingTalk notification",
+ DingTalkSinkConnector.class, props));
+StreamingAgent agent = new StreamingAgent(client, parent, agentId, llm, memory, tools);
+```
+
+## Configuration reference
+
+| Key | Default | Meaning |
+| --- | --- | --- |
+| `agent.tools.sink.` | — | FQCN of a `SinkConnector` exposed as write tool `` |
+| `agent.tools.source.` | — | FQCN of a `SourceConnector` exposed as read tool `` |
+| `agent.tools.props..*` | — | Connector init properties for tool `` |
+| `agent.subscribe.topics` | (empty) | Comma-separated topics consumed as event triggers |
+| `agent.trigger.output.topic` | `agent.triggers` | Topic the trigger answers are published to |
+
+The agent connects to the runtime like any SDK client (`agent.runtime.url`,
+default `http://localhost:10105`).
+
+## Where the code lives
+
+| Piece | Location |
+| --- | --- |
+| LLM SPI + OpenAI default | `eventmesh-agent/.../agent/llm/{LlmClient,OpenAiLlmClient}.java` |
+| Memory SPI + in-memory default | `eventmesh-agent/.../agent/{ConversationMemory,ConversationStore}.java` |
+| Tool SPI + registry + adapter | `eventmesh-agent/.../agent/tool/{AgentTool,ToolRegistry,ConnectorToolAdapter}.java` |
+| Tool loop + event trigger path | `eventmesh-agent/.../agent/StreamingAgent.java` |
+| Boot wiring (tools + triggers) | `eventmesh-agent/.../agent/AgentApplication.java` |
+| Tests | `eventmesh-agent/src/test/.../tool/*`, `.../llm/OpenAiLlmClientChatTest.java` |
diff --git a/docs/index.md b/docs/index.md
index 8e40738099..fcb63b7639 100644
--- a/docs/index.md
+++ b/docs/index.md
@@ -49,6 +49,10 @@ too.
- [A2A — Agent-to-Agent](feature/a2a.md) — the agent collaboration protocol:
task lifecycle, JSON-RPC/MCP bridge, REST + SDK usage
*(Experimental)*.
+- [Agent tools & event triggers](feature/agent-tools.md) — the agent
+ extension points (LLM client, conversation memory, function-calling tools)
+ and connector plugins as the agent tool library; event-driven agents
+ *(Experimental)*.
- [Connectors](feature/connectors/README.md) — bridge EventMesh with external
systems (Kafka, RocketMQ, RabbitMQ, databases, chat platforms): the connector
diff --git a/eventmesh-agent/build.gradle b/eventmesh-agent/build.gradle
index 17858d5e82..08a6703fea 100644
--- a/eventmesh-agent/build.gradle
+++ b/eventmesh-agent/build.gradle
@@ -1,6 +1,6 @@
/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
@@ -16,10 +16,13 @@
*/
// Agent — 独立进程, 仅经 lite topic + CloudEvents 与 EventMesh Runtime 通信.
-// 不依赖 eventmesh-runtime, 复用 sdk-java (CloudEventsClient) + common (stream framing) + 一个
+// 不依赖 eventmesh-runtime, 复用 sdk-java (CloudEventsClient) + common (stream framing) +
+// connector-api (ConnectorToolAdapter 把存量 connector 插件变成 LLM 工具) + 一个
// OpenAI 兼容的 LLM SSE client.
dependencies {
implementation project(':eventmesh-sdks:eventmesh-sdk-java')
+ implementation project(':eventmesh-connector-api')
+ implementation 'io.cloudevents:cloudevents-core'
implementation 'com.fasterxml.jackson.core:jackson-databind'
implementation 'org.slf4j:slf4j-api'
runtimeOnly 'org.apache.logging.log4j:log4j-slf4j2-impl'
diff --git a/eventmesh-agent/conf/agent.properties b/eventmesh-agent/conf/agent.properties
index 78e52adfc1..9314b061c6 100644
--- a/eventmesh-agent/conf/agent.properties
+++ b/eventmesh-agent/conf/agent.properties
@@ -28,6 +28,15 @@
# AGENT_HEARTBEAT_FAILLIMIT -> -Dagent.heartbeat.failLimit (default 6; exits after N consecutive failures)
# AGENT_MAX_CONVERSATIONS -> -Dagent.conversation.maxConversations (default 1000; LRU bound)
#
+# Connector tools (function calling; see docs/feature/agent-tools.md):
+# AGENT_TOOL_SINK_= sink connector exposed as a write tool
+# AGENT_TOOL_SOURCE_= source connector exposed a read tool
+# pass connector init props via AGENT_OPTS: -Dagent.tools.props..=
+#
+# Event-driven triggers:
+# AGENT_SUBSCRIBE_TOPICS -> -Dagent.subscribe.topics (comma list; default none)
+# AGENT_TRIGGER_OUTPUT -> -Dagent.trigger.output.topic (default agent.triggers)
+#
# Note: an empty LLM_API_KEY fails fast unless LLM_API_KEY_OPTIONAL=true is set via AGENT_OPTS.
#
# Example (real LLM):
diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java
index b686737a8a..ebbb9b54d4 100644
--- a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java
+++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java
@@ -18,9 +18,15 @@
package org.apache.eventmesh.agent;
import org.apache.eventmesh.agent.llm.OpenAiLlmClient;
+import org.apache.eventmesh.agent.tool.AgentTool;
+import org.apache.eventmesh.agent.tool.ConnectorToolAdapter;
+import org.apache.eventmesh.agent.tool.ToolRegistry;
import org.apache.eventmesh.client.cloudevents.CloudEventsClient;
+import org.apache.eventmesh.connector.SinkConnector;
+import org.apache.eventmesh.connector.SourceConnector;
import java.util.List;
+import java.util.Properties;
import lombok.extern.slf4j.Slf4j;
@@ -31,7 +37,10 @@
*
* Config keys: {@code agent.runtime.url}, {@code agent.id}, {@code agent.heartbeat.intervalMs},
* {@code agent.capacity}, {@code agent.conversation.maxHistory}, {@code llm.base.url},
- * {@code llm.api.key}, {@code llm.model}.
+ * {@code llm.api.key}, {@code llm.model}; connector tools via
+ * {@code agent.tools.sink.=} / {@code agent.tools.source.=}
+ * (+ {@code agent.tools.props..*}); event triggers via {@code agent.subscribe.topics} +
+ * {@code agent.trigger.output.topic}.
*/
@Slf4j
public class AgentApplication {
@@ -57,6 +66,8 @@ public static void main(String[] args) throws Exception {
// supervisor can restart it; the runtime TTL has evicted the registration by then
// and a silent zombie serves nothing.
final int heartbeatFailLimit = Integer.getInteger("agent.heartbeat.failLimit", 6);
+ String triggerTopics = System.getProperty("agent.subscribe.topics", "");
+ String triggerOutput = System.getProperty("agent.trigger.output.topic", "agent.triggers");
// Step 1: register (gets the assigned agent-parent + client-reply-parent)
AgentControlClient control = new AgentControlClient(runtimeUrl);
@@ -64,13 +75,29 @@ public static void main(String[] args) throws Exception {
String agentParent = reg.parent();
log.info("agent registered: agentId={} parent={} clientParent={}", agentId, agentParent, reg.clientParent());
- // Step 2: subscribe the agent's channel
+ // Step 2: subscribe the agent's channel (with connector tools, if configured)
CloudEventsClient client = CloudEventsClient.builder().runtimeUrl(runtimeUrl).clientId(agentId).build();
OpenAiLlmClient llm = new OpenAiLlmClient(llmBase, llmKey, llmModel);
ConversationStore store = new ConversationStore(maxHistory, maxConversations);
- StreamingAgent agent = new StreamingAgent(client, agentParent, agentId, llm, store);
+ ToolRegistry tools = buildConnectorTools();
+ StreamingAgent agent = new StreamingAgent(client, agentParent, agentId, llm, store, tools);
agent.start(); // subscribe (agentParent, agent.)
+ // Step 2b: optional event-driven triggers (own client so the lite-channel poller is untouched)
+ CloudEventsClient triggerClient = null;
+ if (!triggerTopics.isBlank()) {
+ triggerClient = CloudEventsClient.builder().runtimeUrl(runtimeUrl)
+ .clientId(agentId + "-triggers").build();
+ for (String topic : triggerTopics.split(",")) {
+ String trimmed = topic.trim();
+ if (trimmed.isEmpty()) {
+ continue;
+ }
+ triggerClient.subscribe(trimmed, "LOAD_BALANCE", event -> agent.onEvent(event, triggerOutput));
+ }
+ log.info("agent trigger subscriptions: [{}] -> output topic {}", triggerTopics, triggerOutput);
+ }
+
// Step 3: ready-before-route: only now is this agent eligible for matchmaking
control.ready(agentId);
log.info("agent READY: agentId={} (subscribed agent.{})", agentId, agentId);
@@ -99,6 +126,7 @@ public static void main(String[] args) throws Exception {
}
});
+ final CloudEventsClient finalTriggerClient = triggerClient;
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
log.info("shutting down agent...");
heartbeatThread.interrupt();
@@ -108,11 +136,61 @@ public static void main(String[] args) throws Exception {
log.warn("unregister failed: {}", e.toString());
}
agent.shutdown();
+ if (finalTriggerClient != null) {
+ finalTriggerClient.shutdown();
+ }
client.shutdown();
}, "agent-shutdown"));
- log.info("AgentApplication running: runtime={} agentId={} parent={} model={} (Ctrl+C to stop)",
- runtimeUrl, agentId, agentParent, llmModel);
+ log.info("AgentApplication running: runtime={} agentId={} parent={} model={} tools={} (Ctrl+C to stop)",
+ runtimeUrl, agentId, agentParent, llmModel, tools.size());
Thread.currentThread().join();
}
+
+ /**
+ * Build the connector-backed tool registry from {@code -Dagent.tools.sink.=} /
+ * {@code -Dagent.tools.source.=} plus per-tool connector properties under
+ * {@code agent.tools.props..*}. Empty registry when nothing is configured.
+ */
+ private static ToolRegistry buildConnectorTools() {
+ ToolRegistry registry = new ToolRegistry();
+ for (String key : System.getProperties().stringPropertyNames()) {
+ if (key.startsWith("agent.tools.sink.")) {
+ String name = key.substring("agent.tools.sink.".length());
+ registerConnectorTool(registry, "sink", name, System.getProperty(key));
+ } else if (key.startsWith("agent.tools.source.")) {
+ String name = key.substring("agent.tools.source.".length());
+ registerConnectorTool(registry, "source", name, System.getProperty(key));
+ }
+ }
+ return registry;
+ }
+
+ private static void registerConnectorTool(ToolRegistry registry, String kind, String name, String fqcn) {
+ Properties props = new Properties();
+ String prefix = "agent.tools.props." + name + ".";
+ for (String key : System.getProperties().stringPropertyNames()) {
+ if (key.startsWith(prefix)) {
+ props.put(key.substring(prefix.length()), System.getProperty(key));
+ }
+ }
+ AgentTool tool;
+ if ("sink".equals(kind)) {
+ tool = ConnectorToolAdapter.sinkTool(name, "Deliver a JSON payload through the '" + name + "' connector",
+ loadConnector(fqcn, SinkConnector.class), props);
+ } else {
+ tool = ConnectorToolAdapter.sourceTool(name, "Poll the next event batch from the '" + name
+ + "' connector", loadConnector(fqcn, SourceConnector.class), props);
+ }
+ registry.register(tool);
+ log.info("connector tool registered: {}={} ({})", kind, name, fqcn);
+ }
+
+ private static Class extends T> loadConnector(String fqcn, Class spi) {
+ try {
+ return Class.forName(fqcn).asSubclass(spi);
+ } catch (ClassNotFoundException e) {
+ throw new IllegalArgumentException("connector class not found on agent classpath: " + fqcn, e);
+ }
+ }
}
diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationMemory.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationMemory.java
new file mode 100644
index 0000000000..44872710e5
--- /dev/null
+++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationMemory.java
@@ -0,0 +1,35 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.agent;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Conversation history abstraction keyed by {@code conversationId}. {@link ConversationStore} is
+ * the in-memory default; implement this interface to persist history (Redis, RocksDB, a database,
+ * ...) — hand the instance to {@code StreamingAgent}.
+ */
+public interface ConversationMemory {
+
+ /** Snapshot of the conversation history (empty list if id null/unknown). Caller may mutate. */
+ List