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
2 changes: 2 additions & 0 deletions .github/workflows/prc.yml
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ on:
pull_request:
merge_group:
push:
branches:
- main
workflow_dispatch:
permissions:
contents: read
Expand Down
11 changes: 10 additions & 1 deletion agent/dify/dify_agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import (
"trpc.group/trpc-go/trpc-agent-go/agent"
"trpc.group/trpc-go/trpc-agent-go/event"
"trpc.group/trpc-go/trpc-agent-go/model"
"trpc.group/trpc-go/trpc-agent-go/platform"
"trpc.group/trpc-go/trpc-agent-go/tool"
)

Expand Down Expand Up @@ -97,12 +98,20 @@ func (r *DifyAgent) sendErrorEvent(ctx context.Context, eventChan chan<- *event.
r.name,
event.WithResponse(&model.Response{
Error: &model.ResponseError{
Message: errorMessage,
Message: redactAgentErrorMessage(errorMessage),
},
}),
))
}

func redactAgentErrorMessage(errorMessage string) string {
redactor, err := platform.NewRedactor()
if err != nil {
return errorMessage
}
return redactor.Redact(errorMessage)
}

// Run implements the Agent interface
func (r *DifyAgent) Run(ctx context.Context, invocation *agent.Invocation) (<-chan *event.Event, error) {
cli, err := r.getDifyClient(invocation)
Expand Down
99 changes: 93 additions & 6 deletions agent/dify/dify_agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"fmt"
"io"
"net/http"
"strings"
"testing"

"github.com/cloudernative/dify-sdk-go"
Expand Down Expand Up @@ -535,7 +536,12 @@ func TestDifyAgent_SendErrorEvent(t *testing.T) {
}

eventChan := make(chan *event.Event, 1)
difyAgent.sendErrorEvent(context.Background(), eventChan, invocation, "test error message")
difyAgent.sendErrorEvent(
context.Background(),
eventChan,
invocation,
"request failed Authorization: Bearer raw-token\napi_key=sk-1234567890abcdef\nCookie: session=abc; sid=def",
)
close(eventChan)

evt := <-eventChan
Expand All @@ -548,8 +554,16 @@ func TestDifyAgent_SendErrorEvent(t *testing.T) {
if evt.Response.Error == nil {
t.Fatal("expected error in response")
}
if evt.Response.Error.Message != "test error message" {
t.Errorf("expected error message 'test error message', got: %s", evt.Response.Error.Message)
message := evt.Response.Error.Message
for _, secret := range []string{"raw-token", "sk-1234567890abcdef", "session=abc", "sid=def"} {
if strings.Contains(message, secret) {
t.Errorf("expected error message to redact %q, got: %s", secret, message)
}
}
for _, redacted := range []string{"Authorization: ****", "api_key=****", "Cookie: ****"} {
if !strings.Contains(message, redacted) {
t.Errorf("expected error message to contain %q, got: %s", redacted, message)
}
}
if evt.Author != "test-agent" {
t.Errorf("expected author 'test-agent', got: %s", evt.Author)
Expand All @@ -559,6 +573,66 @@ func TestDifyAgent_SendErrorEvent(t *testing.T) {
}
}

func TestDifyAgent_SendErrorEvent_RedactsSecretFields(t *testing.T) {
tests := []struct {
name string
message string
secret string
redacted string
shouldKeep string
}{
{
name: "token",
message: "provider failed token=raw-token",
secret: "raw-token",
redacted: "token=****",
},
{
name: "secret",
message: "provider failed secret=sk-secret-value",
secret: "sk-secret-value",
redacted: "secret=****",
},
{
name: "password",
message: "provider failed password=plain",
secret: "plain",
redacted: "password=****",
},
{
name: "benign token text",
message: "token limit exceeded",
shouldKeep: "token limit exceeded",
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
difyAgent := &DifyAgent{name: "test-agent"}
invocation := &agent.Invocation{InvocationID: "test-inv"}
eventChan := make(chan *event.Event, 1)

difyAgent.sendErrorEvent(context.Background(), eventChan, invocation, tt.message)
close(eventChan)

evt := <-eventChan
if evt == nil || evt.Response == nil || evt.Response.Error == nil {
t.Fatal("expected response error")
}
message := evt.Response.Error.Message
if tt.secret != "" && strings.Contains(message, tt.secret) {
t.Errorf("expected error message to redact %q, got: %s", tt.secret, message)
}
if tt.redacted != "" && !strings.Contains(message, tt.redacted) {
t.Errorf("expected error message to contain %q, got: %s", tt.redacted, message)
}
if tt.shouldKeep != "" && message != tt.shouldKeep {
t.Errorf("expected error message %q, got: %s", tt.shouldKeep, message)
}
})
}
}

// Test new extracted functions for streaming
func TestDifyAgent_ProcessStreamEvent(t *testing.T) {
t.Run("with default handler", func(t *testing.T) {
Expand Down Expand Up @@ -1129,7 +1203,12 @@ func TestDifyAgent_HelperFunctions(t *testing.T) {

eventChan := make(chan *event.Event, 1)

difyAgent.sendErrorEvent(context.Background(), eventChan, invocation, "test error message")
difyAgent.sendErrorEvent(
context.Background(),
eventChan,
invocation,
"request failed Authorization: Bearer raw-token\napi_key=sk-1234567890abcdef\nCookie: session=abc; sid=def",
)
close(eventChan)

evt := <-eventChan
Expand All @@ -1139,8 +1218,16 @@ func TestDifyAgent_HelperFunctions(t *testing.T) {
if evt.Response == nil || evt.Response.Error == nil {
t.Error("Expected error in response")
}
if evt.Response.Error.Message != "test error message" {
t.Errorf("Expected error message 'test error message', got: %s", evt.Response.Error.Message)
message := evt.Response.Error.Message
for _, secret := range []string{"raw-token", "sk-1234567890abcdef", "session=abc", "sid=def"} {
if strings.Contains(message, secret) {
t.Errorf("expected error message to redact %q, got: %s", secret, message)
}
}
for _, redacted := range []string{"Authorization: ****", "api_key=****", "Cookie: ****"} {
if !strings.Contains(message, redacted) {
t.Errorf("expected error message to contain %q, got: %s", redacted, message)
}
}
})
}
Expand Down
21 changes: 19 additions & 2 deletions internal/flow/processor/content.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@ import (
"trpc.group/trpc-go/trpc-agent-go/graph"
"trpc.group/trpc-go/trpc-agent-go/internal/fileref"
iflow "trpc.group/trpc-go/trpc-agent-go/internal/flow"
itelemetry "trpc.group/trpc-go/trpc-agent-go/internal/telemetry"
itrace "trpc.group/trpc-go/trpc-agent-go/internal/trace"
"trpc.group/trpc-go/trpc-agent-go/internal/util/message"
"trpc.group/trpc-go/trpc-agent-go/log"
"trpc.group/trpc-go/trpc-agent-go/memory"
Expand Down Expand Up @@ -3107,8 +3109,23 @@ func (p *ContentRequestProcessor) getAdaptivePreloadMemoryMessage(
Deduplicate: true,
HybridSearch: true,
}
memories, err := reader.SearchMemories(
ctx,
searchCtx, span, startedSpan := itrace.StartSpan(ctx, inv, itelemetry.NewMemorySearchSpanName())
var memories []*memory.Entry
if startedSpan {
defer func() {
itelemetry.TraceMemorySearch(
span,
searchOpts.MaxResults,
len(memories),
searchOpts.HybridSearch,
searchOpts.Deduplicate,
err,
)
span.End()
}()
}
memories, err = reader.SearchMemories(
searchCtx,
userKey,
query,
memory.WithSearchOptions(searchOpts),
Expand Down
97 changes: 97 additions & 0 deletions internal/flow/processor/content_memory_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,16 @@ import (

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.opentelemetry.io/otel/codes"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
"trpc.group/trpc-go/trpc-agent-go/agent"
"trpc.group/trpc-go/trpc-agent-go/event"
itelemetry "trpc.group/trpc-go/trpc-agent-go/internal/telemetry"
"trpc.group/trpc-go/trpc-agent-go/memory"
"trpc.group/trpc-go/trpc-agent-go/model"
"trpc.group/trpc-go/trpc-agent-go/session"
"trpc.group/trpc-go/trpc-agent-go/session/inmemory"
semconvtrace "trpc.group/trpc-go/trpc-agent-go/telemetry/semconv/trace"
"trpc.group/trpc-go/trpc-agent-go/tool"
)

Expand Down Expand Up @@ -371,6 +375,49 @@ func (m *mockMemoryService) Close() error {
return nil
}

func requireMemorySearchSpan(t *testing.T, recorder interface {
Ended() []sdktrace.ReadOnlySpan
}) sdktrace.ReadOnlySpan {
t.Helper()
for _, span := range recorder.Ended() {
if span.Name() == itelemetry.NewMemorySearchSpanName() {
return span
}
}
t.Fatalf("span %q not recorded; ended spans=%v", itelemetry.NewMemorySearchSpanName(), recorder.Ended())
return nil
}

func requireMemorySearchSpanAttribute(t *testing.T, span sdktrace.ReadOnlySpan, key string, want any) {
t.Helper()
for _, attr := range span.Attributes() {
if string(attr.Key) != key {
continue
}
switch v := want.(type) {
case string:
require.Equal(t, v, attr.Value.AsString())
case int64:
require.Equal(t, v, attr.Value.AsInt64())
case bool:
require.Equal(t, v, attr.Value.AsBool())
default:
t.Fatalf("unsupported expected attribute type %T", want)
}
return
}
t.Fatalf("missing attribute %s=%v; attributes=%v", key, want, span.Attributes())
}

func requireNoMemorySearchSpanAttribute(t *testing.T, span sdktrace.ReadOnlySpan, key string) {
t.Helper()
for _, attr := range span.Attributes() {
if string(attr.Key) == key {
t.Fatalf("unexpected attribute %s present; attributes=%v", key, span.Attributes())
}
}
}

type mockSearchableSessionService struct {
session.Service
searchResults []session.EventSearchResult
Expand Down Expand Up @@ -666,6 +713,34 @@ func TestGetPreloadMemoryMessage(t *testing.T) {
assert.Contains(t, msg.Content, "Relevant memory")
})

t.Run("positive preload records memory search trace contract", func(t *testing.T) {
recorder := useSpanRecorder(t)
p := NewContentRequestProcessor(WithPreloadMemory(2))
mockSvc := &mockMemoryService{
memories: []*memory.Entry{
newTestMemoryEntry("mem-1", "first"),
newTestMemoryEntry("mem-2", "second"),
newTestMemoryEntry("mem-3", "third"),
},
searchResults: []*memory.Entry{
newTestMemoryEntry("mem-search", "Relevant memory"),
},
}
inv := newTestInvocation(model.NewUserMessage("find relevant"), mockSvc)

msg := p.getPreloadMemoryMessage(context.Background(), inv)

require.NotNil(t, msg)
span := requireMemorySearchSpan(t, recorder)
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoTraceSpan, itelemetry.OperationMemorySearch)
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoMemorySearchMaxResults, int64(2))
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoMemorySearchResultCount, int64(1))
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoMemorySearchHybrid, true)
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoMemorySearchDeduplicate, true)
requireNoMemorySearchSpanAttribute(t, span, "trpc.go.agent.memory.search.query")
require.NotEqual(t, codes.Error, span.Status().Code)
})

t.Run("positive preload falls back to recent load when query is empty", func(t *testing.T) {
p := NewContentRequestProcessor(WithPreloadMemory(2))
mockSvc := &mockMemoryService{
Expand Down Expand Up @@ -705,6 +780,28 @@ func TestGetPreloadMemoryMessage(t *testing.T) {
assert.NotContains(t, msg.Content, "third")
})

t.Run("positive preload records memory search trace contract on search error", func(t *testing.T) {
recorder := useSpanRecorder(t)
p := NewContentRequestProcessor(WithPreloadMemory(2))
mockSvc := &mockMemoryService{
memories: []*memory.Entry{
newTestMemoryEntry("mem-1", "first"),
newTestMemoryEntry("mem-2", "second"),
newTestMemoryEntry("mem-3", "third"),
},
searchErr: assert.AnError,
}
inv := newTestInvocation(model.NewUserMessage("hello"), mockSvc)

msg := p.getPreloadMemoryMessage(context.Background(), inv)

require.NotNil(t, msg)
span := requireMemorySearchSpan(t, recorder)
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoTraceSpan, itelemetry.OperationMemorySearch)
requireMemorySearchSpanAttribute(t, span, semconvtrace.KeyTRPCAgentGoMemorySearchResultCount, int64(0))
require.Equal(t, codes.Error, span.Status().Code)
})

t.Run("positive preload falls back to recent load when search is empty", func(t *testing.T) {
p := NewContentRequestProcessor(WithPreloadMemory(2))
mockSvc := &mockMemoryService{
Expand Down
2 changes: 2 additions & 0 deletions internal/flow/processor/functioncall.go
Original file line number Diff line number Diff line change
Expand Up @@ -797,6 +797,7 @@ func (p *FunctionCallResponseProcessor) executeSingleToolCallSequentialResult(
) (toolResult, error) {
ctx, span, startedSpan := itrace.StartSpan(ctx, invocation, itelemetry.NewExecuteToolSpanName(toolCall.Function.Name))
if startedSpan {
itelemetry.MarkToolCallSpan(span)
defer span.End()
}
startTime := time.Now()
Expand Down Expand Up @@ -1101,6 +1102,7 @@ func (p *FunctionCallResponseProcessor) runParallelToolCall(
// Trace the tool execution for observability.
ctx, span, startedSpan := itrace.StartSpan(ctx, invocation, itelemetry.NewExecuteToolSpanName(tc.Function.Name))
if startedSpan {
itelemetry.MarkToolCallSpan(span)
defer span.End()
}
startTime := time.Now()
Expand Down
Loading
Loading