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
24 changes: 22 additions & 2 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -248,7 +248,7 @@ tasks.register('dist-connector') {
// exclude 'eventmesh-*' below, but ConnectorApplication itself imports the
// SPI interfaces — ship it next to the runtime jar in apps/ so a bare
// runtime (no plugins installed) still boots.
def connectorApi = findProject('eventmesh-connector-api')
def connectorApi = findProject('eventmesh-connector-plugin:eventmesh-connector-api')
copy {
from connectorApi.jar.archivePath
into rootProject.file('dist-connector/apps')
Expand Down Expand Up @@ -304,7 +304,7 @@ tasks.register('dist-agent') {
doLast {
// Core: the agent process jar + its third-party deps (sdk-java, common, cloudevents,
// jackson, log4j binding). The agent talks to the runtime only over HTTP + lite topics.
def core = findProject('eventmesh-agent')
def core = findProject('eventmesh-agent-runtime')
logger.lifecycle('Install module: module: {}', core.name)
copy {
from core.jar.archivePath
Expand All @@ -326,6 +326,26 @@ tasks.register('dist-agent') {
duplicatesStrategy = DuplicatesStrategy.EXCLUDE
exclude 'META-INF'
}
// Agent tool plugins (pluginType == agentTool) into dist-agent/plugin/agent/<name>/
String[] agentLibJars = java.util.Optional.ofNullable(file('dist-agent/lib').list()).orElse(new String[0])
def agentPlugins = subprojects.findAll {
it.file('gradle.properties').exists()
&& it.properties.containsKey('pluginType')
&& it.properties.get('pluginType') == 'agentTool'
}
agentPlugins.forEach(subProject -> {
var pluginName = subProject.properties.get('pluginName')
logger.lifecycle('Install agent plugin: {}, module: {}', pluginName, subProject.name)
copy {
from subProject.jar.archivePath
into rootProject.file("dist-agent/plugin/agent/${pluginName}")
}
copy {
from subProject.configurations.runtimeClasspath
into rootProject.file("dist-agent/plugin/agent/${pluginName}")
exclude(agentLibJars)
}
})
copy {
from 'tools/dist-license'
into rootProject.file('dist-agent')
Expand Down
47 changes: 41 additions & 6 deletions docs/feature/agent-tools.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,40 @@ 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.

## In-tree agent plugins (`eventmesh-agent-plugin/`)

Beyond connector-backed tools, custom tools follow the repo-standard plugin
mechanism (same as storage/connector plugins):

1. Implement `agent.tool.AgentTool` (name + description + JSON schema +
`invoke`), no extra annotation needed on your class.
2. In your jar add a service file
`META-INF/eventmesh/org.apache.eventmesh.agent.tool.AgentTool` containing
`mytool=com.example.MyTool`.
3. Drop the jar into the agent's `plugin/agent/` directory (the launcher
puts every jar there on the classpath) and enable it with
`AGENT_TOOLS_SPI=mytool` (i.e. `-Dagent.tools.spi=mytool`; comma-list
supported).

Resolution goes through `EventMeshExtensionFactory` — the same loader that
serves storage and connector plugins, so singleton semantics and the
`META-INF/eventmesh/` convention are identical.

## In-tree agent plugins

The repo ships a plugin tree for agent tools — `eventmesh-agent-plugin/` — mirroring
`eventmesh-connector-plugin/`: each sub-module is a standalone jar with the
`META-INF/eventmesh/org.apache.eventmesh.agent.tool.AgentTool` service file, declared in
`gradle.properties` with `pluginType=agentTool` + `pluginName=<name>`. The `dist-agent` task
installs every agent plugin into `dist-agent/plugin/agent/<name>/` (the directory the agent
launcher puts on the classpath), so custom in-tree tools need zero wiring:

- `eventmesh-agent-plugin-http-fetch` — reference implementation: fetch a URL and return the
(truncated) body to the model. Enable with `AGENT_TOOLS_SPI=http-fetch`.

Third-party jars follow the same shape: implement `AgentTool`, ship the service file, drop the
jar into `plugin/agent/`.

## Event-driven agents (no user in the loop)

Set a subscription list and an output topic:
Expand Down Expand Up @@ -89,6 +123,7 @@ StreamingAgent agent = new StreamingAgent(client, parent, agentId, llm, memory,

| Key | Default | Meaning |
| --- | --- | --- |
| `agent.tools.spi` | (empty) | Comma list of SPI names resolved via `EventMeshExtensionFactory` (jars in `plugin/agent/`) |
| `agent.tools.sink.<name>` | — | FQCN of a `SinkConnector` exposed as write tool `<name>` |
| `agent.tools.source.<name>` | — | FQCN of a `SourceConnector` exposed as read tool `<name>` |
| `agent.tools.props.<name>.*` | — | Connector init properties for tool `<name>` |
Expand All @@ -102,9 +137,9 @@ default `http://localhost:10105`).

| 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` |
| LLM SPI + OpenAI default | `eventmesh-agent-runtime/.../agent/llm/{LlmClient,OpenAiLlmClient}.java` |
| Memory SPI + in-memory default | `eventmesh-agent-runtime/.../agent/{ConversationMemory,ConversationStore}.java` |
| Tool SPI + registry + adapter | `eventmesh-agent-runtime/.../agent/tool/{AgentTool,ToolRegistry,ConnectorToolAdapter}.java` |
| Tool loop + event trigger path | `eventmesh-agent-runtime/.../agent/StreamingAgent.java` |
| Boot wiring (tools + triggers) | `eventmesh-agent-runtime/.../agent/AgentApplication.java` |
| Tests | `eventmesh-agent-runtime/src/test/.../tool/*`, `.../llm/OpenAiLlmClientChatTest.java` |
4 changes: 2 additions & 2 deletions docs/feature/client-java.md
Original file line number Diff line number Diff line change
Expand Up @@ -361,7 +361,7 @@ An agent that participates in Mode 1 follows a four-step contract:
3. on normal completion → emit a terminal frame `{chunk: "", done: true}`
4. on error → emit a terminal error frame `{chunk: "", done: true, error: "..."}`

Reference implementation: `eventmesh-agent/.../StreamingAgent.java`
Reference implementation: `eventmesh-agent-runtime/.../StreamingAgent.java`
(instantiate with an LLM client, an `agentParent` topic, the agent's `agentId`,
and a `ConversationStore`).

Expand Down Expand Up @@ -656,7 +656,7 @@ history for the migration notes.
(`/events/*` endpoints; `withSecurityGate(...)` wiring point)
* A2A HTTP handler — `eventmesh-runtime/.../a2a/A2AGatewayHttpHandler.java`
(`/a2a/*` endpoints; `withSecurityGate(...)` wiring point)
* Streaming agent — `eventmesh-agent/.../StreamingAgent.java`
* Streaming agent — `eventmesh-agent-runtime/.../StreamingAgent.java`
* Storage plugins —
* `eventmesh-storage-plugin/eventmesh-storage-rocketmq/` (SPI key `rocketmq`)
* `eventmesh-storage-plugin/eventmesh-storage-rocketmq5/` (SPI key `rocketmq5`,
Expand Down
16 changes: 8 additions & 8 deletions docs/feature/connector-api-split.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
## Goal

Move the connector SPI interfaces out of `eventmesh-connector-runtime` into a new
`eventmesh-connector-api` module so plugins depend on a stable, minimal API jar and the runtime
`eventmesh-connector-plugin:eventmesh-connector-plugin:eventmesh-connector-api` module so plugins depend on a stable, minimal API jar and the runtime
implements / orchestrates against that API.

## Interfaces to move
Expand All @@ -24,7 +24,7 @@ From `eventmesh-connector-runtime/src/main/java/org/apache/eventmesh/connector/`
## Target layout

```
eventmesh-connector-api/
eventmesh-connector-plugin/eventmesh-connector-api/
src/main/java/org/apache/eventmesh/connector/api/
SourceConnector.java
SinkConnector.java
Expand All @@ -37,20 +37,20 @@ eventmesh-connector-api/
build.gradle (deps: cloudevents-core only — no plugin imports, no HTTP libs)
```

`eventmesh-connector-runtime` depends on `:eventmesh-connector-api` and continues to provide
`eventmesh-connector-runtime` depends on `:eventmesh-connector-plugin:eventmesh-connector-api` and continues to provide
implementations (`EventMeshHttpEndpoint`, `RocksDBConnectorOffsetStore`, `RemoteOffsetStore`,
`InMemoryOffsetStore`, `ConnectorRuntime`, `ConnectorManager`, `ConnectorAdminServer`,
`ConnectorApplication`, `ConnectorDef`).

Each plugin under `eventmesh-connector-plugin/eventmesh-connector-*` should depend on
`:eventmesh-connector-api` instead of `:eventmesh-connector-runtime`.
`:eventmesh-connector-plugin:eventmesh-connector-api` instead of `:eventmesh-connector-runtime`.

## Plugin changes (mechanical)

For each of the 23 plugins:

1. `build.gradle`: replace `implementation project(":eventmesh-connector-runtime")` with
`implementation project(":eventmesh-connector-api")`.
`implementation project(":eventmesh-connector-plugin:eventmesh-connector-api")`.
2. Source code: if the plugin imports `org.apache.eventmesh.connector.ConnectorRuntime` (it should
not — plugins only use the SPI), add `implementation project(":eventmesh-connector-runtime")`
back. Initial audit shows no plugin currently touches runtime internals.
Expand Down Expand Up @@ -90,9 +90,9 @@ will fail the architecture guard.

## Acceptance criteria for the implementation PR(s)

- [ ] `eventmesh-connector-api` jar builds standalone (deps: cloudevents-core only).
- [ ] `eventmesh-connector-runtime` depends on `:eventmesh-connector-api`.
- [ ] All 23 plugin modules depend on `:eventmesh-connector-api`, not on `:eventmesh-connector-runtime`.
- [ ] `eventmesh-connector-plugin:eventmesh-connector-plugin:eventmesh-connector-api` jar builds standalone (deps: cloudevents-core only).
- [ ] `eventmesh-connector-runtime` depends on `:eventmesh-connector-plugin:eventmesh-connector-api`.
- [ ] All 23 plugin modules depend on `:eventmesh-connector-plugin:eventmesh-connector-api`, not on `:eventmesh-connector-runtime`.
- [ ] ArchUnit rule is added and **fails** the build if any plugin reaches into runtime internals.
- [ ] `:eventmesh-architecture-guard:test` passes.
- [ ] Existing runtime + plugin tests stay green (this PR added the baseline tests they will
Expand Down
20 changes: 20 additions & 0 deletions eventmesh-agent-plugin/build.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
/*
* 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.
*/

// Umbrella for agent tool plugins (AgentTool SPI, #5409). Each sub-module ships a jar with a
// META-INF/eventmesh/<AgentTool-FQCN> service file; the dist-agent task installs them into
// dist-agent/plugin/agent/<name>/.
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
/*
* 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.
*/

dependencies {
implementation project(':eventmesh-agent-runtime')
implementation 'com.fasterxml.jackson.core:jackson-databind'

testImplementation 'org.assertj:assertj-core'
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
# 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.
#

pluginType=agentTool
pluginName=http-fetch
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
/*
* 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.plugin.httpfetch;

import org.apache.eventmesh.agent.tool.AgentTool;

import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.time.Duration;
import java.util.Map;

import com.fasterxml.jackson.databind.ObjectMapper;

/**
* Reference {@link AgentTool} implementation shipped as the first in-tree agent plugin: fetches a
* URL and returns the body (truncated) as the tool result for the model. Deployed via the
* META-INF/eventmesh service file; enable with {@code -Dagent.tools.spi=http-fetch}.
*/
public class HttpFetchTool implements AgentTool {

private static final ObjectMapper MAPPER = new ObjectMapper();

private static final int MAX_BODY_CHARS = 4000;

private final HttpClient http = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(10))
.followRedirects(HttpClient.Redirect.NORMAL)
.build();

@Override
public String name() {
return "http-fetch";
}

@Override
public String description() {
return "Fetch an HTTP/HTTPS URL and return the response body (text, truncated)";
}

@Override
public String parametersJsonSchema() {
return MAPPER.valueToTree(Map.of(
"type", "object",
"properties", Map.of(
"url", Map.of("type", "string", "description", "absolute http(s) URL to fetch")),
"required", java.util.List.of("url"))).toString();
}

@Override
public String invoke(Map<String, Object> args) throws Exception {
String url = String.valueOf(args.get("url"));
HttpRequest req = HttpRequest.newBuilder()
.uri(URI.create(url))
.timeout(Duration.ofSeconds(20))
.GET()
.build();
HttpResponse<String> resp = http.send(req, HttpResponse.BodyHandlers.ofString());
String body = resp.body() == null ? "" : resp.body();
String truncated = body.length() > MAX_BODY_CHARS ? body.substring(0, MAX_BODY_CHARS) + "..." : body;
return MAPPER.writeValueAsString(Map.of("status", resp.statusCode(), "body", truncated));
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
# 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.
#
http-fetch=org.apache.eventmesh.agent.plugin.httpfetch.HttpFetchTool
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,17 @@ AGENT_MAX_CONVERSATIONS="${AGENT_MAX_CONVERSATIONS:-1000}"
AGENT_CAPACITY="${AGENT_CAPACITY:-100}"
AGENT_OPTS="${AGENT_OPTS:-}"

# SPI-deployed agent tools: every jar under plugin/agent/ joins the classpath (an AgentTool
# implementation registers itself via META-INF/eventmesh/<AgentTool-FQCN> inside the jar; enable
# with -Dagent.tools.spi=<name>).
PLUGIN_CP=""
if compgen -G "$AGENT_HOME/plugin/agent/*.jar" > /dev/null; then
PLUGIN_CP="$(find "$AGENT_HOME/plugin/agent" -name '*.jar' | paste -sd ':' -)"
fi

ARGS=""
if [ -n "$AGENT_ID" ]; then
ARGS="$ARGS -Dagent.id=${AGENT_ID}"
fi

exec java -Xmx512m $AGENT_OPTS -cp "conf:apps/*:lib/*" -Dagent.runtime.url="${AGENT_RUNTIME_URL}" -Dllm.base.url="${LLM_BASE_URL}" -Dllm.api.key="${LLM_API_KEY}" -Dllm.model="${LLM_MODEL}" -Dagent.heartbeat.intervalMs="${AGENT_HEARTBEAT_MS}" -Dagent.heartbeat.failLimit="${AGENT_HEARTBEAT_FAILLIMIT}" -Dagent.conversation.maxHistory="${AGENT_MAX_HISTORY}" -Dagent.conversation.maxConversations="${AGENT_MAX_CONVERSATIONS}" -Dagent.capacity="${AGENT_CAPACITY}" $ARGS org.apache.eventmesh.agent.AgentApplication
exec java -Xmx512m $AGENT_OPTS -cp "conf:apps/*:lib/*${PLUGIN_CP:+:$PLUGIN_CP}" -Dagent.runtime.url="${AGENT_RUNTIME_URL}" -Dllm.base.url="${LLM_BASE_URL}" -Dllm.api.key="${LLM_API_KEY}" -Dllm.model="${LLM_MODEL}" -Dagent.heartbeat.intervalMs="${AGENT_HEARTBEAT_MS}" -Dagent.heartbeat.failLimit="${AGENT_HEARTBEAT_FAILLIMIT}" -Dagent.conversation.maxHistory="${AGENT_MAX_HISTORY}" -Dagent.conversation.maxConversations="${AGENT_MAX_CONVERSATIONS}" -Dagent.capacity="${AGENT_CAPACITY}" $ARGS org.apache.eventmesh.agent.AgentApplication
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,8 @@
// OpenAI 兼容的 LLM SSE client.
dependencies {
implementation project(':eventmesh-sdks:eventmesh-sdk-java')
implementation project(':eventmesh-connector-api')
implementation project(':eventmesh-connector-plugin:eventmesh-connector-api')
implementation project(':eventmesh-spi')
implementation 'io.cloudevents:cloudevents-core'
implementation 'com.fasterxml.jackson.core:jackson-databind'
implementation 'org.slf4j:slf4j-api'
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@
# AGENT_HEARTBEAT_FAILLIMIT -> -Dagent.heartbeat.failLimit (default 6; exits after N consecutive failures)
# AGENT_MAX_CONVERSATIONS -> -Dagent.conversation.maxConversations (default 1000; LRU bound)
#
# SPI-deployed custom tools (jars under plugin/agent/, META-INF/eventmesh service file):
# AGENT_TOOLS_SPI -> -Dagent.tools.spi (comma list of SPI names; default none)
#
# Connector tools (function calling; see docs/feature/agent-tools.md):
# AGENT_TOOL_SINK_<NAME>=<fqcn> sink connector exposed as a write tool
# AGENT_TOOL_SOURCE_<NAME>=<fqcn> source connector exposed a read tool
Expand Down
Loading
Loading