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]

Reply via email to