joeyutong commented on code in PR #1185:
URL: https://github.com/apache/flink-agents/pull/1185#discussion_r4229023696
##########
api/src/main/java/org/apache/flink/agents/api/chat/messages/ChatMessage.java:
##########
@@ -18,187 +18,139 @@
package org.apache.flink.agents.api.chat.messages;
+import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
-import java.util.ArrayList;
-import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
+import java.util.Set;
import java.util.stream.Collectors;
/**
- * Chat message class that represents all message types (user, system,
assistant, tool) with
- * different roles.
+ * A message in a conversation, consisting of a role, an ordered list of
content blocks, and
+ * optional metadata.
*
- * <p>Message content is an ordered list of typed {@link ContentBlock}s
({@link TextBlock} plus the
- * media blocks); a text-only message simply carries one {@link TextBlock}.
The string convenience
- * constructors and factories preserve the text-message experience, and {@link
#getText()} is the
- * ordered concatenation of the text blocks.
+ * <p>The role identifies the source or purpose of the message. System
messages contain text
+ * instructions, user messages can combine text and media, and assistant
messages can include text,
+ * reasoning, and tool calls. A tool message contains exactly one {@link
ToolResultBlock}, which
+ * associates the result with its tool call.
*
- * <p>Blocks are immutable and the message snapshots every block list it is
handed, so {@link
- * #getBlocks()} is an unmodifiable view: the content changes only by
replacing it through {@link
- * #setBlocks(List)} or {@link #setText(String)}.
+ * <p>Content blocks preserve their order within the message, allowing
different kinds of content to
+ * be represented together. {@link #getText()} concatenates only the top-level
text blocks;
+ * reasoning and content nested inside tool results are excluded. {@link
#getToolCalls()} returns
+ * the tool calls in block order. Tool call IDs must be unique within a
message.
+ *
+ * <p>Metadata holds additional message-level information, such as
provider-specific attributes. It
+ * is separate from the content blocks and is not included by {@link
#getText()}.
*/
-public class ChatMessage {
+public final class ChatMessage {
+ private static final String ROLE_FIELD = "role";
+ private static final String BLOCKS_FIELD = "blocks";
+ private static final String METADATA_FIELD = "metadata";
private static final ObjectMapper MAPPER = new ObjectMapper();
+ private final MessageRole role;
+ private final List<ContentBlock> blocks;
+ private final Map<String, Object> metadata;
- private MessageRole role;
- private List<ContentBlock> blocks;
-
- @JsonProperty("tool_calls")
- private List<Map<String, Object>> toolCalls;
-
- @JsonProperty("extra_args")
- private Map<String, Object> extraArgs;
-
- /** Default constructor with SYSTEM role */
- public ChatMessage() {
- this(MessageRole.SYSTEM, Collections.emptyList(), null, null);
- }
-
- /** Constructor with role and text content */
- public ChatMessage(MessageRole role, String text) {
- this(role, blocksOf(text), null, null);
- }
-
- /** Constructor with role and content blocks */
- public ChatMessage(MessageRole role, List<ContentBlock> blocks) {
- this(role, blocks, null, null);
- }
-
- public ChatMessage(MessageRole role, String text, Map<String, Object>
extraArgs) {
- this(role, blocksOf(text), null, extraArgs);
- }
-
- public ChatMessage(MessageRole role, String text, List<Map<String,
Object>> toolCalls) {
- this(role, blocksOf(text), toolCalls, null);
- }
-
+ @JsonCreator
public ChatMessage(
- MessageRole role,
- String text,
- List<Map<String, Object>> toolCalls,
- Map<String, Object> extraArgs) {
- this(role, blocksOf(text), toolCalls, extraArgs);
+ @JsonProperty(ROLE_FIELD) MessageRole role,
+ @JsonProperty(BLOCKS_FIELD) List<ContentBlock> blocks,
+ @JsonProperty(METADATA_FIELD) Map<String, Object> metadata) {
+ this.role = Objects.requireNonNull(role, ROLE_FIELD);
+ this.blocks = List.copyOf(Objects.requireNonNull(blocks,
BLOCKS_FIELD));
+ this.metadata = metadata == null ? new HashMap<>() : metadata;
+ validate();
}
- /** Full constructor */
- public ChatMessage(
- MessageRole role,
- List<ContentBlock> blocks,
- List<Map<String, Object>> toolCalls,
- Map<String, Object> extraArgs) {
- this.role = role != null ? role : MessageRole.SYSTEM;
- this.blocks = snapshotOf(blocks);
- this.toolCalls = toolCalls != null ? toolCalls : new ArrayList<>();
- this.extraArgs = extraArgs != null ? new HashMap<>(extraArgs) : new
HashMap<>();
+ public ChatMessage(MessageRole role, List<ContentBlock> blocks) {
+ this(role, blocks, null);
}
- /** An empty or null text becomes an empty block list rather than an empty
text block. */
- private static List<ContentBlock> blocksOf(String text) {
- return text == null || text.isEmpty()
- ? Collections.emptyList()
- : Collections.singletonList(new TextBlock(text));
+ public ChatMessage(MessageRole role, String text) {
+ this(role, text == null || text.isEmpty() ? List.of() : List.of(new
TextBlock(text)));
}
- /** An unmodifiable copy — since blocks are immutable, this freezes the
content. */
- private static List<ContentBlock> snapshotOf(List<ContentBlock> blocks) {
- return List.copyOf(Objects.requireNonNull(blocks, "blocks must not be
null"));
+ private void validate() {
+ if (role == MessageRole.TOOL
+ && (blocks.size() != 1 || !(blocks.get(0) instanceof
ToolResultBlock))) {
+ throw new IllegalArgumentException(
+ "A TOOL message requires exactly one ToolResultBlock");
+ }
+ Set<String> ids = new HashSet<>();
+ for (ContentBlock b : blocks) {
+ if (role != MessageRole.TOOL && b instanceof ToolResultBlock) {
+ throw new IllegalArgumentException("ToolResultBlock requires
TOOL role");
+ }
+ if (role == MessageRole.SYSTEM && !(b instanceof TextBlock)) {
+ throw new IllegalArgumentException("SYSTEM messages accept
only text");
+ }
+ if (role == MessageRole.USER && !(b instanceof TextBlock || b
instanceof MediaBlock)) {
+ throw new IllegalArgumentException("USER messages accept only
text and media");
+ }
+ if (b instanceof ToolCallBlock && !ids.add(((ToolCallBlock)
b).getCallId())) {
+ throw new IllegalArgumentException("Duplicate tool call ID in
one message");
+ }
+ }
}
public MessageRole getRole() {
return role;
}
- public void setRole(MessageRole role) {
- this.role = role;
- }
-
- /** The content as an unmodifiable list of immutable blocks. */
public List<ContentBlock> getBlocks() {
return blocks;
}
- public void setBlocks(List<ContentBlock> blocks) {
- this.blocks = snapshotOf(blocks);
+ public Map<String, Object> getMetadata() {
+ return metadata;
}
- /** Replaces the content with a single text block (empty text clears the
content). */
@JsonIgnore
- public void setText(String text) {
- this.blocks = blocksOf(text);
- }
-
- @JsonProperty("tool_calls")
- public List<Map<String, Object>> getToolCalls() {
- return toolCalls;
- }
-
- @JsonProperty("tool_calls")
- public void setToolCalls(List<Map<String, Object>> toolCalls) {
- this.toolCalls = toolCalls;
- }
-
- @JsonProperty("extra_args")
- public Map<String, Object> getExtraArgs() {
- return extraArgs;
- }
-
- @JsonProperty("extra_args")
- public void setExtraArgs(Map<String, Object> extraArgs) {
- this.extraArgs = extraArgs != null ? extraArgs : new HashMap<>();
+ public String getText() {
+ return blocks.stream()
+ .filter(b -> b instanceof TextBlock)
Review Comment:
`LlmJudgeRoutingExecutor.renderConversation()` and `length()` still use
`ChatMessage.getText()`. With the new TOOL shape, tool results become empty and
are omitted from the judge's prompt and character budget. A new routing request
containing prior tool results therefore loses context that the judge received
before this change.
##########
python/flink_agents/api/prompts/prompt.py:
##########
@@ -88,9 +86,8 @@ def format_messages(
else:
msgs = []
for m in self.template:
- msg = ChatMessage(
- role=m.role,
- blocks=[
+ msg = m.with_blocks(
+ [
TextBlock(text=format_string(b.text, **kwargs))
if isinstance(b, TextBlock)
Review Comment:
Both Java and Python prompt formatters still process only top-level text.
For a TOOL template containing `weather: {city}`, formatting with
`city="Paris"` now produces only `tool: ` in string mode and leaves `{city}`
unchanged in message mode. Existing tool-message templates therefore lose their
content or retain unresolved placeholders.
##########
python/flink_agents/runtime/python_java_utils.py:
##########
@@ -248,6 +245,7 @@ def _encode_python_tool_result(result: Any) -> Dict[str,
Any]:
"error": result.error_message,
"execution_time_ms": result.execution_time_ms,
"tool_name": result.tool_name,
+ "blocks": result.model_dump(mode="json")["blocks"],
Review Comment:
This serializes `result` as well as `blocks`, although the bridge still
returns `result.result` unchanged. `ToolResponse.success(b"\xff")` now raises
`UnicodeDecodeError` before returning to Java; the old encoder and the
raw-bytes branch both accept it. Providing explicit text blocks does not avoid
this failure.
##########
integrations/chat-models/openai/src/main/java/org/apache/flink/agents/integrations/chatmodels/openai/OpenAICompletionsConnection.java:
##########
@@ -258,29 +260,21 @@ private ChatMessage doChat(
ChatMessage response =
OpenAIChatCompletionsUtils.convertFromOpenAIMessage(choice.message());
- // ChatCompletion.Choice#finishReason throws
OpenAIInvalidDataException when the member is
- // absent or null, so the value is read through the raw field.
- choice._finishReason()
- .asKnown()
- .ifPresent(
- reason -> response.getExtraArgs().put("finish_reason",
reason.asString()));
-
- // Stash token usage
- if (completion.usage().isPresent()) {
- String modelName = modelParams != null ? (String)
modelParams.get("model") : null;
- if (modelName == null || modelName.isBlank()) {
- modelName = this.defaultModel;
- }
- if (modelName != null && !modelName.isBlank()) {
- response.getExtraArgs().put("model_name", modelName);
- response.getExtraArgs()
- .put("promptTokens",
completion.usage().get().promptTokens());
- response.getExtraArgs()
- .put("completionTokens",
completion.usage().get().completionTokens());
- }
- }
-
- return response;
+ return new ChatResult(
+ response,
+ (modelParams != null && modelParams.get("model") != null
+ ? (String) modelParams.get("model")
+ : this.defaultModel),
Review Comment:
For `model=""` or whitespace, `buildRequest()` falls back to `defaultModel`,
but this expression keeps the blank string. The request succeeds while
`ChatResult.model` and its token metrics are attributed to a blank model (also
through vLLM).
##########
api/src/main/java/org/apache/flink/agents/api/chat/messages/ContentBlock.java:
##########
@@ -27,28 +27,34 @@
/**
* A single, typed part of a {@link ChatMessage}'s content.
*
- * <p>Blocks are ordered within a message and are immutable value objects:
every construction path,
- * including Jackson deserialization, runs the same validation, so sharing a
block instance never
- * shares mutable state. The concrete type answers how providers route the
content ({@link
+ * <p>Blocks are ordered within a message. Block fields are read-only, while
metadata and tool
+ * inputs use ordinary maps. Every construction path, including Jackson
deserialization, runs the
+ * same structural validation. The concrete type answers how providers route
the content ({@link
* TextBlock}, {@link ImageBlock}, {@link AudioBlock}, {@link VideoBlock},
{@link DocumentBlock}),
* while media encoding is carried by the media type on {@link MediaBlock}.
*
* <p>The serialized form carries a {@code type} discriminator with fixed
values ({@code text},
* {@code image}, {@code audio}, {@code video}, {@code document}) shared with
the Python API, so
* blocks cross the Java/Python boundary as plain JSON.
*/
-@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, include = JsonTypeInfo.As.PROPERTY,
property = "type")
+@JsonTypeInfo(
+ use = JsonTypeInfo.Id.NAME,
+ include = JsonTypeInfo.As.EXISTING_PROPERTY,
+ property = "type")
@JsonSubTypes({
@JsonSubTypes.Type(value = TextBlock.class, name = "text"),
@JsonSubTypes.Type(value = ImageBlock.class, name = "image"),
@JsonSubTypes.Type(value = AudioBlock.class, name = "audio"),
@JsonSubTypes.Type(value = VideoBlock.class, name = "video"),
- @JsonSubTypes.Type(value = DocumentBlock.class, name = "document")
+ @JsonSubTypes.Type(value = DocumentBlock.class, name = "document"),
+ @JsonSubTypes.Type(value = ReasoningBlock.class, name = "reasoning"),
+ @JsonSubTypes.Type(value = ToolCallBlock.class, name = "tool_call"),
+ @JsonSubTypes.Type(value = ToolResultBlock.class, name = "tool_result")
Review Comment:
Non-blocking: text/media blocks describe content formats, while reasoning
and tool calls/results describe semantic roles. Mixing these dimensions under
`ContentBlock` makes its meaning less precise; for example, `ToolResultBlock`
accepts `ContentBlock` values but must reject the reasoning and
tool-call/result subtypes.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]