Skip to content
Merged
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
17 changes: 9 additions & 8 deletions go/adk/pkg/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,19 +14,20 @@ Shared types, interfaces, and implementations for the Kagent Go ADK.
- **runner/** - Google ADK `runner.Config` creation from `AgentConfig`
- **session/** - Session management, persistence, and ADK session service adapter
- **skills/** - Agent skills discovery and shell execution
- **taskstore/** - Task storage and A2A result aggregation
- **taskstore/** - A2A task persistence through the kagent controller API
- **telemetry/** - OpenTelemetry tracing utilities

## Event Processing

The executor (`KAgentExecutor`) holds a `*runner.Runner` directly and implements `a2asrv.AgentExecutor`:
The executor (`KAgentExecutor`) is a thin kagent-specific wrapper around the
upstream `adka2a.Executor`:

```
main.go -> CreateGoogleADKRunner -> *runner.Runner
main.go -> CreateRunnerConfig -> runner.Config
|
KAgentExecutor.Execute(ctx, reqCtx, queue)
-> runner.Run(ctx, userID, sessionID, content, runConfig)
-> iterate *adksession.Event
-> ConvertADKEventToA2AEvents -> queue.Write
-> inline aggregation -> final status/artifact
KAgentExecutor.Execute(ctx, reqCtx)
-> kagent auth, telemetry, skills, session state, HITL resume setup
-> adka2a.Executor.Execute(ctx, reqCtx)
-> artifact updates for task output
-> status-only lifecycle and terminal events
```
94 changes: 13 additions & 81 deletions go/adk/pkg/a2a/converter.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,6 @@ package a2a

import (
"context"
"encoding/json"
"maps"

a2atype "github.com/a2aproject/a2a-go/v2/a2a"
"google.golang.org/adk/v2/server/adka2a/v2"
Expand All @@ -19,42 +17,6 @@ func isEmptyDataPart(part *a2atype.Part) bool {
return dp != nil && len(dp) == 0
}

// filterTextParts returns only TextParts from the given parts.
func filterTextParts(parts a2atype.ContentParts) a2atype.ContentParts {
var out a2atype.ContentParts
for _, p := range parts {
if p != nil && p.Text() != "" {
out = append(out, p)
}
}
return out
}

// messageToGenAIContent converts an A2A message to *genai.Content using kagent
// a2aPartConverter logic: handle kagent_type and adk_type DataParts explicitly,
// drop unrecognised DataParts (e.g. HITL decision parts).
func messageToGenAIContent(ctx context.Context, msg *a2atype.Message) (*genai.Content, error) {
if msg == nil {
return nil, nil
}
parts := make([]*genai.Part, 0, len(msg.Parts))
for _, part := range msg.Parts {
genaiPart, err := a2aPartConverter(ctx, msg, part)
if err != nil {
return nil, err
}
if genaiPart == nil {
continue
}
parts = append(parts, genaiPart)
}
var role genai.Role = genai.RoleUser
if msg.Role == a2atype.MessageRoleAgent {
role = genai.RoleModel
}
return genai.NewContentFromParts(parts, role), nil
}

// a2aPartConverter converts inbound A2A parts to GenAI parts.
func a2aPartConverter(_ context.Context, _ a2atype.Event, part *a2atype.Part) (*genai.Part, error) {
dp := asDataPart(part)
Expand Down Expand Up @@ -82,6 +44,19 @@ func a2aPartConverter(_ context.Context, _ a2atype.Event, part *a2atype.Part) (*
return nil, nil
}

// genAIPartConverter lets the upstream executor own artifact construction
// while preserving kagent's part filtering and long-running-tool metadata.
func genAIPartConverter(_ context.Context, event *adksession.Event, part *genai.Part) (*a2atype.Part, error) {
converted, err := adka2a.ToA2APart(part, event.LongRunningToolIDs)
if err != nil {
return nil, err
}
if isEmptyDataPart(converted) {
return nil, nil
}
return converted, nil
}

// convertDataPartToGenAI converts a DataPart with a type metadata key
// (either adk_type or kagent_type) back to GenAI for inbound message processing.
func convertDataPartToGenAI(data map[string]any, metadata map[string]any, typeKey string) (*genai.Part, error) {
Expand Down Expand Up @@ -113,46 +88,3 @@ func convertDataPartToGenAI(data map[string]any, metadata map[string]any, typeKe
}
return adka2a.ToGenAIPart(a2atype.NewDataPart(data))
}

// toA2AMetadataMap converts v to map[string]any via JSON so values placed in A2A
func toA2AMetadataMap(v any) (map[string]any, error) {
if v == nil {
return nil, nil
}
b, err := json.Marshal(v)
if err != nil {
return nil, err
}
var m map[string]any
if err := json.Unmarshal(b, &m); err != nil {
return nil, err
}
return m, nil
}

// buildEventMeta merges the base metadata with per-event fields such as
// invocation_id, author, branch, usage_metadata, etc.
func buildEventMeta(baseMeta map[string]any, adkEvent *adksession.Event) map[string]any {
result := maps.Clone(baseMeta)
if adkEvent == nil {
return result
}
for k, v := range map[string]string{
"invocation_id": adkEvent.InvocationID,
"author": adkEvent.Author,
"branch": adkEvent.Branch,
} {
if v != "" {
result[adka2a.ToA2AMetaKey(k)] = v
}
}
if adkEvent.UsageMetadata != nil {
if um, err := toA2AMetadataMap(adkEvent.UsageMetadata); err == nil && um != nil {
result[adka2a.ToA2AMetaKey("usage_metadata")] = um
}
}
if adkEvent.ErrorCode != "" {
result[adka2a.ToA2AMetaKey("error_code")] = adkEvent.ErrorCode
}
return result
}
103 changes: 29 additions & 74 deletions go/adk/pkg/a2a/converter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (

a2atype "github.com/a2aproject/a2a-go/v2/a2a"
"google.golang.org/adk/v2/server/adka2a/v2"
adksession "google.golang.org/adk/v2/session"
"google.golang.org/genai"
)

Expand Down Expand Up @@ -118,48 +119,35 @@ func TestConvertDataPartToGenAI_UnknownType(t *testing.T) {
}

// ---------------------------------------------------------------------------
// messageToGenAIContent
// a2aPartConverter
// ---------------------------------------------------------------------------

func TestMessageToGenAIContent_TextPart(t *testing.T) {
msg := a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart("hello"))
content, err := messageToGenAIContent(context.Background(), msg)
func TestA2APartConverter_TextPart(t *testing.T) {
part, err := a2aPartConverter(context.Background(), nil, a2atype.NewTextPart("hello"))
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if content == nil {
t.Fatal("expected non-nil content")
return
}
if len(content.Parts) != 1 {
t.Fatalf("expected 1 part, got %d", len(content.Parts))
}
if content.Parts[0].Text != "hello" {
t.Errorf("text = %q, want %q", content.Parts[0].Text, "hello")
if part == nil || part.Text != "hello" {
t.Fatalf("converted part = %#v, want text hello", part)
}
}

func TestMessageToGenAIContent_DropsUnrecognisedDataPart(t *testing.T) {
func TestA2APartConverter_DropsUnrecognisedDataPart(t *testing.T) {
// A DataPart with no recognised kagent_type metadata (e.g. a HITL decision
// payload like {decision_type: "approve"}) should be dropped silently.
msg := a2atype.NewMessage(a2atype.MessageRoleUser,
a2atype.NewTextPart("approving"),
part, err := a2aPartConverter(
context.Background(), nil,
convDataPart(map[string]any{"decision_type": "approve"}, nil),
)
content, err := messageToGenAIContent(context.Background(), msg)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
// Only the TextPart should survive; the unrecognised DataPart is dropped.
if len(content.Parts) != 1 {
t.Fatalf("expected 1 part (DataPart dropped), got %d", len(content.Parts))
}
if content.Parts[0].Text != "approving" {
t.Errorf("remaining part text = %q, want %q", content.Parts[0].Text, "approving")
if part != nil {
t.Fatalf("converted part = %#v, want nil", part)
}
}

func TestMessageToGenAIContent_KagentTypeFunctionResponse(t *testing.T) {
func TestA2APartConverter_KagentTypeFunctionResponse(t *testing.T) {
// A DataPart with kagent_type=function_response should be converted to GenAI.
dp := convDataPart(map[string]any{
"name": "my_func",
Expand All @@ -168,66 +156,33 @@ func TestMessageToGenAIContent_KagentTypeFunctionResponse(t *testing.T) {
}, map[string]any{
GetKAgentMetadataKey(A2ADataPartMetadataTypeKey): A2ADataPartMetadataTypeFunctionResponse,
})
msg := a2atype.NewMessage(a2atype.MessageRoleUser, dp)
content, err := messageToGenAIContent(context.Background(), msg)
part, err := a2aPartConverter(context.Background(), nil, dp)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if len(content.Parts) != 1 {
t.Fatalf("expected 1 part, got %d", len(content.Parts))
}
if content.Parts[0].FunctionResponse == nil {
if part == nil || part.FunctionResponse == nil {
t.Fatal("expected FunctionResponse, got nil")
}
if content.Parts[0].FunctionResponse.Name != "my_func" {
t.Errorf("name = %q, want my_func", content.Parts[0].FunctionResponse.Name)
}
}

func TestMessageToGenAIContent_NilMessage(t *testing.T) {
content, err := messageToGenAIContent(context.Background(), nil)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if content != nil {
t.Errorf("expected nil content for nil message, got %v", content)
if part.FunctionResponse.Name != "my_func" {
t.Errorf("name = %q, want my_func", part.FunctionResponse.Name)
}
}

// ---------------------------------------------------------------------------
// toA2AMetadataMap
// ---------------------------------------------------------------------------

func TestToA2AMetadataMap(t *testing.T) {
t.Parallel()
um := &genai.GenerateContentResponseUsageMetadata{
PromptTokenCount: 10,
CandidatesTokenCount: 20,
}
m, err := toA2AMetadataMap(um)
func TestGenAIPartConverter_PreservesLongRunningMetadata(t *testing.T) {
call := genai.NewPartFromFunctionCall("dangerous_tool", map[string]any{"path": "/tmp/x"})
call.FunctionCall.ID = "call-1"
part, err := genAIPartConverter(
context.Background(),
&adksession.Event{LongRunningToolIDs: []string{"call-1"}},
call,
)
if err != nil {
t.Fatalf("toA2AMetadataMap: %v", err)
}
if m == nil {
t.Fatal("expected non-nil map")
t.Fatalf("genAIPartConverter() error = %v", err)
}
pt, ok := m["promptTokenCount"].(float64)
if !ok || pt != 10 {
t.Fatalf("promptTokenCount: got %v (%T), want float64 10", m["promptTokenCount"], m["promptTokenCount"])
}
ct, ok := m["candidatesTokenCount"].(float64)
if !ok || ct != 20 {
t.Fatalf("candidatesTokenCount: got %v (%T), want float64 20", m["candidatesTokenCount"], m["candidatesTokenCount"])
}
}

func TestToA2AMetadataMap_nil(t *testing.T) {
t.Parallel()
m, err := toA2AMetadataMap(nil)
if err != nil {
t.Fatalf("toA2AMetadataMap(nil): %v", err)
if part == nil {
t.Fatal("genAIPartConverter() returned nil")
}
if m != nil {
t.Fatalf("expected nil map, got %#v", m)
if got, _ := ReadMetadataValue(part.Metadata, A2ADataPartMetadataIsLongRunningKey); got != true {
t.Fatalf("long-running metadata = %#v, want true", got)
}
}
Loading
Loading