From d3b3d31756a6e256b9443d84860e933ea504b636 Mon Sep 17 00:00:00 2001 From: Yaroslav Shevchuk Date: Fri, 28 Aug 2026 06:26:03 +0000 Subject: [PATCH 1/8] allow specifying svc-params using colon separator --- internal/flagparse/svcparams.go | 20 +++++++++++++++++--- internal/flagparse/svcparams_test.go | 25 +++++++++++++++++++++++++ 2 files changed, 42 insertions(+), 3 deletions(-) diff --git a/internal/flagparse/svcparams.go b/internal/flagparse/svcparams.go index 1426445..add874c 100644 --- a/internal/flagparse/svcparams.go +++ b/internal/flagparse/svcparams.go @@ -16,7 +16,6 @@ package flagparse import ( "fmt" - "strings" "github.com/spf13/pflag" @@ -70,9 +69,9 @@ func (s *ServiceParams) Auth() string { type svcParamValue struct{ s *ServiceParams } func (v *svcParamValue) Set(kv string) error { - k, val, ok := strings.Cut(kv, "=") + k, val, ok := cutServiceParam(kv) if !ok { - return fmt.Errorf("expected key=value, got %q", kv) + return fmt.Errorf("expected key=value or key:value, got %q", kv) } if k == "" { return fmt.Errorf("empty key in %q", kv) @@ -81,6 +80,21 @@ func (v *svcParamValue) Set(kv string) error { return nil } +// cutServiceParam splits a --svc-param argument on whichever of ':' or '=' comes first. +func cutServiceParam(kv string) (key, value string, ok bool) { + sep := -1 + for i := 0; i < len(kv); i++ { + if kv[i] == ':' || kv[i] == '=' { + sep = i + break + } + } + if sep < 0 { + return "", "", false + } + return kv[:sep], kv[sep+1:], true +} + func (v *svcParamValue) String() string { return "" } func (v *svcParamValue) Type() string { return "key=value" } diff --git a/internal/flagparse/svcparams_test.go b/internal/flagparse/svcparams_test.go index b5a88ba..b5871b5 100644 --- a/internal/flagparse/svcparams_test.go +++ b/internal/flagparse/svcparams_test.go @@ -43,6 +43,31 @@ func TestServiceParamsParse(t *testing.T) { args: []string{"--svc-param", "k=a=b"}, want: a2aclient.ServiceParams{"k": {"a=b"}}, }, + { + name: "colon separator", + args: []string{"--svc-param", "x-trace:abc"}, + want: a2aclient.ServiceParams{"x-trace": {"abc"}}, + }, + { + name: "colon value may contain colon", + args: []string{"--svc-param", "redirect:http://example.com"}, + want: a2aclient.ServiceParams{"redirect": {"http://example.com"}}, + }, + { + name: "equals before colon splits on equals", + args: []string{"--svc-param", "url=http://example.com"}, + want: a2aclient.ServiceParams{"url": {"http://example.com"}}, + }, + { + name: "colon before equals splits on colon", + args: []string{"--svc-param", "x-trace:a=b"}, + want: a2aclient.ServiceParams{"x-trace": {"a=b"}}, + }, + { + name: "empty key with colon is an error", + args: []string{"--svc-param", ":value"}, + wantErr: true, + }, { name: "repeated keys append in order", args: []string{"--svc-param", "k=1", "--svc-param", "k=2"}, From 00292a4d3648a5f06c5db9d1fe5f837009ed1d0f Mon Sep 17 00:00:00 2001 From: Yaroslav Shevchuk Date: Fri, 28 Aug 2026 15:40:07 +0000 Subject: [PATCH 2/8] rename polling-interval to poll-interval, add task resume hint, implement version selector --- internal/cli/cli_test.go | 104 +++++++++++++++++++++++++++++++++++++- internal/cli/client.go | 30 +++++++---- internal/cli/root.go | 3 ++ internal/cli/send.go | 24 ++++----- internal/output/output.go | 15 +++++- 5 files changed, 152 insertions(+), 24 deletions(-) diff --git a/internal/cli/cli_test.go b/internal/cli/cli_test.go index 122b5bb..97764c7 100644 --- a/internal/cli/cli_test.go +++ b/internal/cli/cli_test.go @@ -438,7 +438,7 @@ func TestSendStreamingFallbackUsesDefaultPoller(t *testing.T) { nonStreamingURL := startTestServerWith(t, a2a.AgentCapabilities{Streaming: false}) out, err := runCMDWithConfig(t, deps{cfgLoader: clicfg.LoadEmpty}, - "send", "-a", nonStreamingURL, "-o", "json", "--stream", "stream me", "--polling-interval", "5ms") + "send", "-a", nonStreamingURL, "-o", "json", "--stream", "stream me", "--poll-interval", "5ms") if err != nil { t.Fatalf("runCMDWithConfig() error = %v", err) } @@ -457,6 +457,108 @@ func TestSendStreamingFallbackUsesDefaultPoller(t *testing.T) { } } +func TestSend_ResumeHintForInputRequiredTask(t *testing.T) { + t.Parallel() + + var taskID a2a.TaskID + server := httptest.NewServer(a2asrv.NewRESTHandler(a2asrv.NewHandler( + a2asrv.AgentExecutorFunc(func(ctx context.Context, ec *a2asrv.ExecutorContext) iter.Seq2[a2a.Event, error] { + return func(yield func(a2a.Event, error) bool) { + taskID = ec.TaskID + task := &a2a.Task{ + ID: ec.TaskID, + ContextID: ec.ContextID, + Status: a2a.TaskStatus{State: a2a.TaskStateInputRequired}, + } + yield(task, nil) + } + }), + ))) + t.Cleanup(server.Close) + + out := mustRunCMD(t, "send", "-e", server.URL, "--transport", "rest", "hello") + if !strings.Contains(out, "a2a send --task-id "+string(taskID)) { + t.Fatalf("send text output missing the resume hint:\n%s", out) + } +} + +func TestSendWithVersionSelector(t *testing.T) { + t.Parallel() + url := startTestServer(t) + legacyURL := startLegacyTestServer(t) + + testCases := []struct { + name string + connect []string + version string + wantErr bool + }{ + { + name: "new server success", + connect: []string{"-a", url}, + version: "1.0", + }, + { + name: "old server success", + connect: []string{"-a", legacyURL}, + version: "0.3", + }, + { + name: "new server direct success", + connect: []string{"-e", url, "--transport", "rest"}, + version: "1.0", + }, + { + name: "old server direct success", + connect: []string{"-e", legacyURL, "--transport", "jsonrpc"}, + version: "0.3", + }, + { + name: "new server failure", + connect: []string{"-a", url}, + version: "0.3", + wantErr: true, + }, + { + name: "new server direct failure", + connect: []string{"-e", url, "--transport", "rest"}, + version: "0.3", + wantErr: true, + }, + { + name: "old server failure", + connect: []string{"-a", legacyURL}, + version: "1.0", + wantErr: true, + }, + { + name: "old server direct failure", + connect: []string{"-e", legacyURL, "--transport", "jsonrpc"}, + version: "1.0", + wantErr: true, + }, + { + name: "unknown version failure", + connect: []string{"-e", url}, + version: "3.0", + wantErr: true, + }, + } + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + command := []string{"send", "--a2a-version", tc.version, "-o", "json", "hi"} + command = append(command, tc.connect...) + _, err := runCMD(t, command...) + if err != nil && !tc.wantErr { + t.Fatalf("send error = %v", err) + } + if err == nil && tc.wantErr { + t.Fatal("send error = nil, wanted a failure") + } + }) + } +} + func TestGetTask(t *testing.T) { t.Parallel() url := startTestServer(t) diff --git a/internal/cli/client.go b/internal/cli/client.go index 41df7b4..3d65ca8 100644 --- a/internal/cli/client.go +++ b/internal/cli/client.go @@ -67,6 +67,9 @@ func newClientFromEndpoint(ctx context.Context, cfg *globalConfig, ref string, e cfg.logf("connecting directly to %s via %s (skipping card resolution)", endpointURL, protocol) endpoint := a2a.NewAgentInterface(endpointURL, protocol) + if cfg.a2aVersion != "" { + endpoint.ProtocolVersion = a2a.ProtocolVersion(cfg.a2aVersion) + } client, err := a2aclient.NewFromEndpoints(ctx, []*a2a.AgentInterface{endpoint}, append(clientFactoryOpts(cfg), extraOpts...)...) return client, hintInsecure(err) } @@ -108,19 +111,28 @@ func hintInsecure(err error) error { } func clientFactoryOpts(cfg *globalConfig) []a2aclient.FactoryOption { - factoryOpts := []a2aclient.FactoryOption{ - a2av0.WithRESTTransport(a2av0.RESTTransportConfig{}), - a2av0.WithJSONRPCTransport(a2av0.JSONRPCTransportConfig{}), - } var grpcOpts []grpc.DialOption if cfg.insecureGRPC { grpcOpts = append(grpcOpts, grpc.WithTransportCredentials(insecure.NewCredentials())) } - factoryOpts = append(factoryOpts, - a2agrpcv0.WithGRPCTransport(grpcOpts...), - a2agrpc.WithGRPCTransport(grpcOpts...), - ) - return factoryOpts + opts := []a2aclient.FactoryOption{a2aclient.WithDefaultsDisabled()} + if cfg.a2aVersion == "" || cfg.a2aVersion == "1.0" { + opts = append( + opts, + a2aclient.WithRESTTransport(nil), + a2aclient.WithJSONRPCTransport(nil), + a2agrpc.WithGRPCTransport(grpcOpts...), + ) + } + if cfg.a2aVersion == "" || cfg.a2aVersion == "0.3" { + opts = append( + opts, + a2av0.WithRESTTransport(a2av0.RESTTransportConfig{}), + a2av0.WithJSONRPCTransport(a2av0.JSONRPCTransportConfig{}), + a2agrpcv0.WithGRPCTransport(grpcOpts...), + ) + } + return opts } func stripHTTPScheme(raw string) string { diff --git a/internal/cli/root.go b/internal/cli/root.go index 5f74d56..d99fd5f 100644 --- a/internal/cli/root.go +++ b/internal/cli/root.go @@ -50,6 +50,8 @@ type globalConfig struct { url string transports []string svcParams *flagparse.ServiceParams + bearer string + a2aVersion string tenant string timeout time.Duration verbose bool @@ -116,6 +118,7 @@ func newRootCmd(cfg *globalConfig, deps deps) *cobra.Command { pf.StringVarP(&cfg.agentCard, "agent-card", "a", "", "Agent Card reference: host/origin, full card URL, or local file path") pf.StringVarP(&cfg.url, "endpoint", "e", "", "Agent interface URL for a direct connection; skips card resolution and requires a single --transport flag") pf.StringArrayVar(&cfg.transports, "transport", nil, "Transport preference: rest, jsonrpc, grpc (repeatable, highest preference first)") + pf.StringVar(&cfg.a2aVersion, "a2a-version", "", "Controls which a2a-protocol version client will advertise to the server.") cfg.svcParams.Attach(pf) pf.StringVar(&cfg.tenant, "tenant", "", "Tenant identifier") pf.DurationVar(&cfg.timeout, "timeout", 30*time.Second, "Request timeout") diff --git a/internal/cli/send.go b/internal/cli/send.go index 2c0b557..0132a67 100644 --- a/internal/cli/send.go +++ b/internal/cli/send.go @@ -31,15 +31,15 @@ import ( ) type sendFlags struct { - stream bool - async bool - payload string - taskID string - contextID string - history int - pollingInterval time.Duration - parts flagparse.Parts - meta flagparse.Metadata + stream bool + async bool + payload string + taskID string + contextID string + history int + pollInterval time.Duration + parts flagparse.Parts + meta flagparse.Metadata } type pollerFunc func(ctx context.Context, client *a2aclient.Client, req *a2a.SendMessageRequest, interval time.Duration) iter.Seq2[a2a.Event, error] @@ -81,9 +81,9 @@ func newSendCmd(cfg *globalConfig, poller pollerFunc) *cobra.Command { return utils.UnpackCause(ctx, err) } - cfg.logf("falling back to polling (%v): %v", flags.pollingInterval, err) + cfg.logf("falling back to polling (%v): %v", flags.pollInterval, err) - for event, err := range poller(ctx, client, req, flags.pollingInterval) { + for event, err := range poller(ctx, client, req, flags.pollInterval) { debounceTimeout() if err := handleStreamEntry(cfg, event, err); err != nil { return utils.UnpackCause(ctx, err) @@ -113,7 +113,7 @@ func newSendCmd(cfg *globalConfig, poller pollerFunc) *cobra.Command { f.StringVar(&flags.taskID, "task-id", "", "Task ID to continue an existing task") f.StringVar(&flags.contextID, "context-id", "", "Context ID to group this turn under") f.IntVar(&flags.history, "history", 0, "Request n history messages in the response") - f.DurationVar(&flags.pollingInterval, "polling-interval", 5*time.Second, "Duration between GetTask requests in polling fallback mode.") + f.DurationVar(&flags.pollInterval, "poll-interval", 2*time.Second, "Duration between GetTask requests in polling fallback mode.") flags.parts.Attach(f) flags.meta.Attach(f, "metadata", "Attach request metadata as a JSON object (repeatable)") diff --git a/internal/output/output.go b/internal/output/output.go index 7a50df5..77a0e33 100644 --- a/internal/output/output.go +++ b/internal/output/output.go @@ -78,7 +78,7 @@ func (p *Printer) PrintTask(task *a2a.Task) error { if p.Mode == ModeJson { return p.PrintJSON(task) } - _, err := io.WriteString(p.Out, formatTask(task)) + _, err := io.WriteString(p.Out, formatTask(task)+formatResumeHint(task)) return err } @@ -121,7 +121,7 @@ func (p *Printer) PrintSendResult(result a2a.SendMessageResult) error { } switch r := result.(type) { case *a2a.Task: - _, err := io.WriteString(p.Out, formatTask(r)) + _, err := io.WriteString(p.Out, formatTask(r)+formatResumeHint(r)) return err case *a2a.Message: _, err := io.WriteString(p.Out, formatMessage(r)) @@ -206,6 +206,17 @@ func formatTask(task *a2a.Task) string { return sb.String() } +// formatResumeHint returns a copy-pasteable command to continue or reply to a task. +func formatResumeHint(task *a2a.Task) string { + if task.ID == "" { + return "" + } + if task.Status.State != a2a.TaskStateInputRequired && task.Status.State != a2a.TaskStateAuthRequired { + return "" + } + return fmt.Sprintf("\n\nResume: a2a send --task-id %s %q\n", task.ID, "") +} + func formatMessage(msg *a2a.Message) string { role := "user" if msg.Role == a2a.MessageRoleAgent { From 9daa3a4fd328135ddd8a59f6e64dd5170420d678 Mon Sep 17 00:00:00 2001 From: Yaroslav Shevchuk Date: Fri, 28 Aug 2026 15:53:35 +0000 Subject: [PATCH 3/8] remove strange requirement --- specification/SPEC.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/specification/SPEC.md b/specification/SPEC.md index f93bc8f..bf101b3 100644 --- a/specification/SPEC.md +++ b/specification/SPEC.md @@ -492,7 +492,7 @@ A client MAY express its own preference with `--transport`, which is **repeatabl Cross-cutting options such as `--insecure` apply to whichever transport is negotiated. Where a future option is meaningful only for one binding (for example a gRPC keepalive setting that HTTP has no analogue for), a tool SHOULD namespace it per transport rather than overloading a global flag; the reserved convention is `---