Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions components/model/agenticark/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -432,6 +432,24 @@ func main() {
}
```

#### Multimodal Function Tool Results

Function tool results may contain text and images when the selected Ark endpoint supports vision input:

```go
toolResultMsg := &schema.AgenticMessage{
Role: schema.AgenticRoleTypeUser,
ContentBlocks: []*schema.ContentBlock{schema.NewContentBlock(&schema.FunctionToolResult{
CallID: toolCall.CallID,
Name: toolCall.Name,
Content: []*schema.FunctionToolResultContentBlock{
{Type: schema.FunctionToolResultContentBlockTypeText, Text: &schema.UserInputText{Text: "reference image"}},
{Type: schema.FunctionToolResultContentBlockTypeImage, Image: &schema.UserInputImage{URL: "https://example.com/image.png", Detail: schema.ImageURLDetailHigh}},
},
})},
}
```


#### Server Tool Example

Expand Down
18 changes: 18 additions & 0 deletions components/model/agenticark/README.zh_CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -431,6 +431,24 @@ func main() {
}
```

#### 多模态函数工具结果

当所选 Ark endpoint 支持视觉输入时,函数工具结果可以同时包含文本和图片:

```go
toolResultMsg := &schema.AgenticMessage{
Role: schema.AgenticRoleTypeUser,
ContentBlocks: []*schema.ContentBlock{schema.NewContentBlock(&schema.FunctionToolResult{
CallID: toolCall.CallID,
Name: toolCall.Name,
Content: []*schema.FunctionToolResultContentBlock{
{Type: schema.FunctionToolResultContentBlockTypeText, Text: &schema.UserInputText{Text: "参考图片"}},
{Type: schema.FunctionToolResultContentBlockTypeImage, Image: &schema.UserInputImage{URL: "https://example.com/image.png", Detail: schema.ImageURLDetailHigh}},
},
})},
}
```


#### 服务器工具示例

Expand Down
69 changes: 64 additions & 5 deletions components/model/agenticark/convertor.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,11 +22,12 @@ import (
"sync"

"github.com/bytedance/sonic"
"github.com/cloudwego/eino/schema"
"github.com/eino-contrib/jsonschema"
"github.com/volcengine/volcengine-go-sdk/service/arkruntime/model/responses"
"golang.org/x/sync/errgroup"
"google.golang.org/protobuf/types/known/structpb"

"github.com/cloudwego/eino/schema"
)

func toSystemRoleInputItems(msg *schema.AgenticMessage) (items []*responses.InputItem, err error) {
Expand Down Expand Up @@ -447,6 +448,9 @@ func userInputTextToInputItem(role responses.MessageRole_Enum, block *schema.Use
}

func userInputImageToInputItem(role responses.MessageRole_Enum, block *schema.UserInputImage) (inputItem *responses.InputItem, err error) {
if isFileURL(block.URL) {
return nil, fmt.Errorf("file:// URLs are not supported for image input")
}
imageURL, err := resolveURL(block.URL, block.Base64Data, block.MIMEType)
if err != nil {
return nil, err
Expand Down Expand Up @@ -494,6 +498,9 @@ func toContentItemImageDetail(detail schema.ImageURLDetail) (*responses.ContentI
}

func userInputVideoToInputItem(role responses.MessageRole_Enum, block *schema.UserInputVideo) (inputItem *responses.InputItem, err error) {
if isFileURL(block.URL) {
return nil, fmt.Errorf("file:// URLs are not supported for video input")
}
videoURL, err := resolveURL(block.URL, block.Base64Data, block.MIMEType)
if err != nil {
return nil, err
Expand Down Expand Up @@ -556,6 +563,9 @@ func userInputAudioToInputItem(role responses.MessageRole_Enum, block *schema.Us
}

func userInputFileToInputItem(role responses.MessageRole_Enum, block *schema.UserInputFile) (inputItem *responses.InputItem, err error) {
if isFileURL(block.URL) {
return nil, fmt.Errorf("file:// URLs are not supported for file input")
}
fileItem := &responses.ContentItemFile{
Type: responses.ContentItemType_input_file,
Filename: &block.Name,
Expand Down Expand Up @@ -608,18 +618,56 @@ func functionToolResultToInputItem(block *schema.FunctionToolResult) (item *resp
}

func functionToolResultContentToText(content []*schema.FunctionToolResultContentBlock) (string, error) {
if len(content) > 1 {
return "", fmt.Errorf("multiple function tool result content blocks are not supported, got %d", len(content))
if len(content) == 0 {
return "", nil
}
if len(content) == 1 && content[0] != nil && content[0].Type == schema.FunctionToolResultContentBlockTypeText {
if content[0].Text == nil {
return "", fmt.Errorf("function tool result text block is nil")
}
return escapeToolResultText(content[0].Text.Text), nil
}

items := make([]map[string]any, 0, len(content))
for _, block := range content {
if block == nil {
return "", fmt.Errorf("function tool result content block is nil")
}
switch block.Type {
case schema.FunctionToolResultContentBlockTypeText:
return block.Text.Text, nil
if block.Text == nil {
return "", fmt.Errorf("function tool result text block is nil")
}
items = append(items, map[string]any{
"type": "input_text",
"text": block.Text.Text,
})
case schema.FunctionToolResultContentBlockTypeImage:
if block.Image == nil {
return "", fmt.Errorf("function tool result image block is nil")
}
imageURL, err := resolveURL(block.Image.URL, block.Image.Base64Data, block.Image.MIMEType)
if err != nil {
return "", fmt.Errorf("resolve function tool result image: %w", err)
}
item := map[string]any{
"type": "input_image",
"image_url": imageURL,
}
if block.Image.Detail != "" {
item["detail"] = string(block.Image.Detail)
}
items = append(items, item)
default:
return "", fmt.Errorf("unsupported function tool result content block type: %s", block.Type)
}
}
return "", nil

b, err := sonic.Marshal(items)
if err != nil {
return "", fmt.Errorf("marshal multimodal function tool result: %w", err)
}
return multimodalToolOutputPrefix + string(b), nil
}

func assistantGenTextToInputItem(block *schema.ContentBlock) (item *responses.InputItem, err error) {
Expand Down Expand Up @@ -1878,6 +1926,17 @@ func resolveURL(url string, base64Data string, mimeType string) (real string, er
return real, nil
}

func isFileURL(raw string) bool {
return strings.HasPrefix(strings.ToLower(raw), "file://")
}

func escapeToolResultText(text string) string {
if strings.HasPrefix(text, multimodalToolOutputPrefix) || strings.HasPrefix(text, escapedToolOutputPrefix) {
return escapedToolOutputPrefix + text
}
return text
}

func ensureDataURL(base64Data, mimeType string) (string, error) {
if strings.HasPrefix(base64Data, "data:") {
return "", fmt.Errorf("base64Data field must be a raw base64 string, but got a string with prefix 'data:'")
Expand Down
28 changes: 24 additions & 4 deletions components/model/agenticark/convertor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,12 @@ import (
"testing"

"github.com/bytedance/mockey"
"github.com/cloudwego/eino/schema"
"github.com/eino-contrib/jsonschema"
"github.com/stretchr/testify/assert"
"github.com/volcengine/volcengine-go-sdk/service/arkruntime/model/responses"
"google.golang.org/protobuf/types/known/structpb"

"github.com/cloudwego/eino/schema"
)

func TestToSystemRoleInputItems(t *testing.T) {
Expand Down Expand Up @@ -69,10 +70,10 @@ func TestToAssistantRoleInputItems(t *testing.T) {
setItemID(msg.ContentBlocks[1], "id-1")
setItemStatus(msg.ContentBlocks[1], responses.ItemStatus_completed.String())
msg.ResponseMeta = &schema.AgenticResponseMeta{
Extension: &ResponseMetaExtension{},
}
Extension: &ResponseMetaExtension{},
}

items, err := toAssistantRoleInputItems(msg)
items, err := toAssistantRoleInputItems(msg)
assert.NoError(t, err)
assert.Equal(t, 2, len(items))
assert.Equal(t, responses.MessageRole_assistant, items[0].GetInputMessage().Role)
Expand Down Expand Up @@ -274,6 +275,25 @@ func TestFunctionToolResultToInputItem(t *testing.T) {
assert.Equal(t, "r1", out.Output)
}

func TestUserInputRejectsFileURL(t *testing.T) {
_, err := userInputImageToInputItem(responses.MessageRole_user, &schema.UserInputImage{
URL: "file:///tmp/image.png",
Detail: schema.ImageURLDetailAuto,
})
assert.ErrorContains(t, err, "file:// URLs are not supported")

_, err = userInputVideoToInputItem(responses.MessageRole_user, &schema.UserInputVideo{
URL: "file:///tmp/video.mp4",
})
assert.ErrorContains(t, err, "file:// URLs are not supported")

_, err = userInputFileToInputItem(responses.MessageRole_user, &schema.UserInputFile{
URL: "file:///tmp/file.txt",
Name: "file.txt",
})
assert.ErrorContains(t, err, "file:// URLs are not supported")
}

func TestAssistantGenTextToInputItem(t *testing.T) {
block := schema.NewContentBlock(&schema.AssistantGenText{
Text: "answer"},
Expand Down
49 changes: 46 additions & 3 deletions components/model/agenticark/event_convertor.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@ import (
"fmt"
"io"

"github.com/volcengine/volcengine-go-sdk/service/arkruntime/model/responses"

"github.com/cloudwego/eino/components/model"
"github.com/cloudwego/eino/schema"
"github.com/volcengine/volcengine-go-sdk/service/arkruntime/model/responses"
"github.com/volcengine/volcengine-go-sdk/service/arkruntime/utils"
)

func receivedStreamResponse(streamReader *utils.ResponsesStreamReader,
func receivedStreamResponse(streamReader responseStreamReader,
config *model.AgenticConfig, sw *schema.StreamWriter[*model.AgenticCallbackOutput]) {

receiver := newStreamReceiver()
Expand Down Expand Up @@ -110,6 +110,10 @@ func receivedStreamResponse(streamReader *utils.ResponsesStreamReader,
block := receiver.reasoningSummaryTextDeltaEventToContentBlock(ev.ReasoningText)
sender.sendBlock(block, nil)

case *responses.Event_ReasoningRawTextDelta:
block := receiver.reasoningTextDeltaEventToContentBlock(ev.ReasoningRawTextDelta)
sender.sendBlock(block, nil)

case *responses.Event_FunctionCallArguments:
block := receiver.functionCallArgumentsDeltaEventToContentBlock(ev.FunctionCallArguments)
sender.sendBlock(block, nil)
Expand Down Expand Up @@ -300,6 +304,9 @@ type streamReceiver struct {
MaxReasoningSummaryIndex map[string]int
ReasoningSummaryIndexMapper map[string]int

MaxReasoningContentIndex map[string]int
ReasoningContentIndexMapper map[string]int

MaxTextAnnotationIndex map[string]int
TextAnnotationIndexMapper map[string]int

Expand All @@ -313,6 +320,8 @@ func newStreamReceiver() *streamReceiver {
IndexMapper: map[string]int{},
MaxReasoningSummaryIndex: map[string]int{},
ReasoningSummaryIndexMapper: map[string]int{},
MaxReasoningContentIndex: map[string]int{},
ReasoningContentIndexMapper: map[string]int{},
TextAnnotationIndexMapper: map[string]int{},
MaxTextAnnotationIndex: map[string]int{},
ItemAddedEventCache: map[string]any{},
Expand Down Expand Up @@ -348,6 +357,24 @@ func (r *streamReceiver) isNewReasoningSummaryIndex(outputIdx, summaryIdx int64)
return true
}

func (r *streamReceiver) isNewReasoningContentIndex(outputIdx, contentIdx int64) bool {
maxContentIndex := -1
if idx, ok := r.MaxReasoningContentIndex[int64ToStr(outputIdx)]; ok {
maxContentIndex = idx
}

idxKey := fmt.Sprintf("%d:%d", outputIdx, contentIdx)
if _, ok := r.ReasoningContentIndexMapper[idxKey]; ok {
return false
}

maxContentIndex++
r.ReasoningContentIndexMapper[idxKey] = maxContentIndex
r.MaxReasoningContentIndex[int64ToStr(outputIdx)] = maxContentIndex

return true
}

func (r *streamReceiver) getTextAnnotationIndex(outputIdx, contentIdx, annotationIdx int64) int {
maxAnnotationIndex := -1

Expand Down Expand Up @@ -797,6 +824,22 @@ func (r *streamReceiver) reasoningSummaryTextDeltaEventToContentBlock(ev *respon
return block
}

func (r *streamReceiver) reasoningTextDeltaEventToContentBlock(ev *responses.ReasoningTextDeltaEvent) *schema.ContentBlock {
text := ev.GetDelta()
if r.isNewReasoningContentIndex(ev.OutputIndex, ev.ContentIndex) && ev.ContentIndex != 0 {
text = "\n" + text
}

meta := &schema.StreamingMeta{
Index: r.getBlockIndex(makeReasoningIndexKey(ev.OutputIndex)),
}
block := schema.NewContentBlockChunk(&schema.Reasoning{Text: text}, meta)

setItemID(block, ev.ItemId)

return block
}

func (r *streamReceiver) functionCallArgumentsDeltaEventToContentBlock(ev *responses.FunctionCallArgumentsEvent) *schema.ContentBlock {
meta := &schema.StreamingMeta{
Index: r.getBlockIndex(makeFunctionToolCallIndexKey(ev.OutputIndex)),
Expand Down
Loading
Loading