介绍一下merged pr:langchain4j#6306
为langchain4j agentic a2a增加了多artifact响应的支持,并添加了流式事件监听器。
Note (pr url)
a2a协议支持的交互模式
非streaming消息模式
The primary operation for initiating agent interactions. Clients send a message to an agent and receive either a task that tracks the processing or a direct response message. The agent MAY create a new Task to process the provided message asynchronously or MAY return a direct Message response for simple interactions. The operation MUST return immediately with either task information or response message. Task processing MAY continue asynchronously after the response when a Task is returned.
这种最基本的非streaming类的消息,客户端发送请求给agent之后,agent会返回一个task和一个message,一个任务仅仅对应一条消息,在a2a-java中,此类消息为TaskEvent。在TaskEvent中,一次性包含所有的artifact和history。此类请求一般不能够多次请求,一次task对应一次请求。
对于服务端,只需要在a2a card中设置streaming(false)即可。可以在里面使用单一message或update模式。
streaming消息模式 单一message
Message-only stream: If the agent returns a
Message, the stream MUST contain exactly oneMessageobject and then close immediately. No task tracking or updates are provided.
对于这类请求每次请求同样只返回一个message,但是可以对同一个taskId进行多次请求。也就是说message的状态可以设置为TASK_STATE_AUTH_REQUIRED或者TASK_STATE_INPUT_REQUIRED要求客户端再次请求,提供相应的信息。客户端对于单一的message包装为MessageEvent。
对于服务端,使用emitter.sendMessage(msg)来发送一条完整的Message,且发送Message前后都不可以使用其他update模式的消息。
streaming消息模式 update模式
Task lifecycle stream: If the agent returns a
Task, the stream MUST begin with the Task object, followed by zero or moreTaskStatusUpdateEventorTaskArtifactUpdateEventobjects. The stream MUST close when the task reaches a terminal state (TASK_STATE_COMPLETED,TASK_STATE_FAILED,TASK_STATE_CANCELED,TASK_STATE_REJECTED).
对于同一个task的同一次请求,a2a server可以实时向客户端更新状态,包括task的状态更新和artifact的更新。客户端对于这种模式包装为TaskUpdateEvent,其中包含了task的状态和artifact的状态,分别包装为TaskStatusUpdateEvent和TaskArtifactUpdateEvent。客户端需要在受到此次请求的interrupt state或terminal state之前,持续监听服务器发送的TaskUpdateEvent。
对于服务端,使用emitter.addArtifact(parts)和emitter.updateStatus(parts)来发送task的状态更新和artifact的更新。
langchain4j原有的bug
langchain4j原有的时间消费器的实现如下所示,忽略了task的状态,当收到artifact时,会立即返回结果到上层,不会继续监听task的状态,缺失了后续的状态更新和artifacts。
另外用户对event的生命周期是无法修改的,例如不能快速退出,不能实时检测task的状态,从而做出相应的处理。
List<BiConsumer<ClientEvent, AgentCard>> consumers = List.of((event, card) -> { if (event instanceof MessageEvent messageEvent) { Message msg = messageEvent.getMessage(); responseContextId.set(msg.contextId()); responseTaskId.set(msg.taskId()); messageResponse.complete(extractTextFromParts(msg.parts())); } else if (event instanceof TaskEvent taskEvent) { captureTaskIds(taskEvent.getTask(), responseContextId, responseTaskId); completeFromTask(taskEvent.getTask(), messageResponse); } else if (event instanceof TaskUpdateEvent updateEvent) { captureTaskIds(updateEvent.getTask(), responseContextId, responseTaskId); completeFromTask(updateEvent.getTask(), messageResponse); } else { messageResponse.completeExceptionally( new IllegalArgumentException("The event expected should be of type " + event.getClass())); }});static void completeFromTask(Task task, CompletableFuture<String> messageResponse) { TaskState state = task.status().state(); if (state.isInterrupted()) { // The remote agent paused the task (input-required/auth-required) and will not advance it // on its own, so this must complete exceptionally rather than being left pending forever. Message statusMessage = task.status().message(); String reason = statusMessage != null ? extractTextFromParts(statusMessage.parts()) : ""; messageResponse.completeExceptionally( new A2ATaskInterruptedException(task.id(), task.contextId(), state, reason)); return; } if (!isTerminalState(state) && task.artifacts().isEmpty()) { return; } if (isFailureState(state)) { Message statusMessage = task.status().message(); String reason = statusMessage != null ? extractTextFromParts(statusMessage.parts()) : ""; messageResponse.completeExceptionally(new RuntimeException("A2A task " + task.id() + " ended in terminal state " + state + (reason.isEmpty() ? "" : ": " + reason))); return; } messageResponse.complete(extractTextFromParts( task.artifacts().stream().flatMap(a -> a.parts().stream()).toList()));}此次修改
理顺三种事件的生命周期
为三种事件分别做出区分,非流式或message模式受到消息后立刻返还给上层。
private void handleTaskEvent(TaskEvent taskEvent, CompletableFuture<String> messageResponse) { completeFromTask(taskEvent.getTask(), messageResponse);}
private void handleMessageEvent(Message message, CompletableFuture<String> messageResponse) { messageResponse.complete(extractTextFromParts(message.parts()));}
private void handleUpdateEvent( TaskUpdateEvent taskUpdateEvent, CompletableFuture<String> messageResponse, AtomicBoolean stopped) { if (streamingClientListener != null) { A2AStreamingClientListenerResult listenerResult = streamingClientListener.onUpdateEvent(taskUpdateEvent); if (listenerResult.stop()) { if (listenerResult.withCurrentArtifacts()) { completeArtifact(taskUpdateEvent.getTask().artifacts(), messageResponse); } else { String response = listenerResult.response(); messageResponse.complete(response == null ? "" : response); } stopped.set(true); return; } } completeFromTask(taskUpdateEvent.getTask(), messageResponse);}private static void completeArtifact(List<Artifact> artifacts, CompletableFuture<String> messageResponse) { if (artifacts == null || artifacts.isEmpty()) { messageResponse.complete(""); return; } messageResponse.complete(extractTextFromParts(artifacts.stream() .flatMap(artifact -> artifact.parts().stream()) .toList()));}为流式事件添加监听器
@FunctionalInterfacepublic interface A2AStreamingClientListener<E> { A2AStreamingClientListenerResult onUpdateEvent(E event);}public record A2AStreamingClientListenerResult(boolean stop, String response, boolean withCurrentArtifacts) { public static A2AStreamingClientListenerResult continueStreaming() { return new A2AStreamingClientListenerResult(false, null, false); }
public static A2AStreamingClientListenerResult stopWithResponse(String response) { return new A2AStreamingClientListenerResult(true, response, false); }
public static A2AStreamingClientListenerResult stopWithCurrentArtifacts() { return new A2AStreamingClientListenerResult(true, null, true); }}如上,为流式事件添加了一个监听器,用户可以在监听器中对事件进行处理,并决定是否停止流式事件的继续处理。只需要在a2aAgentBuilder当中,设置A2AStreamingClientListener即可。其返回可以为上述三种情况。
文档中的使用方法介绍为:
When invoking a remote A2A agent with streaming enabled, you can use
A2AStreamingClientListenerto observe events received from the remote agent and control when the client should stop consuming the stream.
This is useful when you only need to react to specific events instead of waiting for the remote task to finish. For example, you can stop listening when the remote agent requests additional input and return the message to the caller immediately.
The listener is invoked for each event received from the remote A2A agent. Return
continueStreaming()to keep consuming events,stopWithResponse(response)to stop consuming the stream and return the specified response to the caller, orstopWithCurrentArtifacts()to stop consuming the stream and return the artifacts received so far, using the default A2A client artifact-to-text extraction logic.
Stopping the client-side stream does not cancel the remote A2A task. The remote task may continue executing asynchronously.
UntypedAgent creativeWriter = AgenticServices.a2aBuilder(A2A_SERVER_URL) .inputKeys("topic") .outputKey("story") .streamingClientListener((TaskUpdateEvent event) ->{ UpdateEvent updateEvent = event.getUpdateEvent(); if (updateEvent instanceof TaskStatusUpdateEvent taskStatusUpdateEvent && taskStatusUpdateEvent.status().state() == TaskState.TASK_STATE_WORKING) { return A2AStreamingClientListenerResult.stopWithResponse( "stop when status update to TASK_STATE_WORKING, and return this message to the caller"); } return A2AStreamingClientListenerResult.continueStreaming(); }) .build();