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
13 changes: 13 additions & 0 deletions agent/a2aagent/a2a_agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import (
itrace "trpc.group/trpc-go/trpc-agent-go/internal/trace"
"trpc.group/trpc-go/trpc-agent-go/log"
"trpc.group/trpc-go/trpc-agent-go/model"
"trpc.group/trpc-go/trpc-agent-go/platform"
semconvtrace "trpc.group/trpc-go/trpc-agent-go/telemetry/semconv/trace"
"trpc.group/trpc-go/trpc-agent-go/tool"
)
Expand Down Expand Up @@ -155,6 +156,7 @@ func (r *A2AAgent) sendErrorEvent(
err error,
) *model.ResponseError {
respErr := model.ResponseErrorFromError(err, model.ErrorTypeRunError)
redactAgentResponseError(respErr)
agent.EmitEvent(ctx, invocation, eventChan, event.New(
invocation.InvocationID,
r.name,
Expand All @@ -166,6 +168,17 @@ func (r *A2AAgent) sendErrorEvent(
return respErr
}

func redactAgentResponseError(respErr *model.ResponseError) {
if respErr == nil {
return
}
redactor, err := platform.NewRedactor()
if err != nil {
return
}
respErr.Message = redactor.Redact(respErr.Message)
}

// validateA2ARequestOptions validates that all A2A request options are of the correct type
func (r *A2AAgent) validateA2ARequestOptions(invocation *agent.Invocation) error {
if invocation.RunOptions.A2ARequestOptions == nil {
Expand Down
27 changes: 27 additions & 0 deletions agent/a2aagent/a2a_agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2914,6 +2914,33 @@ func TestA2AAgent_sendErrorEvent_UsesRunErrorType(t *testing.T) {
require.Equal(t, evt.Response.Error, respErr)
}

func TestA2AAgent_sendErrorEvent_RedactsSensitiveMessage(t *testing.T) {
a := &A2AAgent{name: "remote-agent"}
eventCh := make(chan *event.Event, 1)
invocation := &agent.Invocation{InvocationID: "inv-test"}

respErr := a.sendErrorEvent(
context.Background(),
eventCh,
invocation,
fmt.Errorf("request failed Authorization: Bearer raw-token\napi_key=sk-1234567890abcdef\nCookie: session=abc; sid=def"),
)

require.NotNil(t, respErr)
require.Equal(t, model.ErrorTypeRunError, respErr.Type)
for _, secret := range []string{"raw-token", "sk-1234567890abcdef", "session=abc", "sid=def"} {
require.NotContains(t, respErr.Message, secret)
}
for _, redacted := range []string{"Authorization: ****", "api_key=****", "Cookie: ****"} {
require.Contains(t, respErr.Message, redacted)
}

evt := <-eventCh
require.NotNil(t, evt)
require.NotNil(t, evt.Response)
require.Equal(t, respErr, evt.Response.Error)
}

func TestA2AAgent_aggregateEventContent_IgnoresErrorResponses(t *testing.T) {
a := &A2AAgent{name: "remote-agent"}
builder := &strings.Builder{}
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