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 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> get(String conversationId); + + /** Append a completed turn (user prompt + assistant answer). No-op if id is null. */ + void appendTurn(String conversationId, String userPrompt, String assistantAnswer); +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java index 44b1dfb76f..c1b4d7bf2f 100644 --- a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java @@ -32,10 +32,11 @@ * number of conversations themselves ({@code maxConversations}, access-order eviction) — * without the outer bound a long-lived agent leaks memory one session at a time (#5405). * - *

Lost on agent restart (MVP — persistence to Redis/DB is a TODO). Thread-safe per conversation + *

Lost on agent restart (persistence to Redis/DB is a TODO — implement {@link ConversationMemory} + * and inject it into {@code StreamingAgent} to add one). Thread-safe per conversation * (synchronized on the per-id list).

*/ -public class ConversationStore { +public class ConversationStore implements ConversationMemory { private final int maxMessages; private final ConcurrentHashMap>> history = new ConcurrentHashMap<>(); @@ -55,6 +56,7 @@ public ConversationStore(int maxMessages, int maxConversations) { private final int maxConversations; /** Snapshot of the conversation history (empty list if convId null/unknown). Caller may mutate. */ + @Override public List> get(String conversationId) { if (conversationId == null) { return new ArrayList<>(); @@ -94,6 +96,7 @@ private synchronized void touch(String id, boolean admit) { } /** Append a completed turn (user prompt + assistant answer). No-op if convId is null. */ + @Override public void appendTurn(String conversationId, String userPrompt, String assistantAnswer) { if (conversationId == null) { return; diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/StreamingAgent.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/StreamingAgent.java index acac3dac99..c13c9d679a 100644 --- a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/StreamingAgent.java +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/StreamingAgent.java @@ -17,18 +17,28 @@ package org.apache.eventmesh.agent; -import org.apache.eventmesh.agent.llm.OpenAiLlmClient; +import org.apache.eventmesh.agent.llm.LlmClient; +import org.apache.eventmesh.agent.llm.LlmCompletion; +import org.apache.eventmesh.agent.llm.ToolCall; +import org.apache.eventmesh.agent.tool.ToolRegistry; import org.apache.eventmesh.client.cloudevents.CloudEventsClient; import org.apache.eventmesh.common.stream.StreamChunk; import org.apache.eventmesh.common.stream.StreamRequest; +import java.nio.charset.StandardCharsets; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.UUID; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicInteger; +import io.cloudevents.CloudEvent; +import io.cloudevents.core.builder.CloudEventBuilder; + +import com.fasterxml.jackson.databind.ObjectMapper; + import lombok.extern.slf4j.Slf4j; /** @@ -37,29 +47,51 @@ * lite). Each request runs on its own virtual thread. Multi-turn history is keyed by * {@code sessionId} (sessionId IS the conversation key). * + *

Extension points (constructor-injected): {@link LlmClient} (any chat backend — the default is + * the OpenAI-compatible client), {@link ConversationMemory} (history persistence), and an optional + * {@link ToolRegistry} of {@code AgentTool}s. With tools registered the agent runs a + * function-calling loop: the model may request tools (e.g. connector-backed write/read tools) and + * receives their results before producing the final answer; the answer is then published as one + * chunk (per-token streaming of the tool-loop answer is a follow-up).

+ * + *

{@link #onEvent(CloudEvent, String)} is the event-driven trigger path: a subscribed external + * event is turned into a prompt, answered (with tools if registered), and the answer is published + * as a CloudEvent onto an output topic — where sink connectors can deliver it to external systems.

+ * *

Mode 2 (publish/subscribe) has no agent involvement — a gateway publishes chunks directly to a * per-session lite topic; consumers subscribe via the runtime. This class is mode-1-only.

*/ @Slf4j public class StreamingAgent { + private static final ObjectMapper MAPPER = new ObjectMapper(); + + private static final int MAX_TOOL_ITERATIONS = 5; + private final CloudEventsClient client; private final String agentParent; private final String agentId; - private final OpenAiLlmClient llm; - private final ConversationStore store; + private final LlmClient llm; + private final ConversationMemory store; + private final ToolRegistry tools; /** One virtual thread per in-flight stream. */ private final ExecutorService streamExecutor = Executors.newThreadPerTaskExecutor( Thread.ofVirtual().name("em-agent-stream-", 1).factory()); private final AtomicInteger activeSessions = new AtomicInteger(0); public StreamingAgent(CloudEventsClient client, String agentParent, String agentId, - OpenAiLlmClient llm, ConversationStore store) { + LlmClient llm, ConversationMemory store) { + this(client, agentParent, agentId, llm, store, new ToolRegistry()); + } + + public StreamingAgent(CloudEventsClient client, String agentParent, String agentId, + LlmClient llm, ConversationMemory store, ToolRegistry tools) { this.client = client; this.agentParent = agentParent; this.agentId = agentId; this.llm = llm; this.store = store; + this.tools = tools == null ? new ToolRegistry() : tools; } /** Subscribe to this agent's channel; runs until {@link #shutdown()}. */ @@ -68,7 +100,8 @@ public void start() { // agent.; decode each frame to a StreamRequest. client.subscribeLiteBytes(agentParent, "agent." + agentId, frame -> streamExecutor.submit(() -> handleRequest(frame))); - log.info("StreamingAgent subscribed (private-wire): parent={} lite=agent.{}", agentParent, agentId); + log.info("StreamingAgent subscribed (private-wire): parent={} lite=agent.{} tools={}", + agentParent, agentId, tools.size()); } /** In-flight stream count, reported to the runtime via heartbeat. */ @@ -82,16 +115,22 @@ private void handleRequest(byte[] frame) { String replyTo = req.getReplyTo(); int[] seq = {0}; activeSessions.incrementAndGet(); - log.info("stream request: sessionId={} replyTo={} model={}", sessionId, replyTo, req.getModel()); + log.info("stream request: sessionId={} replyTo={} model={} tools={}", sessionId, replyTo, req.getModel(), + tools.size()); // Multi-turn: prepend conversation history (empty for a new sessionId), then this turn's prompt. List> messages = store.get(sessionId); messages.add(message("user", req.getPrompt())); StringBuilder answer = new StringBuilder(); try { - llm.stream(messages, req.getModel(), token -> { - answer.append(token); - publish(replyTo, sessionId, seq, token, false, null); - }); + if (tools.isEmpty()) { + llm.stream(messages, req.getModel(), token -> { + answer.append(token); + publish(replyTo, sessionId, seq, token, false, null); + }); + } else { + answer.append(runToolLoop(messages, req.getModel())); + publish(replyTo, sessionId, seq, answer.toString(), false, null); + } publish(replyTo, sessionId, seq, "", true, null); log.info("stream completed: sessionId={} chunks={}", sessionId, seq[0] - 1); store.appendTurn(sessionId, req.getPrompt(), answer.toString()); @@ -103,6 +142,76 @@ private void handleRequest(byte[] frame) { } } + /** + * Event-driven trigger path: render the incoming CloudEvent as a prompt, answer it (with the + * registered tools, if any), and publish the answer as a CloudEvent onto {@code outputTopic} — + * the topic sink connectors subscribe to for external delivery. Each event is its own + * conversation ({@code trigger:}) so triggers never bleed into user sessions. + */ + public void onEvent(CloudEvent event, String outputTopic) { + streamExecutor.submit(() -> { + String data = event.getData() == null ? "{}" + : new String(event.getData().toBytes(), StandardCharsets.UTF_8); + String prompt = "Incoming event (type=" + event.getType() + "):\n" + data + + "\n\nProcess this event and decide the follow-up action."; + String conversationId = "trigger:" + event.getId(); + List> messages = store.get(conversationId); + messages.add(message("user", prompt)); + try { + String answer = tools.isEmpty() ? aggregate(messages) : runToolLoop(messages, null); + store.appendTurn(conversationId, prompt, answer); + CloudEvent out = CloudEventBuilder.v1() + .withId(UUID.randomUUID().toString()) + .withSource(java.net.URI.create("urn:eventmesh:agent:" + agentId)) + .withType("org.apache.eventmesh.agent.trigger.output") + .withDataContentType("application/json") + .withData(MAPPER.writeValueAsBytes(Map.of("agentId", agentId, "answer", answer))) + .build(); + client.publish(outputTopic, out); + log.info("trigger processed: eventId={} -> topic={} answerChars={}", event.getId(), outputTopic, + answer.length()); + } catch (Exception e) { + log.warn("trigger failed: eventId={} err={}", event.getId(), e.toString()); + } + }); + } + + /** Function-calling loop: chat → (tool → feed result back)* → final answer text. */ + private String runToolLoop(List> messages, String model) throws Exception { + for (int i = 0; i < MAX_TOOL_ITERATIONS; i++) { + LlmCompletion completion = llm.chat(messages, tools.specs(), model); + if (completion.text() != null) { + return completion.text(); + } + ToolCall call = completion.toolCall(); + String result; + try { + result = tools.invoke(call.function(), parseArgs(call.argumentsJson())); + } catch (Exception e) { + result = "tool error: " + e; + } + // Tool results are fed back as user messages: keeps the LlmClient message shape + // (role+content maps) provider-neutral. + messages.add(message("user", "Tool result for `" + call.function() + "`: " + result)); + log.debug("tool loop: iter={} tool={} resultChars={}", i, call.function(), result.length()); + } + return "(max tool iterations reached without a final answer)"; + } + + private String aggregate(List> messages) throws Exception { + StringBuilder sb = new StringBuilder(); + llm.stream(messages, null, sb::append); + return sb.toString(); + } + + @SuppressWarnings("unchecked") + private static Map parseArgs(String argumentsJson) throws Exception { + if (argumentsJson == null || argumentsJson.isBlank()) { + return Map.of(); + } + return MAPPER.readValue(argumentsJson, Map.class); + } + private void publish(String replyTo, String sessionId, int[] seq, String chunk, boolean done, String error) { StreamChunk c = StreamChunk.builder() .sessionId(sessionId).seq(seq[0]++).chunk(chunk).done(done).error(error).build(); diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/LlmClient.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/LlmClient.java new file mode 100644 index 0000000000..d144fcb545 --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/LlmClient.java @@ -0,0 +1,50 @@ +/* + * 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.llm; + +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; + +/** + * SPI for chat-completion backends. {@link OpenAiLlmClient} (OpenAI-compatible SSE) is the default + * implementation; implement this interface and hand the instance to {@code StreamingAgent} to plug + * in any non-OpenAI-protocol LLM (Anthropic native, Bedrock, an internal inference service, ...). + * + *

Two operations: {@link #stream} for token streaming (no tools), {@link #chat} for a + * single-shot completion that may request a function call.

+ */ +public interface LlmClient { + + /** + * Stream the completion for the given message list (each map holds {@code role} + {@code + * content}); {@code chunkCb} receives token fragments in order. Blocks until the stream ends. + * + * @param model overrides the client default when non-null/non-empty + */ + void stream(List> messages, String model, Consumer chunkCb) throws Exception; + + /** + * Single-shot completion with optional function-calling tools. Returns either assistant text + * ({@link LlmCompletion#text()}) or one tool call the model wants performed first + * ({@link LlmCompletion#toolCall()}); the caller executes the tool and loops. + * + * @param tools may be empty (plain completion, no tool advertisement) + */ + LlmCompletion chat(List> messages, List tools, String model) throws Exception; +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/LlmCompletion.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/LlmCompletion.java new file mode 100644 index 0000000000..f2a207ea2d --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/LlmCompletion.java @@ -0,0 +1,34 @@ +/* + * 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.llm; + +/** + * Either {@code text} (final answer, no tool call) or {@code toolCall} (execute the tool, loop + * again). + */ +public record LlmCompletion(String text, ToolCall toolCall) { + + public static LlmCompletion ofText(String text) { + return new LlmCompletion(text, null); + } + + public static LlmCompletion ofToolCall(ToolCall call) { + return new LlmCompletion(null, call); + } +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/OpenAiLlmClient.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/OpenAiLlmClient.java index bf2b718337..1f839a4cdb 100644 --- a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/OpenAiLlmClient.java +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/OpenAiLlmClient.java @@ -19,6 +19,8 @@ import java.net.URI; import java.time.Duration; +import java.util.List; +import java.util.Map; import java.util.function.Consumer; import java.util.stream.Stream; @@ -30,14 +32,15 @@ import lombok.extern.slf4j.Slf4j; /** - * Minimal OpenAI-compatible chat-completion streaming client. Issues {@code POST {base}/v1/chat/ - * completions} with {@code stream:true} and {@code Authorization: Bearer }, reads the SSE token - * stream and invokes {@code chunkCb} per {@code choices[0].delta.content}. Compatible with OpenAI, - * Azure-OpenAI, vLLM, Ollama (OpenAI mode), DeepSeek, Moonshot and most internal gateways. Blocking; - * intended to run on a virtual thread. + * Minimal OpenAI-compatible chat-completion client — the default {@link LlmClient} implementation. + * Issues {@code POST {base}/v1/chat/completions}; with {@code stream:true} it reads the SSE token + * stream and invokes {@code chunkCb} per {@code choices[0].delta.content}; via {@link #chat} it + * performs a single-shot completion that may return a function call. Compatible with OpenAI, + * Azure-OpenAI, vLLM, Ollama (OpenAI mode), DeepSeek, Moonshot and most internal gateways. + * Blocking; intended to run on a virtual thread. */ @Slf4j -public class OpenAiLlmClient { +public class OpenAiLlmClient implements LlmClient { private static final ObjectMapper MAPPER = new ObjectMapper(); @@ -74,28 +77,12 @@ public void stream(String prompt, String model, Consumer chunkCb) throws * Stream a completion for the given message list (supports multi-turn: pass prior history + * the current user message). Each map must contain {@code role} + {@code content}. */ + @Override public void stream(java.util.List> messages, String model, Consumer chunkCb) throws Exception { - String useModel = (model != null && !model.isEmpty()) ? model : defaultModel; - ObjectNode body = MAPPER.createObjectNode(); - body.put("model", useModel); - ArrayNode msgs = body.putArray("messages"); - for (java.util.Map message : messages) { - ObjectNode msg = msgs.addObject(); - msg.put("role", message.getOrDefault("role", "user")); - msg.put("content", message.getOrDefault("content", "")); - } + ObjectNode body = baseBody(messages, model); body.put("stream", true); - - java.net.http.HttpRequest req = java.net.http.HttpRequest.newBuilder() - .uri(URI.create(trimSlash(this.baseUrl) + "/v1/chat/completions")) - .timeout(Duration.ofMinutes(3)) - .header("Content-Type", "application/json") - .header("Accept", "text/event-stream") - .header("Authorization", "Bearer " + apiKey) - .POST(java.net.http.HttpRequest.BodyPublishers.ofString(MAPPER.writeValueAsString(body))) - .build(); - + java.net.http.HttpRequest req = request(body); java.net.http.HttpResponse> resp = http.send(req, java.net.http.HttpResponse.BodyHandlers.ofLines()); int status = resp.statusCode(); @@ -127,6 +114,66 @@ public void stream(java.util.List> messages, Strin } } + /** + * Single-shot completion with optional function-calling tools (non-streaming): parses + * {@code choices[0].message} into either assistant text or the first requested tool call. + */ + @Override + public LlmCompletion chat(List> messages, List tools, + String model) throws Exception { + ObjectNode body = baseBody(messages, model); + if (tools != null && !tools.isEmpty()) { + ArrayNode toolsArr = body.putArray("tools"); + for (ToolSpec spec : tools) { + ObjectNode fn = toolsArr.addObject().put("type", "function").putObject("function"); + fn.put("name", spec.name()); + fn.put("description", spec.description()); + fn.set("parameters", MAPPER.readTree( + spec.parametersJsonSchema() == null ? "{\"type\":\"object\"}" : spec.parametersJsonSchema())); + } + } + java.net.http.HttpResponse resp = + http.send(request(body), java.net.http.HttpResponse.BodyHandlers.ofString()); + int status = resp.statusCode(); + if (status != 200) { + throw new RuntimeException("LLM HTTP " + status + " (check llm.base.url / llm.api.key / llm.model)"); + } + JsonNode message = MAPPER.readTree(resp.body()).path("choices").path(0).path("message"); + JsonNode toolCalls = message.get("tool_calls"); + if (toolCalls != null && toolCalls.isArray() && toolCalls.size() > 0) { + JsonNode call = toolCalls.get(0).path("function"); + return LlmCompletion.ofToolCall(new ToolCall( + toolCalls.get(0).path("id").asText(""), + call.path("name").asText(""), + call.path("arguments").asText("{}"))); + } + return LlmCompletion.ofText(message.path("content").asText("")); + } + + private ObjectNode baseBody(List> messages, String model) { + String useModel = (model != null && !model.isEmpty()) ? model : defaultModel; + ObjectNode body = MAPPER.createObjectNode(); + body.put("model", useModel); + ArrayNode msgs = body.putArray("messages"); + for (Map message : messages) { + ObjectNode msg = msgs.addObject(); + msg.put("role", message.getOrDefault("role", "user")); + msg.put("content", message.getOrDefault("content", "")); + } + return body; + } + + private java.net.http.HttpRequest request(ObjectNode body) throws Exception { + return java.net.http.HttpRequest.newBuilder() + .uri(URI.create(trimSlash(this.baseUrl) + "/v1/chat/completions")) + .timeout(Duration.ofMinutes(3)) + .header("Content-Type", "application/json") + .header("Accept", "application/json") + .header("Authorization", "Bearer " + apiKey) + .POST(java.net.http.HttpRequest.BodyPublishers.ofString(MAPPER.writeValueAsString(body))) + .build(); + } + private static String trimSlash(String url) { return url.endsWith("/") ? url.substring(0, url.length() - 1) : url; } diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/ToolCall.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/ToolCall.java new file mode 100644 index 0000000000..58c91efa80 --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/ToolCall.java @@ -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. + */ + + +package org.apache.eventmesh.agent.llm; + +/** A tool call the model requested: provider call id, function name, raw JSON args. */ +public record ToolCall(String id, String function, String argumentsJson) { +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/ToolSpec.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/ToolSpec.java new file mode 100644 index 0000000000..c3578a4264 --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/llm/ToolSpec.java @@ -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. + */ + + +package org.apache.eventmesh.agent.llm; + +/** A function-calling tool advertisement: name + description + JSON-schema parameters. */ +public record ToolSpec(String name, String description, String parametersJsonSchema) { +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/AgentTool.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/AgentTool.java new file mode 100644 index 0000000000..a40017b48b --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/AgentTool.java @@ -0,0 +1,41 @@ +/* + * 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.tool; + +import java.util.Map; + +/** + * An invocable capability the LLM may call via function calling during a + * {@code StreamingAgent} tool loop. Register instances on a {@link ToolRegistry} and pass the + * registry to the agent; the agent advertises every tool to the model, executes the calls the + * model requests, and feeds results back until a final answer. + */ +public interface AgentTool { + + /** Stable tool name the model addresses this tool by. */ + String name(); + + /** One-sentence description telling the model when to use this tool. */ + String description(); + + /** JSON schema (type "object") describing the arguments object; may be permissive. */ + String parametersJsonSchema(); + + /** Invoke with the model-produced arguments object; return a text result for the model. */ + String invoke(Map args) throws Exception; +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/ConnectorToolAdapter.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/ConnectorToolAdapter.java new file mode 100644 index 0000000000..316a36ea80 --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/ConnectorToolAdapter.java @@ -0,0 +1,179 @@ +/* + * 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.tool; + +import org.apache.eventmesh.connector.SinkConnector; +import org.apache.eventmesh.connector.SourceConnector; + +import java.net.URI; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.UUID; + +import io.cloudevents.CloudEvent; +import io.cloudevents.core.builder.CloudEventBuilder; + +import com.fasterxml.jackson.databind.ObjectMapper; + +/** + * Adapts the existing connector SPI ({@link SinkConnector} / {@link SourceConnector}) into + * {@link AgentTool}s, turning every shipped connector plugin into a capability the LLM can call. + * + *
    + *
  • {@link #sinkTool}: a write tool — each invocation wraps the arguments object as one + * CloudEvent and pushes it through the sink ({@code put} + {@code commit}).
  • + *
  • {@link #sourceTool}: a read tool — each invocation polls one batch + * ({@code poll()}) and returns it as a JSON array of events.
  • + *
+ * + *

The connector instance is initialized once and reused across invocations (same lifecycle as + * inside the connector runtime).

+ */ +public final class ConnectorToolAdapter { + + private static final ObjectMapper MAPPER = new ObjectMapper(); + + private static final String PERMISSIVE_SCHEMA = "{\"type\":\"object\",\"properties\":{}," + + "\"additionalProperties\":true,\"description\":\"Arguments delivered as the event payload\"}"; + + private ConnectorToolAdapter() { + } + + /** Build a write tool on top of a {@link SinkConnector}; the connector is initialized eagerly. */ + public static AgentTool sinkTool(String name, String description, Class sinkClass, + Properties props) { + return new SinkTool(name, description, newInstance(SinkConnector.class, sinkClass), props); + } + + /** Build a write tool on top of a pre-built {@link SinkConnector} instance (already initialized). */ + public static AgentTool sinkTool(String name, String description, SinkConnector sink) { + return new SinkTool(name, description, sink, null); + } + + /** Build a read tool on top of a {@link SourceConnector}; the connector is initialized eagerly. */ + public static AgentTool sourceTool(String name, String description, Class sourceClass, + Properties props) { + return new SourceTool(name, description, newInstance(SourceConnector.class, sourceClass), props); + } + + /** Build a read tool on top of a pre-built {@link SourceConnector} instance (already initialized). */ + public static AgentTool sourceTool(String name, String description, SourceConnector source) { + return new SourceTool(name, description, source, null); + } + + private static T newInstance(Class spi, Class impl) { + try { + return impl.getDeclaredConstructor().newInstance(); + } catch (ReflectiveOperationException e) { + throw new IllegalArgumentException("cannot instantiate connector " + impl.getName(), e); + } + } + + private static final class SinkTool implements AgentTool { + + private final String name; + private final String description; + private final SinkConnector sink; + + private SinkTool(String name, String description, SinkConnector sink, Properties props) { + this.name = name; + this.description = description; + this.sink = sink; + if (props != null) { + sink.init(props); + } + } + + @Override + public String name() { + return name; + } + + @Override + public String description() { + return description; + } + + @Override + public String parametersJsonSchema() { + return PERMISSIVE_SCHEMA; + } + + @Override + public String invoke(Map args) throws Exception { + CloudEvent event = CloudEventBuilder.v1() + .withId(UUID.randomUUID().toString()) + .withSource(URI.create("urn:eventmesh:agent-tool:" + name)) + .withType("org.apache.eventmesh.agent.tool.invoke") + .withDataContentType("application/json") + .withData(MAPPER.writeValueAsBytes(args == null ? Map.of() : args)) + .build(); + sink.put(List.of(event)); + sink.commit(List.of(event)); + return "delivered"; + } + } + + private static final class SourceTool implements AgentTool { + + private final String name; + private final String description; + private final SourceConnector source; + + private SourceTool(String name, String description, SourceConnector source, Properties props) { + this.name = name; + this.description = description; + this.source = source; + if (props != null) { + source.init(props); + } + } + + @Override + public String name() { + return name; + } + + @Override + public String description() { + return description; + } + + @Override + public String parametersJsonSchema() { + return PERMISSIVE_SCHEMA; + } + + @Override + public String invoke(Map args) throws Exception { + List batch = source.poll(); + List payloads = new ArrayList<>(); + for (CloudEvent event : batch) { + payloads.add(event.getData() == null ? "{}" + : new String(event.getData().toBytes(), StandardCharsets.UTF_8)); + } + if (!batch.isEmpty()) { + source.commit(batch.get(batch.size() - 1)); + } + return MAPPER.writeValueAsString(payloads); + } + } +} diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/ToolRegistry.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/ToolRegistry.java new file mode 100644 index 0000000000..a72028c066 --- /dev/null +++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/tool/ToolRegistry.java @@ -0,0 +1,62 @@ +/* + * 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.tool; + +import org.apache.eventmesh.agent.llm.ToolSpec; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** Named, ordered set of {@link AgentTool}s; renders the tool list for the LLM tool loop. */ +public class ToolRegistry { + + private final Map tools = new LinkedHashMap<>(); + + public ToolRegistry register(AgentTool tool) { + tools.put(tool.name(), tool); + return this; + } + + public boolean isEmpty() { + return tools.isEmpty(); + } + + public int size() { + return tools.size(); + } + + /** Run one tool by name; unknown names throw IllegalArgumentException. */ + public String invoke(String name, Map args) throws Exception { + AgentTool tool = tools.get(name); + if (tool == null) { + throw new IllegalArgumentException("unknown tool: " + name); + } + return tool.invoke(args); + } + + /** Render the registered tools as LLM tool advertisements (OpenAI function-calling shape). */ + public List specs() { + List specs = new ArrayList<>(); + for (AgentTool tool : tools.values()) { + specs.add(new ToolSpec(tool.name(), tool.description(), tool.parametersJsonSchema())); + } + return specs; + } +} diff --git a/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/llm/OpenAiLlmClientChatTest.java b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/llm/OpenAiLlmClientChatTest.java new file mode 100644 index 0000000000..87ebb479c6 --- /dev/null +++ b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/llm/OpenAiLlmClientChatTest.java @@ -0,0 +1,118 @@ +/* + * 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.llm; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; + +/** Drives {@link OpenAiLlmClient#chat} (function calling) against an in-process mock. Hermetic. */ +class OpenAiLlmClientChatTest { + + private HttpServer server; + private OpenAiLlmClient client; + private final List capturedBodies = new ArrayList<>(); + + @BeforeEach + void setUp() throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.start(); + client = new OpenAiLlmClient("http://127.0.0.1:" + server.getAddress().getPort(), "k", "m"); + } + + @AfterEach + void tearDown() { + if (server != null) { + server.stop(0); + } + } + + @Test + void chatReturnsToolCallWhenModelRequestsOne() throws Exception { + server.createContext("/v1/chat/completions", ex -> { + capture(ex); + respond(ex, "{\"choices\":[{\"message\":{\"role\":\"assistant\",\"content\":null," + + "\"tool_calls\":[{\"id\":\"call-1\",\"type\":\"function\"," + + "\"function\":{\"name\":\"notify\",\"arguments\":\"{\\\"text\\\":\\\"hi\\\"}\"}}]}}]}"); + }); + LlmCompletion completion = client.chat(List.of(Map.of("role", "user", "content", "go")), + List.of(new ToolSpec("notify", "send a notification", "{\"type\":\"object\"}")), null); + assertThat(completion.text()).isNull(); + assertThat(completion.toolCall().id()).isEqualTo("call-1"); + assertThat(completion.toolCall().function()).isEqualTo("notify"); + assertThat(completion.toolCall().argumentsJson()).isEqualTo("{\"text\":\"hi\"}"); + } + + @Test + void chatReturnsTextWhenNoToolCall() throws Exception { + server.createContext("/v1/chat/completions", ex -> { + capture(ex); + respond(ex, "{\"choices\":[{\"message\":{\"role\":\"assistant\",\"content\":\"all done\"}}]}"); + }); + LlmCompletion completion = client.chat(List.of(Map.of("role", "user", "content", "go")), + List.of(), null); + assertThat(completion.toolCall()).isNull(); + assertThat(completion.text()).isEqualTo("all done"); + } + + @Test + void chatAdvertisesToolsInTheRequestBody() throws Exception { + server.createContext("/v1/chat/completions", ex -> { + capture(ex); + respond(ex, "{\"choices\":[{\"message\":{\"role\":\"assistant\",\"content\":\"ok\"}}]}"); + }); + client.chat(List.of(Map.of("role", "user", "content", "go")), + List.of(new ToolSpec("poll-orders", "poll a batch", "{\"type\":\"object\"}")), "override-m"); + assertThat(capturedBodies).hasSize(1); + JsonNode body = capturedBodies.get(0); + assertThat(body.get("model").asText()).isEqualTo("override-m"); + assertThat(body.get("stream")).isNull(); + JsonNode fn = body.path("tools").path(0).path("function"); + assertThat(fn.get("name").asText()).isEqualTo("poll-orders"); + assertThat(fn.get("description").asText()).isEqualTo("poll a batch"); + assertThat(fn.path("parameters").get("type").asText()).isEqualTo("object"); + } + + private void capture(HttpExchange exchange) throws IOException { + capturedBodies.add(new ObjectMapper().readTree(exchange.getRequestBody().readAllBytes())); + } + + private void respond(HttpExchange exchange, String body) throws IOException { + byte[] bytes = body.getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().add("Content-Type", "application/json"); + exchange.sendResponseHeaders(200, bytes.length); + try (OutputStream os = exchange.getResponseBody()) { + os.write(bytes); + } + } +} diff --git a/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/tool/ConnectorToolAdapterTest.java b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/tool/ConnectorToolAdapterTest.java new file mode 100644 index 0000000000..09719c0593 --- /dev/null +++ b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/tool/ConnectorToolAdapterTest.java @@ -0,0 +1,117 @@ +/* + * 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.tool; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.apache.eventmesh.connector.SinkConnector; +import org.apache.eventmesh.connector.SourceConnector; + +import java.nio.charset.StandardCharsets; +import java.util.List; +import java.util.Map; +import java.util.Properties; + +import org.junit.jupiter.api.Test; + +import io.cloudevents.CloudEvent; +import io.cloudevents.core.builder.CloudEventBuilder; + +/** Covers the sink (write) and source (read) adapter shapes against fake connectors. */ +class ConnectorToolAdapterTest { + + private static final java.util.UUID ID = java.util.UUID.randomUUID(); + + static final class RecordingSink implements SinkConnector { + + volatile List lastWritten; + + @Override + public void init(Properties props) { + } + + @Override + public void put(List events) { + lastWritten = events; + } + + @Override + public void commit(List written) { + } + } + + static final class OneBatchSource implements SourceConnector { + + volatile boolean polled; + + @Override + public void init(Properties props) { + } + + @Override + public List poll() { + if (polled) { + return List.of(); + } + polled = true; + return List.of( + CloudEventBuilder.v1().withId(ID.toString()).withSource(java.net.URI.create("urn:test")) + .withType("t1").withDataContentType("application/json") + .withData("{\"k\":\"v1\"}".getBytes(StandardCharsets.UTF_8)).build(), + CloudEventBuilder.v1().withId(ID.toString()).withSource(java.net.URI.create("urn:test")) + .withType("t2").withDataContentType("application/json") + .withData("{\"k\":\"v2\"}".getBytes(StandardCharsets.UTF_8)).build()); + } + + @Override + public void commit(CloudEvent lastPublished) { + } + } + + @Test + void sinkToolDeliversArgsAsCloudEvent() throws Exception { + AgentTool tool = ConnectorToolAdapter.sinkTool("notify", + "Deliver a payload", RecordingSink.class, new Properties()); + assertThat(tool.name()).isEqualTo("notify"); + assertThat(tool.parametersJsonSchema()).contains("\"object\""); + + String result = tool.invoke(Map.of("text", "high risk order")); + assertThat(result).isEqualTo("delivered"); + } + + @Test + void sinkToolWritesTheArgumentsObjectAsEventData() throws Exception { + RecordingSink sink = new RecordingSink(); + AgentTool tool = ConnectorToolAdapter.sinkTool("notify", "d", sink); + tool.invoke(Map.of("order", 42, "risk", "high")); + assertThat(sink.lastWritten).hasSize(1); + assertThat(new String(sink.lastWritten.get(0).getData().toBytes(), StandardCharsets.UTF_8)) + .contains("\"order\":42").contains("\"risk\":\"high\""); + } + + @Test + void sourceToolReturnsBatchAsJsonArray() throws Exception { + AgentTool tool = ConnectorToolAdapter.sourceTool("poll-orders", + "Poll the next batch", OneBatchSource.class, new Properties()); + String first = tool.invoke(Map.of()); + // payloads are JSON strings inside a JSON array, so inner quotes are escaped + assertThat(first).contains("{\\\"k\\\":\\\"v1\\\"}").contains("{\\\"k\\\":\\\"v2\\\"}"); + String second = tool.invoke(Map.of()); + assertThat(second).isEqualTo("[]"); + } +} diff --git a/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/tool/ToolRegistryTest.java b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/tool/ToolRegistryTest.java new file mode 100644 index 0000000000..cb2a2a90f5 --- /dev/null +++ b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/tool/ToolRegistryTest.java @@ -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.tool; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.util.Map; + +import org.junit.jupiter.api.Test; + +/** Covers registry lookup, spec rendering and the unknown-tool error path. */ +class ToolRegistryTest { + + private static AgentTool fixedTool(String name, String result) { + return new AgentTool() { + + @Override + public String name() { + return name; + } + + @Override + public String description() { + return "test tool " + name; + } + + @Override + public String parametersJsonSchema() { + return "{\"type\":\"object\"}"; + } + + @Override + public String invoke(Map args) { + return result; + } + }; + } + + @Test + void invokesByName() throws Exception { + ToolRegistry registry = new ToolRegistry().register(fixedTool("echo", "ok")); + assertThat(registry.invoke("echo", Map.of())).isEqualTo("ok"); + assertThat(registry.size()).isEqualTo(1); + assertThat(registry.isEmpty()).isFalse(); + } + + @Test + void unknownToolThrows() { + ToolRegistry registry = new ToolRegistry(); + assertThatThrownBy(() -> registry.invoke("nope", Map.of())) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("unknown tool"); + } + + @Test + void specsRenderNameDescriptionSchema() { + ToolRegistry registry = new ToolRegistry().register(fixedTool("a", "x")).register(fixedTool("b", "y")); + assertThat(registry.specs()).hasSize(2); + assertThat(registry.specs().get(0).name()).isEqualTo("a"); + assertThat(registry.specs().get(1).description()).isEqualTo("test tool b"); + assertThat(registry.specs().get(0).parametersJsonSchema()).isEqualTo("{\"type\":\"object\"}"); + } +} diff --git a/eventmesh-connector-runtime/bin/start-connector.sh b/eventmesh-connector-runtime/bin/start-connector.sh index 47f48167d7..779cf96a96 100644 --- a/eventmesh-connector-runtime/bin/start-connector.sh +++ b/eventmesh-connector-runtime/bin/start-connector.sh @@ -23,7 +23,7 @@ # (unlike the runtime) connector jars MUST be on the -cp. # # Env vars (all optional): -# EVENTMESH_RUNTIME_URL EventMesh runtime URL (default http://localhost:8080) +# EVENTMESH_RUNTIME_URL EventMesh runtime URL (default http://localhost:10105) # CONNECTOR_ADMIN_PORT connector admin HTTP port (default 0 = off) # CONNECTOR_OFFSET_MODE remote | rocksdb | inmemory (default remote) # CONNECTOR_OFFSET_PATH rocksdb offset path (default $EVENTMESH_HOME/data/connector-offset) @@ -36,7 +36,7 @@ SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" EVENTMESH_HOME="${EVENTMESH_HOME:-$(cd "$SCRIPT_DIR/.." && pwd)}" cd "$EVENTMESH_HOME" -EVENTMESH_RUNTIME_URL="${EVENTMESH_RUNTIME_URL:-http://localhost:8080}" +EVENTMESH_RUNTIME_URL="${EVENTMESH_RUNTIME_URL:-http://localhost:10105}" CONNECTOR_ADMIN_PORT="${CONNECTOR_ADMIN_PORT:-0}" CONNECTOR_OFFSET_MODE="${CONNECTOR_OFFSET_MODE:-remote}" CONNECTOR_OFFSET_PATH="${CONNECTOR_OFFSET_PATH:-$EVENTMESH_HOME/data/connector-offset}" diff --git a/eventmesh-connector-runtime/conf/connector.properties b/eventmesh-connector-runtime/conf/connector.properties index bed5bde44d..284809af4e 100644 --- a/eventmesh-connector-runtime/conf/connector.properties +++ b/eventmesh-connector-runtime/conf/connector.properties @@ -20,7 +20,7 @@ # translates the ENV vars below into -D flags; pass connector-specific -D flags through the # CONNECTOR_OPTS env var. # -# EVENTMESH_RUNTIME_URL -> -Deventmesh.runtime.url (default http://localhost:8080) +# EVENTMESH_RUNTIME_URL -> -Deventmesh.runtime.url (default http://localhost:10105) # CONNECTOR_ADMIN_PORT -> -Dconnector.admin.port (default 0 = off) # CONNECTOR_OFFSET_MODE -> -Dconnector.offset.mode (remote | rocksdb | inmemory; default remote) # CONNECTOR_OFFSET_PATH -> -Dconnector.offset.path (rocksdb path) diff --git a/eventmesh-connector-runtime/src/main/java/org/apache/eventmesh/connector/ConnectorApplication.java b/eventmesh-connector-runtime/src/main/java/org/apache/eventmesh/connector/ConnectorApplication.java index e7a46067cd..f3ef1f7317 100644 --- a/eventmesh-connector-runtime/src/main/java/org/apache/eventmesh/connector/ConnectorApplication.java +++ b/eventmesh-connector-runtime/src/main/java/org/apache/eventmesh/connector/ConnectorApplication.java @@ -54,7 +54,7 @@ public class ConnectorApplication { private static final ObjectMapper MAPPER = new ObjectMapper(); public static void main(String[] args) throws Exception { - String runtimeUrl = System.getProperty("eventmesh.runtime.url", "http://localhost:8080"); + String runtimeUrl = System.getProperty("eventmesh.runtime.url", "http://localhost:10105"); String offsetPath = System.getProperty("connector.offset.path"); final int adminPort = Integer.getInteger("connector.admin.port", 0); final String workerId = System.getProperty("connector.worker.id", ""); diff --git a/eventmesh-examples/src/main/java/org/apache/eventmesh/cloudevents/demo/stream/StreamingCallDemo.java b/eventmesh-examples/src/main/java/org/apache/eventmesh/cloudevents/demo/stream/StreamingCallDemo.java index fe35846c8f..43b95a5bfb 100644 --- a/eventmesh-examples/src/main/java/org/apache/eventmesh/cloudevents/demo/stream/StreamingCallDemo.java +++ b/eventmesh-examples/src/main/java/org/apache/eventmesh/cloudevents/demo/stream/StreamingCallDemo.java @@ -36,7 +36,7 @@ * *

Prerequisites: *

    - *
  1. A running EventMesh Runtime on {@code runtimeUrl} (default {@code http://localhost:8080}), + *
  2. A running EventMesh Runtime on {@code runtimeUrl} (default {@code http://localhost:10105}), * wired with a {@code SessionRouter} (see docs/sdk-streaming-call-guide.md §9).
  3. *
  4. At least one streaming Agent registered + heartbeating against that runtime.
  5. *
@@ -44,7 +44,7 @@ *

Run: *

  *   java org.apache.eventmesh.cloudevents.demo.stream.StreamingCallDemo [runtimeUrl] [prompt]
- *   # defaults: runtimeUrl=http://localhost:8080  prompt="Introduce Apache EventMesh in three sentences."
+ *   # defaults: runtimeUrl=http://localhost:10105  prompt="Introduce Apache EventMesh in three sentences."
  * 
* *

The demo exercises {@code forEach}-only consumption (the sole posture) on a single-turn call, @@ -54,7 +54,7 @@ @Slf4j public class StreamingCallDemo { - private static final String DEFAULT_RUNTIME_URL = "http://localhost:8080"; + private static final String DEFAULT_RUNTIME_URL = "http://localhost:10105"; private static final String DEFAULT_PROMPT = "Introduce Apache EventMesh in three sentences."; public static void main(String[] args) throws Exception {