diff --git a/go/core/cli/cmd/kagent/main.go b/go/core/cli/cmd/kagent/main.go index 1e1eb4339..09c3b57c2 100644 --- a/go/core/cli/cmd/kagent/main.go +++ b/go/core/cli/cmd/kagent/main.go @@ -2,386 +2,20 @@ package main import ( "context" - "errors" "fmt" "os" "os/signal" - "strconv" "syscall" - "time" - cli "github.com/kagent-dev/kagent/go/core/cli/internal/cli/agent" - agentinstancecli "github.com/kagent-dev/kagent/go/core/cli/internal/cli/agentinstance" - agenttemplatecli "github.com/kagent-dev/kagent/go/core/cli/internal/cli/agenttemplate" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/envdoc" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/mcp" - "github.com/kagent-dev/kagent/go/core/cli/internal/profiles" - "github.com/kagent-dev/kagent/go/core/cli/internal/tui" - dbcli "github.com/kagent-dev/kagent/go/core/pkg/cli/db" - dbmigrate "github.com/kagent-dev/kagent/go/core/pkg/cli/db/migrate" - "github.com/kagent-dev/kagent/go/core/pkg/migrations" - "github.com/spf13/cobra" - "golang.org/x/term" - corev1 "k8s.io/api/core/v1" - "k8s.io/client-go/tools/clientcmd" - "sigs.k8s.io/controller-runtime/pkg/client" + "github.com/kagent-dev/kagent/go/core/cli" ) func main() { - ctx, cancel := context.WithCancel(context.Background()) + ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer cancel() - // listen for signals to cancel the context throughout the application - done := make(chan os.Signal, 1) - signal.Notify(done, os.Interrupt, syscall.SIGTERM) - - go func() { - <-done - - fmt.Fprintf(os.Stderr, "kagent aborted.\n") - fmt.Fprintf(os.Stderr, "Exiting.\n") - - cancel() - }() - rootCmd := newRootCommand(ctx, defaultRootOptions()) - if err := rootCmd.ExecuteContext(ctx); err != nil { + if err := cli.Root().ExecuteContext(ctx); err != nil { fmt.Fprintf(os.Stderr, "Error: %v\n", err) - os.Exit(1) } } - -type rootOptions struct { - Connection connection.Options - OutputFormat string -} - -func defaultRootOptions() *rootOptions { - return &rootOptions{Connection: connection.DefaultOptions(), OutputFormat: "table"} -} - -func newRootCommand(ctx context.Context, opts *rootOptions) *cobra.Command { - cfg := &opts.Connection - rootCmd := &cobra.Command{ - Use: "kagent", - Short: "kagent is a CLI for kagent", - Long: "kagent is a CLI for kagent", - SilenceErrors: true, - SilenceUsage: true, - RunE: func(cmd *cobra.Command, _ []string) error { - return runInteractive(cmd, cfg) - }, - } - rootCmd.SetContext(ctx) - - rootCmd.PersistentFlags().StringVar(&cfg.KAgentURL, "kagent-url", cfg.KAgentURL, "KAgent REST URL") - rootCmd.PersistentFlags().StringVar(&cfg.KAgentGRPCURL, "grpc-url", cfg.KAgentGRPCURL, "KAgent gRPC target") - rootCmd.PersistentFlags().BoolVar(&cfg.KAgentGRPCTLS, "grpc-tls", cfg.KAgentGRPCTLS, "Use TLS for KAgent gRPC") - rootCmd.PersistentFlags().StringVar(&cfg.KAgentGRPCCAFile, "grpc-ca-file", cfg.KAgentGRPCCAFile, "CA certificate file for KAgent gRPC") - rootCmd.PersistentFlags().StringVar(&cfg.KAgentGRPCServerName, "grpc-server-name", cfg.KAgentGRPCServerName, "TLS server name for KAgent gRPC") - rootCmd.PersistentFlags().StringVarP(&cfg.Namespace, "namespace", "n", cfg.Namespace, "Namespace") - rootCmd.PersistentFlags().StringVarP(&opts.OutputFormat, "output-format", "o", opts.OutputFormat, "Output format") - rootCmd.PersistentFlags().BoolVarP(&cfg.Verbose, "verbose", "v", cfg.Verbose, "Verbose output") - rootCmd.PersistentFlags().DurationVar(&cfg.Timeout, "timeout", cfg.Timeout, "Timeout") - rootCmd.PersistentFlags().StringVar(&cfg.UserID, "user-id", cfg.UserID, "Caller identity used to select the server-side data partition") - installCfg := &cli.InstallCfg{ - Connection: cfg, - } - - installCmd := &cobra.Command{ - Use: "install", - Short: "Install kagent", - Long: `Install kagent`, - Run: func(cmd *cobra.Command, args []string) { - cli.InstallCmd(cmd.Context(), installCfg) - }, - } - installCmd.Flags().StringVar(&installCfg.Profile, "profile", "", "Installation profile (minimal|demo)") - _ = installCmd.RegisterFlagCompletionFunc("profile", func(cmd *cobra.Command, args []string, toComplete string) ([]string, cobra.ShellCompDirective) { - return profiles.Profiles, cobra.ShellCompDirectiveNoFileComp - }) - - uninstallCmd := &cobra.Command{ - Use: "uninstall", - Short: "Uninstall kagent", - Long: `Uninstall kagent`, - Run: func(cmd *cobra.Command, args []string) { - cli.UninstallCmd(cmd.Context(), cfg.Namespace) - }, - } - - invokeCfg := &agentinstancecli.InvokeCfg{ - Connection: cfg, - } - - invokeCmd := &cobra.Command{ - Use: "invoke", - Short: "Invoke an AgentInstance", - Long: `Invoke an existing AgentInstance through the A2A API.`, - Args: cobra.NoArgs, - RunE: func(cmd *cobra.Command, args []string) error { - invokeCfg.OutputFormat = opts.OutputFormat - return agentinstancecli.InvokeCmd(cmd.Context(), invokeCfg, cmd.InOrStdin(), cmd.OutOrStdout()) - }, - Example: `kagent invoke --agent-instance 8bd650a8-9775-488f-8bc1-0d52bf7bdcab --task "Get all the pods"`, - } - - invokeCmd.Flags().StringVar(&invokeCfg.AgentInstance, "agent-instance", "", "AgentInstance ID") - invokeCmd.Flags().StringVarP(&invokeCfg.Task, "task", "t", "", "Task text") - invokeCmd.Flags().StringVarP(&invokeCfg.File, "file", "f", "", "Read task text from a file or - for stdin") - invokeCmd.Flags().BoolVarP(&invokeCfg.Stream, "stream", "S", false, "Stream the response") - invokeCmd.Flags().StringVar(&invokeCfg.Token, "token", "", "Model API key passed through as an A2A Bearer token") - _ = invokeCmd.MarkFlagRequired("agent-instance") - invokeCmd.MarkFlagsOneRequired("task", "file") - invokeCmd.MarkFlagsMutuallyExclusive("task", "file") - - bugReportCmd := &cobra.Command{ - Use: "bug-report", - Short: "Generate a bug report", - Long: `Generate a bug report`, - Run: func(cmd *cobra.Command, args []string) { - pf, err := connection.Connect(cmd.Context(), cfg) - if err != nil { - fmt.Fprintf(os.Stderr, "Error connecting to server: %v\n", err) - return - } - if pf != nil { - defer pf.Stop() - } - cli.BugReportCmd(cfg.Namespace, cfg.Verbose) - }, - } - - versionCmd := &cobra.Command{ - Use: "version", - Short: "Print the kagent version", - Long: `Print the kagent version`, - Run: func(cmd *cobra.Command, args []string) { - // print out kagent CLI version regardless if a port-forward to kagent server succeeds - // versions unable to obtain from the remote kagent will be reported as "unknown" - clientSet := cfg.Client() - defer clientSet.Close() //nolint:errcheck - defer cli.VersionCmd(clientSet) - - if pf, _ := connection.Connect(cmd.Context(), cfg); pf != nil { - defer pf.Stop() - } - }, - } - - dashboardCmd := &cobra.Command{ - Use: "dashboard", - Short: "Open the kagent dashboard", - Long: `Open the kagent dashboard`, - Run: func(cmd *cobra.Command, args []string) { - cli.DashboardCmd(cmd.Context(), cfg.Namespace) - }, - } - - getCmd := &cobra.Command{ - Use: "get", - Short: "Get a kagent resource", - Long: `Get a kagent resource`, - Args: cobra.NoArgs, - RunE: func(_ *cobra.Command, _ []string) error { - return fmt.Errorf("resource type is required") - }, - } - agentInstanceGetCfg := &agentinstancecli.GetCfg{Connection: cfg} - getAgentInstanceCmd := &cobra.Command{ - Use: "agent-instance [ID]", - Short: "Get an AgentInstance or list your AgentInstances", - Args: cobra.MaximumNArgs(1), - RunE: func(cmd *cobra.Command, args []string) error { - agentInstanceGetCfg.OutputFormat = opts.OutputFormat - agentInstanceGetCfg.InstanceID = "" - if len(args) == 1 { - agentInstanceGetCfg.InstanceID = args[0] - } - return agentinstancecli.GetCmd(cmd.Context(), agentInstanceGetCfg, cmd.OutOrStdout()) - }, - } - getAgentInstanceCmd.Flags().Int32Var(&agentInstanceGetCfg.PageSize, "page-size", 0, "Number of AgentInstances to return (default 50, maximum 100)") - getAgentInstanceCmd.Flags().StringVar(&agentInstanceGetCfg.PageToken, "page-token", "", "Token returned by the previous page") - - agentTemplateGetCfg := &agenttemplatecli.GetCfg{} - getAgentTemplateCmd := &cobra.Command{ - Use: "agent-template [NAME]", - Short: "Get an AgentTemplate or list AgentTemplates", - Args: cobra.MaximumNArgs(1), - RunE: func(cmd *cobra.Command, args []string) error { - agentTemplateGetCfg.Namespace = cfg.Namespace - agentTemplateGetCfg.OutputFormat = opts.OutputFormat - agentTemplateGetCfg.Name = "" - if len(args) == 1 { - agentTemplateGetCfg.Name = args[0] - } - return agenttemplatecli.GetCmd(cmd.Context(), agentTemplateGetCfg, cmd.OutOrStdout()) - }, - } - getAgentTemplateCmd.Flags().Int64Var(&agentTemplateGetCfg.PageSize, "page-size", 0, "Number of AgentTemplates per page (0 uses 100; maximum 100)") - getAgentTemplateCmd.Flags().StringVar(&agentTemplateGetCfg.PageToken, "page-token", "", "Token returned by the previous page") - - getCmd.AddCommand(getAgentInstanceCmd, getAgentTemplateCmd) - - createCmd := &cobra.Command{ - Use: "create", - Short: "Create a kagent resource", - Args: cobra.NoArgs, - RunE: func(_ *cobra.Command, _ []string) error { - return fmt.Errorf("resource type is required") - }, - } - createAgentInstanceCfg := &agentinstancecli.CreateCfg{Connection: cfg} - createAgentInstanceCmd := &cobra.Command{ - Use: "agent-instance", - Short: "Create an AgentInstance", - Args: cobra.NoArgs, - RunE: func(cmd *cobra.Command, _ []string) error { - createAgentInstanceCfg.OutputFormat = opts.OutputFormat - return agentinstancecli.CreateCmd(cmd.Context(), createAgentInstanceCfg, cmd.OutOrStdout()) - }, - } - createAgentInstanceCmd.Flags().StringVar(&createAgentInstanceCfg.Harness, "harness", "", "Harness name") - createAgentInstanceCmd.Flags().StringVar(&createAgentInstanceCfg.AgentTemplate, "agent-template", "", "AgentTemplate name") - createAgentInstanceCmd.Flags().StringVar(&createAgentInstanceCfg.RequestID, "request-id", "", "Idempotency key (generated when omitted)") - _ = createAgentInstanceCmd.MarkFlagRequired("harness") - _ = createAgentInstanceCmd.MarkFlagRequired("agent-template") - createCmd.AddCommand(createAgentInstanceCmd) - - deleteCmd := &cobra.Command{ - Use: "delete", - Short: "Delete a kagent resource", - Args: cobra.NoArgs, - RunE: func(_ *cobra.Command, _ []string) error { - return fmt.Errorf("resource type is required") - }, - } - deleteAgentInstanceCfg := &agentinstancecli.DeleteCfg{Connection: cfg} - deleteAgentInstanceCmd := &cobra.Command{ - Use: "agent-instance ID", - Short: "Delete an AgentInstance", - Args: cobra.ExactArgs(1), - RunE: func(cmd *cobra.Command, args []string) error { - deleteAgentInstanceCfg.OutputFormat = opts.OutputFormat - deleteAgentInstanceCfg.InstanceID = args[0] - return agentinstancecli.DeleteCmd(cmd.Context(), deleteAgentInstanceCfg, cmd.OutOrStdout()) - }, - } - deleteCmd.AddCommand(deleteAgentInstanceCmd) - - rootCmd.AddCommand(installCmd, uninstallCmd, invokeCmd, bugReportCmd, versionCmd, dashboardCmd, getCmd, createCmd, deleteCmd, mcp.NewMCPCmd(), envdoc.NewEnvCmd(), dbcli.NewCommandFromFunc(migrationSources(opts))) - - return rootCmd -} - -// vectorEnabledKey names two lookups that deliberately share it: the CLI's -// own DATABASE_VECTOR_ENABLED env var (a local operator override), and the -// controller-configmap key the chart renders — the value the controller pod -// itself consumes via envFrom. Same name, two different places. -const vectorEnabledKey = "DATABASE_VECTOR_ENABLED" - -// migrationSources resolves the built-in migration tracks when a db -// subcommand runs (never during command construction, so unrelated commands -// do no work and print no warnings). The vector track is gated, in order of -// precedence, on: the DATABASE_VECTOR_ENABLED env var in the CLI's own -// environment (explicit operator intent, works without a cluster), the -// controller's configmap on the live cluster (the same value the server -// reads), and finally the controller's default (enabled). -func migrationSources(opts *rootOptions) dbmigrate.SourcesFunc { - return func(ctx context.Context) ([]migrations.Source, error) { - vectorEnabled := true - if v := os.Getenv(vectorEnabledKey); v != "" { - b, err := strconv.ParseBool(v) - if err != nil { - fmt.Fprintf(os.Stderr, "warning: invalid %s=%q; assuming true\n", vectorEnabledKey, v) - } else { - vectorEnabled = b - } - } else if b, ok := clusterVectorEnabled(ctx, opts.Connection.Namespace); ok { - vectorEnabled = b - } - return migrations.BuiltinSources(vectorEnabled), nil - } -} - -// clusterVectorEnabled reads the vectorEnabledKey entry from the controller -// configmap in the given namespace (the same "kagent-controller" default -// naming the rest of the CLI assumes) — the cluster-side counterpart of the -// env-var override in migrationSources. When the value is used it says so on -// stderr, naming the kubeconfig context it was read from — the lookup follows -// the *current* context, so this is the operator's cue that the cluster and -// their --db-url had better be the same install. Best-effort: reports -// ok=false when no cluster is reachable, the configmap is absent, or the -// value doesn't parse — callers fall back to the default. -func clusterVectorEnabled(ctx context.Context, namespace string) (enabled, ok bool) { - restConfig, err := clientcmd.NewNonInteractiveDeferredLoadingClientConfig( - clientcmd.NewDefaultClientConfigLoadingRules(), - &clientcmd.ConfigOverrides{}, - ).ClientConfig() - if err != nil { - return false, false - } - k8sClient, err := client.New(restConfig, client.Options{}) - if err != nil { - return false, false - } - ctx, cancel := context.WithTimeout(ctx, 3*time.Second) - defer cancel() - var cm corev1.ConfigMap - if err := k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: "kagent-controller"}, &cm); err != nil { - return false, false - } - b, err := strconv.ParseBool(cm.Data[vectorEnabledKey]) - if err != nil { - return false, false - } - // Trailing blank line separates the notice from the command's stdout - // when both land on a terminal; piped stdout is unaffected. - fmt.Fprintf(os.Stderr, "resolved vector track from cluster context %q: configmap %s/kagent-controller has %s=%t (set %s to override)\n\n", - currentKubeContext(), namespace, vectorEnabledKey, b, vectorEnabledKey) - return b, true -} - -// currentKubeContext names the kubeconfig context the CLI's Kubernetes client -// dials, for operator-facing messages. Best-effort. -func currentKubeContext() string { - raw, err := clientcmd.NewDefaultClientConfigLoadingRules().Load() - if err != nil || raw.CurrentContext == "" { - return "(current kubeconfig context)" - } - return raw.CurrentContext -} - -// runInteractive launches the workspace; the TUI reads raw keys, so a redirected stream is an error. -func runInteractive(cmd *cobra.Command, cfg *connection.Options) (err error) { - if !isTerminal(cmd.InOrStdin()) || !isTerminal(cmd.OutOrStdout()) { - return errors.New("kagent requires a terminal; use `kagent get agent-instance` and `kagent invoke` for non-interactive use") - } - - client := cfg.Client() - defer func() { - err = errors.Join(err, client.Close()) - }() - - portForward, connectErr := connection.Connect(cmd.Context(), cfg) - if connectErr != nil { - return fmt.Errorf("connect to kagent: %w", connectErr) - } - if portForward != nil { - defer portForward.Stop() - } - - workspace := tui.Options{Namespace: cfg.Namespace} - if runErr := tui.RunWorkspace(cmd.Context(), workspace, client, cfg.Verbose); runErr != nil { - return fmt.Errorf("run kagent workspace: %w", runErr) - } - return nil -} - -// isTerminal reports whether a stream is backed by a TTY; a non-*os.File never is. -func isTerminal(stream any) bool { - file, ok := stream.(*os.File) - return ok && term.IsTerminal(int(file.Fd())) -} diff --git a/go/core/cli/interactive.go b/go/core/cli/interactive.go new file mode 100644 index 000000000..7e0bd45b2 --- /dev/null +++ b/go/core/cli/interactive.go @@ -0,0 +1,43 @@ +package cli + +import ( + "errors" + "fmt" + "os" + + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + "github.com/kagent-dev/kagent/go/core/cli/internal/tui" + "github.com/spf13/cobra" + "golang.org/x/term" +) + +// runInteractive launches the workspace; the TUI reads raw keys, so a redirected stream is an error. +func runInteractive(cmd *cobra.Command, _ []string) (err error) { + if !isTerminal(cmd.InOrStdin()) || !isTerminal(cmd.OutOrStdout()) { + return errors.New("kagent requires a terminal; use `kagent get agent-instance` and `kagent invoke` for non-interactive use") + } + + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + session, err := connection.Open(cmd.Context(), options) + if err != nil { + return err + } + defer func() { + err = errors.Join(err, session.Close()) + }() + + workspace := tui.Options{Namespace: session.Namespace} + if runErr := tui.RunWorkspace(cmd.Context(), workspace, session.Client, options.Verbose); runErr != nil { + return fmt.Errorf("run kagent workspace: %w", runErr) + } + return nil +} + +// isTerminal reports whether a stream is backed by a TTY; a non-*os.File never is. +func isTerminal(stream any) bool { + file, ok := stream.(*os.File) + return ok && term.IsTerminal(int(file.Fd())) +} diff --git a/go/core/cli/internal/cli/agent/version.go b/go/core/cli/internal/cli/agent/version.go deleted file mode 100644 index c9cd2adc9..000000000 --- a/go/core/cli/internal/cli/agent/version.go +++ /dev/null @@ -1,29 +0,0 @@ -package cli - -import ( - "context" - "encoding/json" - "os" - "time" - - "github.com/kagent-dev/kagent/go/api/client" - "github.com/kagent-dev/kagent/go/core/internal/version" -) - -func VersionCmd(clientSet *client.ClientSet) { - versionInfo := map[string]string{ - "kagent_version": version.Version, - "git_commit": version.GitCommit, - "build_date": version.BuildDate, - } - ctx, cancel := context.WithTimeout(context.Background(), time.Second*5) - defer cancel() - serverVersion, err := clientSet.Version.GetVersion(ctx) - if err != nil { - versionInfo["backend_version"] = "unknown" - } else { - versionInfo["backend_version"] = serverVersion.KAgentVersion - } - - json.NewEncoder(os.Stdout).Encode(versionInfo) //nolint:errcheck -} diff --git a/go/core/cli/internal/cli/agenttemplate/get.go b/go/core/cli/internal/commands/agent_template.go similarity index 54% rename from go/core/cli/internal/cli/agenttemplate/get.go rename to go/core/cli/internal/commands/agent_template.go index 9fb294f2d..2d38521b0 100644 --- a/go/core/cli/internal/cli/agenttemplate/get.go +++ b/go/core/cli/internal/commands/agent_template.go @@ -1,5 +1,4 @@ -// Package agenttemplate implements AgentTemplate CLI commands. -package agenttemplate +package commands import ( "context" @@ -12,16 +11,18 @@ import ( "github.com/jedib0t/go-pretty/v6/table" typedapiv1alpha3 "github.com/kagent-dev/kagent/go/api/clientset/versioned/typed/api/v1alpha3" apiv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" commonk8s "github.com/kagent-dev/kagent/go/core/cli/internal/common/k8s" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" + "github.com/spf13/cobra" "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) -const maxPageSize = 100 +const agentTemplateMaxPageSize = 100 -// GetCfg configures AgentTemplate get and list operations. -type GetCfg struct { +// AgentTemplateGetCfg configures AgentTemplate get and list operations. +type AgentTemplateGetCfg struct { Namespace string OutputFormat string Name string @@ -29,13 +30,13 @@ type GetCfg struct { PageToken string } -// GetCmd gets one AgentTemplate or lists AgentTemplates through Kubernetes. -func GetCmd(ctx context.Context, cfg *GetCfg, out io.Writer) error { +// runGetAgentTemplate gets one AgentTemplate or lists AgentTemplates through Kubernetes. +func runGetAgentTemplate(ctx context.Context, cfg *AgentTemplateGetCfg, out io.Writer) error { format, err := clioutput.Parse(cfg.OutputFormat) if err != nil { return err } - if err := validateGetCfg(cfg); err != nil { + if err := validateAgentTemplateGetCfg(cfg); err != nil { return err } @@ -43,12 +44,12 @@ func GetCmd(ctx context.Context, cfg *GetCfg, out io.Writer) error { if err != nil { return err } - return get(ctx, clients.ApiV1alpha3().AgentTemplates(cfg.Namespace), cfg, format, out) + return getAgentTemplates(ctx, clients.ApiV1alpha3().AgentTemplates(cfg.Namespace), cfg, format, out) } -func validateGetCfg(cfg *GetCfg) error { - if cfg.PageSize < 0 || cfg.PageSize > maxPageSize { - return fmt.Errorf("page size must be between 1 and %d, or 0 for the default of %d", maxPageSize, maxPageSize) +func validateAgentTemplateGetCfg(cfg *AgentTemplateGetCfg) error { + if cfg.PageSize < 0 || cfg.PageSize > agentTemplateMaxPageSize { + return fmt.Errorf("page size must be between 1 and %d, or 0 for the default of %d", agentTemplateMaxPageSize, agentTemplateMaxPageSize) } if cfg.Name != "" && (cfg.PageSize != 0 || cfg.PageToken != "") { return errors.New("pagination flags cannot be used when getting one AgentTemplate") @@ -56,10 +57,10 @@ func validateGetCfg(cfg *GetCfg) error { return nil } -func get( +func getAgentTemplates( ctx context.Context, client typedapiv1alpha3.AgentTemplateInterface, - cfg *GetCfg, + cfg *AgentTemplateGetCfg, format clioutput.Format, out io.Writer, ) error { @@ -71,12 +72,12 @@ func get( if format == clioutput.FormatJSON { return clioutput.WriteJSON(out, template) } - return writeTemplatesTable(out, []apiv1alpha3.AgentTemplate{*template}, false, "") + return writeAgentTemplatesTable(out, []apiv1alpha3.AgentTemplate{*template}, false, "") } pageSize := cfg.PageSize if pageSize == 0 { - pageSize = maxPageSize + pageSize = agentTemplateMaxPageSize } templates, err := client.List(ctx, metav1.ListOptions{Limit: pageSize, Continue: cfg.PageToken}) if err != nil { @@ -85,10 +86,10 @@ func get( if format == clioutput.FormatJSON { return clioutput.WriteJSON(out, templates) } - return writeTemplatesTable(out, templates.Items, true, templates.Continue) + return writeAgentTemplatesTable(out, templates.Items, true, templates.Continue) } -func writeTemplatesTable(w io.Writer, templates []apiv1alpha3.AgentTemplate, list bool, nextPageToken string) error { +func writeAgentTemplatesTable(w io.Writer, templates []apiv1alpha3.AgentTemplate, list bool, nextPageToken string) error { tw := table.NewWriter() tw.AppendHeader(table.Row{"NAME", "HARNESS", "READY", "CREATED"}) for i := range templates { @@ -122,3 +123,34 @@ func writeTemplatesTable(w io.Writer, templates []apiv1alpha3.AgentTemplate, lis } return nil } + +// NewGetAgentTemplateCmd constructs the AgentTemplate get/list command. +func NewGetAgentTemplateCmd() *cobra.Command { + cfg := &AgentTemplateGetCfg{} + cmd := &cobra.Command{ + Use: "agent-template [NAME]", + Short: "Get an AgentTemplate or list AgentTemplates", + Args: cobra.MaximumNArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + format, err := clioutput.FromCommand(cmd) + if err != nil { + return err + } + var name string + if len(args) == 1 { + name = args[0] + } + cfg.Namespace = options.Namespace + cfg.OutputFormat = format + cfg.Name = name + return runGetAgentTemplate(cmd.Context(), cfg, cmd.OutOrStdout()) + }, + } + cmd.Flags().Int64Var(&cfg.PageSize, "page-size", 0, "Number of AgentTemplates per page (0 uses 100; maximum 100)") + cmd.Flags().StringVar(&cfg.PageToken, "page-token", "", "Token returned by the previous page") + return cmd +} diff --git a/go/core/cli/internal/cli/agenttemplate/get_test.go b/go/core/cli/internal/commands/agent_template_test.go similarity index 79% rename from go/core/cli/internal/cli/agenttemplate/get_test.go rename to go/core/cli/internal/commands/agent_template_test.go index 743af6444..95a4aec65 100644 --- a/go/core/cli/internal/cli/agenttemplate/get_test.go +++ b/go/core/cli/internal/commands/agent_template_test.go @@ -1,4 +1,4 @@ -package agenttemplate +package commands import ( "bytes" @@ -8,7 +8,7 @@ import ( clientfake "github.com/kagent-dev/kagent/go/api/clientset/versioned/fake" apiv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -16,23 +16,23 @@ import ( k8stesting "k8s.io/client-go/testing" ) -func TestValidateGetCfg(t *testing.T) { +func TestValidateAgentTemplateGetCfg(t *testing.T) { tests := []struct { name string - cfg GetCfg + cfg AgentTemplateGetCfg wantErr string }{ {name: "list"}, - {name: "list page", cfg: GetCfg{PageSize: 10, PageToken: "next"}}, - {name: "get", cfg: GetCfg{Name: "template"}}, - {name: "negative page size", cfg: GetCfg{PageSize: -1}, wantErr: "page size"}, - {name: "large page size", cfg: GetCfg{PageSize: 101}, wantErr: "page size"}, - {name: "get with pagination", cfg: GetCfg{Name: "template", PageSize: 10}, wantErr: "pagination"}, + {name: "list page", cfg: AgentTemplateGetCfg{PageSize: 10, PageToken: "next"}}, + {name: "get", cfg: AgentTemplateGetCfg{Name: "template"}}, + {name: "negative page size", cfg: AgentTemplateGetCfg{PageSize: -1}, wantErr: "page size"}, + {name: "large page size", cfg: AgentTemplateGetCfg{PageSize: 101}, wantErr: "page size"}, + {name: "get with pagination", cfg: AgentTemplateGetCfg{Name: "template", PageSize: 10}, wantErr: "pagination"}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - err := validateGetCfg(&tt.cfg) + err := validateAgentTemplateGetCfg(&tt.cfg) if tt.wantErr != "" { require.Error(t, err) assert.Contains(t, err.Error(), tt.wantErr) @@ -60,7 +60,7 @@ func TestGetAgentTemplatesTableReportsHarnessReadiness(t *testing.T) { }) var output bytes.Buffer - err := get(context.Background(), clientSet.ApiV1alpha3().AgentTemplates("kagent"), &GetCfg{ + err := getAgentTemplates(context.Background(), clientSet.ApiV1alpha3().AgentTemplates("kagent"), &AgentTemplateGetCfg{ Namespace: "kagent", PageSize: 3, PageToken: "previous-page", }, clioutput.FormatTable, &output) require.NoError(t, err) @@ -79,7 +79,7 @@ func TestGetAgentTemplatesJSONPreservesListMetadata(t *testing.T) { clientSet := clientfake.NewSimpleClientset() clientSet.PrependReactor("list", "agenttemplates", func(action k8stesting.Action) (bool, runtime.Object, error) { options := action.(interface{ GetListOptions() metav1.ListOptions }).GetListOptions() - assert.Equal(t, int64(maxPageSize), options.Limit) + assert.Equal(t, int64(agentTemplateMaxPageSize), options.Limit) return true, &apiv1alpha3.AgentTemplateList{ ListMeta: metav1.ListMeta{Continue: "next-page"}, Items: []apiv1alpha3.AgentTemplate{ @@ -91,7 +91,7 @@ func TestGetAgentTemplatesJSONPreservesListMetadata(t *testing.T) { }) var output bytes.Buffer - err := get(context.Background(), clientSet.ApiV1alpha3().AgentTemplates("kagent"), &GetCfg{ + err := getAgentTemplates(context.Background(), clientSet.ApiV1alpha3().AgentTemplates("kagent"), &AgentTemplateGetCfg{ Namespace: "kagent", }, clioutput.FormatJSON, &output) require.NoError(t, err) diff --git a/go/core/cli/internal/cli/agentinstance/get.go b/go/core/cli/internal/commands/agentinstance/get.go similarity index 72% rename from go/core/cli/internal/cli/agentinstance/get.go rename to go/core/cli/internal/commands/agentinstance/get.go index 7da3049e2..2954250ec 100644 --- a/go/core/cli/internal/cli/agentinstance/get.go +++ b/go/core/cli/internal/commands/agentinstance/get.go @@ -12,8 +12,9 @@ import ( "github.com/google/uuid" "github.com/jedib0t/go-pretty/v6/table" apiv1alpha1 "github.com/kagent-dev/kagent/go/api/gen/kagent/api/v1alpha1" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" + "github.com/spf13/cobra" "google.golang.org/protobuf/types/known/timestamppb" ) @@ -26,15 +27,19 @@ type getClient interface { // GetCfg configures AgentInstance get and list operations. type GetCfg struct { - Connection *connection.Options OutputFormat string InstanceID string PageSize int32 PageToken string } -// GetCmd gets one AgentInstance or lists the caller's AgentInstances. -func GetCmd(ctx context.Context, cfg *GetCfg, out io.Writer) (err error) { +// runGet gets one AgentInstance or lists the caller's AgentInstances. +func runGet( + ctx context.Context, + options connection.Options, + cfg *GetCfg, + out io.Writer, +) (err error) { format, err := clioutput.Parse(cfg.OutputFormat) if err != nil { return err @@ -43,19 +48,14 @@ func GetCmd(ctx context.Context, cfg *GetCfg, out io.Writer) (err error) { return err } - portForward, err := connection.Connect(ctx, cfg.Connection) + session, err := connection.Open(ctx, options) if err != nil { - return fmt.Errorf("connect to kagent: %w", err) - } - if portForward != nil { - defer portForward.Stop() + return err } - - clientSet := cfg.Connection.Client() defer func() { - err = errors.Join(err, clientSet.Close()) + err = errors.Join(err, session.Close()) }() - return get(ctx, clientSet.AgentInstance, cfg.Connection.Namespace, cfg, format, out) + return get(ctx, session.Client.AgentInstance, session.Namespace, cfg, format, out) } func validateGetCfg(cfg *GetCfg) error { @@ -154,3 +154,33 @@ func formatTimestamp(timestamp *timestamppb.Timestamp) string { } return timestamp.AsTime().UTC().Format(time.RFC3339) } + +// NewGetCmd constructs the AgentInstance get/list command. +func NewGetCmd() *cobra.Command { + cfg := &GetCfg{} + cmd := &cobra.Command{ + Use: "agent-instance [ID]", + Short: "Get an AgentInstance or list your AgentInstances", + Args: cobra.MaximumNArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + format, err := clioutput.FromCommand(cmd) + if err != nil { + return err + } + var instanceID string + if len(args) == 1 { + instanceID = args[0] + } + cfg.OutputFormat = format + cfg.InstanceID = instanceID + return runGet(cmd.Context(), options, cfg, cmd.OutOrStdout()) + }, + } + cmd.Flags().Int32Var(&cfg.PageSize, "page-size", 0, "Number of AgentInstances to return (default 50, maximum 100)") + cmd.Flags().StringVar(&cfg.PageToken, "page-token", "", "Token returned by the previous page") + return cmd +} diff --git a/go/core/cli/internal/cli/agentinstance/get_test.go b/go/core/cli/internal/commands/agentinstance/get_test.go similarity index 98% rename from go/core/cli/internal/cli/agentinstance/get_test.go rename to go/core/cli/internal/commands/agentinstance/get_test.go index 187ff01f8..a45952f64 100644 --- a/go/core/cli/internal/cli/agentinstance/get_test.go +++ b/go/core/cli/internal/commands/agentinstance/get_test.go @@ -8,7 +8,7 @@ import ( "time" apiv1alpha1 "github.com/kagent-dev/kagent/go/api/gen/kagent/api/v1alpha1" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "google.golang.org/protobuf/types/known/timestamppb" diff --git a/go/core/cli/internal/cli/agentinstance/invoke.go b/go/core/cli/internal/commands/agentinstance/invoke.go similarity index 82% rename from go/core/cli/internal/cli/agentinstance/invoke.go rename to go/core/cli/internal/commands/agentinstance/invoke.go index 9fc6d111c..05c69e249 100644 --- a/go/core/cli/internal/cli/agentinstance/invoke.go +++ b/go/core/cli/internal/commands/agentinstance/invoke.go @@ -14,14 +14,14 @@ import ( "github.com/a2aproject/a2a-go/v2/a2apb/v1/pbconv" "github.com/google/uuid" clia2a "github.com/kagent-dev/kagent/go/core/cli/internal/a2a" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" + "github.com/spf13/cobra" ) var errTruncatedA2AStream = errors.New("a2a stream ended before returning a final result") type InvokeCfg struct { - Connection *connection.Options OutputFormat string Task string File string @@ -30,7 +30,13 @@ type InvokeCfg struct { Token string } -func InvokeCmd(ctx context.Context, cfg *InvokeCfg, in io.Reader, out io.Writer) (err error) { +func runInvoke( + ctx context.Context, + options connection.Options, + cfg *InvokeCfg, + in io.Reader, + out io.Writer, +) (err error) { format, err := clioutput.Parse(cfg.OutputFormat) if err != nil { return err @@ -47,19 +53,14 @@ func InvokeCmd(ctx context.Context, cfg *InvokeCfg, in io.Reader, out io.Writer) return errors.New("model API key must not contain whitespace") } - portForward, err := connection.Connect(ctx, cfg.Connection) + session, err := connection.Open(ctx, options) if err != nil { - return fmt.Errorf("connect to kagent: %w", err) - } - if portForward != nil { - defer portForward.Stop() + return err } - - clientSet := cfg.Connection.Client() defer func() { - err = errors.Join(err, clientSet.Close()) + err = errors.Join(err, session.Close()) }() - a2aClient, err := clientSet.A2A.ForAgentInstance(ctx, cfg.Connection.Namespace, instanceID.String()) + a2aClient, err := session.Client.A2A.ForAgentInstance(ctx, session.Namespace, instanceID.String()) if err != nil { return fmt.Errorf("create AgentInstance A2A client: %w", err) } @@ -327,3 +328,36 @@ func sendResultError(result a2atype.SendMessageResult) error { return fmt.Errorf("AgentInstance task %s returned before reaching a final state: %s", task.ID, task.Status.State) } } + +// NewInvokeCmd constructs the AgentInstance invoke command. +func NewInvokeCmd() *cobra.Command { + cfg := &InvokeCfg{} + cmd := &cobra.Command{ + Use: "invoke", + Short: "Invoke an AgentInstance", + Long: `Invoke an existing AgentInstance through the A2A API.`, + Args: cobra.NoArgs, + Example: `kagent invoke --agent-instance 8bd650a8-9775-488f-8bc1-0d52bf7bdcab --task "Get all the pods"`, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + format, err := clioutput.FromCommand(cmd) + if err != nil { + return err + } + cfg.OutputFormat = format + return runInvoke(cmd.Context(), options, cfg, cmd.InOrStdin(), cmd.OutOrStdout()) + }, + } + cmd.Flags().StringVar(&cfg.AgentInstance, "agent-instance", "", "AgentInstance ID") + cmd.Flags().StringVarP(&cfg.Task, "task", "t", "", "Task text") + cmd.Flags().StringVarP(&cfg.File, "file", "f", "", "Read task text from a file or - for stdin") + cmd.Flags().BoolVarP(&cfg.Stream, "stream", "S", false, "Stream the response") + cmd.Flags().StringVar(&cfg.Token, "token", "", "Model API key passed through as an A2A Bearer token") + _ = cmd.MarkFlagRequired("agent-instance") + cmd.MarkFlagsOneRequired("task", "file") + cmd.MarkFlagsMutuallyExclusive("task", "file") + return cmd +} diff --git a/go/core/cli/internal/cli/agentinstance/invoke_test.go b/go/core/cli/internal/commands/agentinstance/invoke_test.go similarity index 99% rename from go/core/cli/internal/cli/agentinstance/invoke_test.go rename to go/core/cli/internal/commands/agentinstance/invoke_test.go index 1ab4a371f..d2469e11b 100644 --- a/go/core/cli/internal/cli/agentinstance/invoke_test.go +++ b/go/core/cli/internal/commands/agentinstance/invoke_test.go @@ -13,7 +13,7 @@ import ( a2atype "github.com/a2aproject/a2a-go/v2/a2a" "github.com/a2aproject/a2a-go/v2/a2aclient" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) diff --git a/go/core/cli/internal/cli/agentinstance/lifecycle.go b/go/core/cli/internal/commands/agentinstance/lifecycle.go similarity index 54% rename from go/core/cli/internal/cli/agentinstance/lifecycle.go rename to go/core/cli/internal/commands/agentinstance/lifecycle.go index f1fce72c0..7659f89a9 100644 --- a/go/core/cli/internal/cli/agentinstance/lifecycle.go +++ b/go/core/cli/internal/commands/agentinstance/lifecycle.go @@ -8,8 +8,9 @@ import ( "github.com/google/uuid" apiv1alpha1 "github.com/kagent-dev/kagent/go/api/gen/kagent/api/v1alpha1" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" + "github.com/spf13/cobra" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "google.golang.org/protobuf/proto" @@ -22,7 +23,6 @@ type lifecycleClient interface { // CreateCfg configures AgentInstance creation. type CreateCfg struct { - Connection *connection.Options OutputFormat string Harness string AgentTemplate string @@ -31,32 +31,31 @@ type CreateCfg struct { // DeleteCfg configures AgentInstance deletion. type DeleteCfg struct { - Connection *connection.Options OutputFormat string InstanceID string } -// CreateCmd creates an AgentInstance. -func CreateCmd(ctx context.Context, cfg *CreateCfg, out io.Writer) (err error) { +// runCreate creates an AgentInstance. +func runCreate( + ctx context.Context, + options connection.Options, + cfg *CreateCfg, + out io.Writer, +) (err error) { format, err := clioutput.Parse(cfg.OutputFormat) if err != nil { return err } ensureRequestID(cfg) - portForward, err := connection.Connect(ctx, cfg.Connection) + session, err := connection.Open(ctx, options) if err != nil { - return fmt.Errorf("connect to kagent: %w", err) - } - if portForward != nil { - defer portForward.Stop() + return err } - - clientSet := cfg.Connection.Client() defer func() { - err = errors.Join(err, clientSet.Close()) + err = errors.Join(err, session.Close()) }() - return create(ctx, clientSet.AgentInstance, cfg.Connection.Namespace, cfg, format, out) + return create(ctx, session.Client.AgentInstance, session.Namespace, cfg, format, out) } func ensureRequestID(cfg *CreateCfg) { @@ -65,25 +64,25 @@ func ensureRequestID(cfg *CreateCfg) { } } -// DeleteCmd deletes an AgentInstance. -func DeleteCmd(ctx context.Context, cfg *DeleteCfg, out io.Writer) (err error) { +// runDelete deletes an AgentInstance. +func runDelete( + ctx context.Context, + options connection.Options, + cfg *DeleteCfg, + out io.Writer, +) (err error) { format, err := clioutput.Parse(cfg.OutputFormat) if err != nil { return err } - portForward, err := connection.Connect(ctx, cfg.Connection) + session, err := connection.Open(ctx, options) if err != nil { - return fmt.Errorf("connect to kagent: %w", err) - } - if portForward != nil { - defer portForward.Stop() + return err } - - clientSet := cfg.Connection.Client() defer func() { - err = errors.Join(err, clientSet.Close()) + err = errors.Join(err, session.Close()) }() - return deleteAgentInstance(ctx, clientSet.AgentInstance, cfg.Connection.Namespace, cfg, format, out) + return deleteAgentInstance(ctx, session.Client.AgentInstance, session.Namespace, cfg, format, out) } func create( @@ -141,3 +140,55 @@ func writeLifecycleResult( } return writeInstancesTable(w, []*apiv1alpha1.AgentInstance{instance}, "") } + +// NewCreateCmd constructs the AgentInstance create command. +func NewCreateCmd() *cobra.Command { + cfg := &CreateCfg{} + cmd := &cobra.Command{ + Use: "agent-instance", + Short: "Create an AgentInstance", + Args: cobra.NoArgs, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + format, err := clioutput.FromCommand(cmd) + if err != nil { + return err + } + cfg.OutputFormat = format + return runCreate(cmd.Context(), options, cfg, cmd.OutOrStdout()) + }, + } + cmd.Flags().StringVar(&cfg.Harness, "harness", "", "Harness name") + cmd.Flags().StringVar(&cfg.AgentTemplate, "agent-template", "", "AgentTemplate name") + cmd.Flags().StringVar(&cfg.RequestID, "request-id", "", "Idempotency key (generated when omitted)") + _ = cmd.MarkFlagRequired("harness") + _ = cmd.MarkFlagRequired("agent-template") + return cmd +} + +// NewDeleteCmd constructs the AgentInstance delete command. +func NewDeleteCmd() *cobra.Command { + cfg := &DeleteCfg{} + cmd := &cobra.Command{ + Use: "agent-instance ID", + Short: "Delete an AgentInstance", + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + format, err := clioutput.FromCommand(cmd) + if err != nil { + return err + } + cfg.OutputFormat = format + cfg.InstanceID = args[0] + return runDelete(cmd.Context(), options, cfg, cmd.OutOrStdout()) + }, + } + return cmd +} diff --git a/go/core/cli/internal/cli/agentinstance/lifecycle_test.go b/go/core/cli/internal/commands/agentinstance/lifecycle_test.go similarity index 98% rename from go/core/cli/internal/cli/agentinstance/lifecycle_test.go rename to go/core/cli/internal/commands/agentinstance/lifecycle_test.go index 9e51ebdf7..cf60f3bef 100644 --- a/go/core/cli/internal/cli/agentinstance/lifecycle_test.go +++ b/go/core/cli/internal/commands/agentinstance/lifecycle_test.go @@ -8,7 +8,7 @@ import ( "github.com/google/uuid" apiv1alpha1 "github.com/kagent-dev/kagent/go/api/gen/kagent/api/v1alpha1" - clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/cli/output" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "google.golang.org/grpc/codes" diff --git a/go/core/cli/internal/cli/agent/bug_report.go b/go/core/cli/internal/commands/bug_report.go similarity index 83% rename from go/core/cli/internal/cli/agent/bug_report.go rename to go/core/cli/internal/commands/bug_report.go index 3c1083178..73170ddc9 100644 --- a/go/core/cli/internal/cli/agent/bug_report.go +++ b/go/core/cli/internal/commands/bug_report.go @@ -1,4 +1,4 @@ -package cli +package commands import ( "fmt" @@ -8,9 +8,11 @@ import ( "time" commonexec "github.com/kagent-dev/kagent/go/core/cli/internal/common/exec" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + "github.com/spf13/cobra" ) -func BugReportCmd(namespace string, verbose bool) { +func runBugReport(namespace string, verbose bool) { // Create a temporary directory for bug report timestamp := time.Now().Format("20060102-150405") reportDir := fmt.Sprintf("kagent-bug-report-%s", timestamp) @@ -116,3 +118,26 @@ func BugReportCmd(namespace string, verbose bool) { fmt.Printf("Bug report generated in directory: %s\n", reportDir) fmt.Println("WARNING: Please review and scrub any sensitive information from agent.yaml before sharing the bug report.") } + +// NewBugReportCmd constructs the kagent bug-report command. +func NewBugReportCmd() *cobra.Command { + return &cobra.Command{ + Use: "bug-report", + Short: "Generate a bug report", + Long: `Generate a bug report`, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + session, err := connection.Open(cmd.Context(), options) + if err != nil { + fmt.Fprintf(os.Stderr, "Error connecting to server: %v\n", err) + return nil + } + defer session.Close() //nolint:errcheck + runBugReport(options.Namespace, options.Verbose) + return nil + }, + } +} diff --git a/go/core/cli/internal/cli/agent/const.go b/go/core/cli/internal/commands/const.go similarity index 99% rename from go/core/cli/internal/cli/agent/const.go rename to go/core/cli/internal/commands/const.go index 08edd89a5..7aceb9efc 100644 --- a/go/core/cli/internal/cli/agent/const.go +++ b/go/core/cli/internal/commands/const.go @@ -1,4 +1,4 @@ -package cli +package commands import ( "os" diff --git a/go/core/cli/internal/commands/dashboard.go b/go/core/cli/internal/commands/dashboard.go new file mode 100644 index 000000000..f4bf591a0 --- /dev/null +++ b/go/core/cli/internal/commands/dashboard.go @@ -0,0 +1,23 @@ +package commands + +import ( + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + "github.com/spf13/cobra" +) + +// NewDashboardCmd constructs the kagent dashboard command. +func NewDashboardCmd() *cobra.Command { + return &cobra.Command{ + Use: "dashboard", + Short: "Open the kagent dashboard", + Long: `Open the kagent dashboard`, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + runDashboard(cmd.Context(), options.Namespace) + return nil + }, + } +} diff --git a/go/core/cli/internal/cli/agent/dashboard_darwin.go b/go/core/cli/internal/commands/dashboard_darwin.go similarity index 93% rename from go/core/cli/internal/cli/agent/dashboard_darwin.go rename to go/core/cli/internal/commands/dashboard_darwin.go index eddf428f4..c08ecf143 100644 --- a/go/core/cli/internal/cli/agent/dashboard_darwin.go +++ b/go/core/cli/internal/commands/dashboard_darwin.go @@ -1,6 +1,6 @@ //go:build darwin -package cli +package commands import ( "context" @@ -11,7 +11,7 @@ import ( "time" ) -func DashboardCmd(ctx context.Context, namespace string) { +func runDashboard(ctx context.Context, namespace string) { ctx, cancel := context.WithCancel(ctx) cmd := exec.CommandContext(ctx, "kubectl", "-n", namespace, "port-forward", "service/kagent-ui", "8082:8080") diff --git a/go/core/cli/internal/cli/agent/dashboard.go b/go/core/cli/internal/commands/dashboard_other.go similarity index 83% rename from go/core/cli/internal/cli/agent/dashboard.go rename to go/core/cli/internal/commands/dashboard_other.go index 286ce1123..11bbfb349 100644 --- a/go/core/cli/internal/cli/agent/dashboard.go +++ b/go/core/cli/internal/commands/dashboard_other.go @@ -1,6 +1,6 @@ //go:build !darwin -package cli +package commands import ( "context" @@ -8,7 +8,7 @@ import ( "os" ) -func DashboardCmd(ctx context.Context, namespace string) { +func runDashboard(ctx context.Context, namespace string) { fmt.Fprintln(os.Stderr, "Dashboard is not available on this platform") fmt.Fprintln(os.Stderr, "You can easily start the dashboard by running:") fmt.Fprintf(os.Stderr, "kubectl port-forward -n %s service/kagent-ui 8082:8080\n", namespace) diff --git a/go/core/cli/internal/commands/db/db.go b/go/core/cli/internal/commands/db/db.go new file mode 100644 index 000000000..b7fdb1ada --- /dev/null +++ b/go/core/cli/internal/commands/db/db.go @@ -0,0 +1,114 @@ +// Package db wires the shared database subcommand to the CLI's migration tracks. +package db + +import ( + "context" + "fmt" + "os" + "strconv" + "time" + + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + dbcli "github.com/kagent-dev/kagent/go/core/pkg/cli/db" + dbmigrate "github.com/kagent-dev/kagent/go/core/pkg/cli/db/migrate" + "github.com/kagent-dev/kagent/go/core/pkg/migrations" + "github.com/spf13/cobra" + corev1 "k8s.io/api/core/v1" + "k8s.io/client-go/tools/clientcmd" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// vectorEnabledKey names two lookups that deliberately share it: the CLI's +// own DATABASE_VECTOR_ENABLED env var (a local operator override), and the +// controller-configmap key the chart renders — the value the controller pod +// itself consumes via envFrom. Same name, two different places. +const vectorEnabledKey = "DATABASE_VECTOR_ENABLED" + +// NewDBCmd constructs the kagent db command. The source callback the shared db +// package takes carries no command, so the namespace is captured just before a +// subcommand runs, when the root's flags have been parsed. +func NewDBCmd() *cobra.Command { + var namespace string + cmd := dbcli.NewCommandFromFunc(migrationSources(&namespace)) + cmd.PersistentPreRunE = func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + namespace = options.Namespace + return nil + } + return cmd +} + +// migrationSources resolves the built-in migration tracks when a db +// subcommand runs (never during command construction, so unrelated commands +// do no work and print no warnings). The vector track is gated, in order of +// precedence, on: the DATABASE_VECTOR_ENABLED env var in the CLI's own +// environment (explicit operator intent, works without a cluster), the +// controller's configmap on the live cluster (the same value the server +// reads), and finally the controller's default (enabled). +func migrationSources(namespace *string) dbmigrate.SourcesFunc { + return func(ctx context.Context) ([]migrations.Source, error) { + vectorEnabled := true + if v := os.Getenv(vectorEnabledKey); v != "" { + b, err := strconv.ParseBool(v) + if err != nil { + fmt.Fprintf(os.Stderr, "warning: invalid %s=%q; assuming true\n", vectorEnabledKey, v) + } else { + vectorEnabled = b + } + } else if b, ok := clusterVectorEnabled(ctx, *namespace); ok { + vectorEnabled = b + } + return migrations.BuiltinSources(vectorEnabled), nil + } +} + +// clusterVectorEnabled reads the vectorEnabledKey entry from the controller +// configmap in the given namespace (the same "kagent-controller" default +// naming the rest of the CLI assumes) — the cluster-side counterpart of the +// env-var override in migrationSources. When the value is used it says so on +// stderr, naming the kubeconfig context it was read from — the lookup follows +// the *current* context, so this is the operator's cue that the cluster and +// their --db-url had better be the same install. Best-effort: reports +// ok=false when no cluster is reachable, the configmap is absent, or the +// value doesn't parse — callers fall back to the default. +func clusterVectorEnabled(ctx context.Context, namespace string) (enabled, ok bool) { + restConfig, err := clientcmd.NewNonInteractiveDeferredLoadingClientConfig( + clientcmd.NewDefaultClientConfigLoadingRules(), + &clientcmd.ConfigOverrides{}, + ).ClientConfig() + if err != nil { + return false, false + } + k8sClient, err := client.New(restConfig, client.Options{}) + if err != nil { + return false, false + } + ctx, cancel := context.WithTimeout(ctx, 3*time.Second) + defer cancel() + var cm corev1.ConfigMap + if err := k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: "kagent-controller"}, &cm); err != nil { + return false, false + } + b, err := strconv.ParseBool(cm.Data[vectorEnabledKey]) + if err != nil { + return false, false + } + // Trailing blank line separates the notice from the command's stdout + // when both land on a terminal; piped stdout is unaffected. + fmt.Fprintf(os.Stderr, "resolved vector track from cluster context %q: configmap %s/kagent-controller has %s=%t (set %s to override)\n\n", + currentKubeContext(), namespace, vectorEnabledKey, b, vectorEnabledKey) + return b, true +} + +// currentKubeContext names the kubeconfig context the CLI's Kubernetes client +// dials, for operator-facing messages. Best-effort. +func currentKubeContext() string { + raw, err := clientcmd.NewDefaultClientConfigLoadingRules().Load() + if err != nil || raw.CurrentContext == "" { + return "(current kubeconfig context)" + } + return raw.CurrentContext +} diff --git a/go/core/cli/internal/cli/envdoc/envdoc.go b/go/core/cli/internal/commands/env.go similarity index 94% rename from go/core/cli/internal/cli/envdoc/envdoc.go rename to go/core/cli/internal/commands/env.go index a949efa82..c0a680638 100644 --- a/go/core/cli/internal/cli/envdoc/envdoc.go +++ b/go/core/cli/internal/commands/env.go @@ -1,4 +1,4 @@ -package envdoc +package commands import ( "fmt" @@ -7,13 +7,9 @@ import ( "github.com/spf13/cobra" ) -var ( - format string - component string -) - // NewEnvCmd returns a cobra command that generates environment variable documentation. func NewEnvCmd() *cobra.Command { + var format, component string cmd := &cobra.Command{ Use: "env", Hidden: true, diff --git a/go/core/cli/internal/commands/env_test.go b/go/core/cli/internal/commands/env_test.go new file mode 100644 index 000000000..95282b5fa --- /dev/null +++ b/go/core/cli/internal/commands/env_test.go @@ -0,0 +1,23 @@ +package commands + +import ( + "bytes" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestEnvCommandsOwnIndependentFlagState(t *testing.T) { + first := NewEnvCmd() + second := NewEnvCmd() + first.SetArgs([]string{"--format", "json", "--component", "cli"}) + first.SetOut(&bytes.Buffer{}) + secondOutput := &bytes.Buffer{} + second.SetOut(secondOutput) + + require.NoError(t, first.ExecuteContext(t.Context())) + require.NoError(t, second.ExecuteContext(t.Context())) + assert.Contains(t, secondOutput.String(), "# Kagent Environment Variables") + assert.Contains(t, secondOutput.String(), "## controller") +} diff --git a/go/core/cli/internal/cli/agent/install.go b/go/core/cli/internal/commands/install.go similarity index 86% rename from go/core/cli/internal/cli/agent/install.go rename to go/core/cli/internal/commands/install.go index 2702f163a..9d7ec8b16 100644 --- a/go/core/cli/internal/cli/agent/install.go +++ b/go/core/cli/internal/commands/install.go @@ -1,4 +1,4 @@ -package cli +package commands import ( "context" @@ -12,15 +12,15 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" "github.com/kagent-dev/kagent/go/core/internal/version" "github.com/kagent-dev/kagent/go/core/pkg/env" + "github.com/spf13/cobra" "github.com/briandowns/spinner" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" "github.com/kagent-dev/kagent/go/core/cli/internal/profiles" ) type InstallCfg struct { - Connection *connection.Options - Profile string + Profile string } // installChart installs or upgrades a Helm chart with the given parameters @@ -64,7 +64,7 @@ func installChart(ctx context.Context, chartName string, namespace string, regis return "", nil } -func InstallCmd(ctx context.Context, cfg *InstallCfg) *connection.PortForward { +func runInstall(ctx context.Context, options connection.Options, cfg *InstallCfg) *connection.PortForward { if version.Version == "dev" { fmt.Fprintln(os.Stderr, "Installation requires released version of kagent") return nil @@ -101,7 +101,7 @@ func InstallCmd(ctx context.Context, cfg *InstallCfg) *connection.PortForward { helmConfig.inlineValues = profiles.GetProfileYaml(cfg.Profile) } - return install(ctx, cfg.Connection, helmConfig, modelProvider) + return install(ctx, &options, helmConfig, modelProvider) } // helmConfig is the config for the kagent chart @@ -231,7 +231,7 @@ func deleteCRDs(ctx context.Context) error { return nil } -func UninstallCmd(ctx context.Context, namespace string) { +func runUninstall(ctx context.Context, namespace string) { // Check if helm is available if err := checkHelmAvailable(); err != nil { fmt.Fprintln(os.Stderr, err) @@ -303,3 +303,43 @@ func checkHelmAvailable() error { } return nil } + +// NewInstallCmd constructs the kagent install command. +func NewInstallCmd() *cobra.Command { + cfg := &InstallCfg{} + cmd := &cobra.Command{ + Use: "install", + Short: "Install kagent", + Long: `Install kagent`, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + runInstall(cmd.Context(), options, cfg) + return nil + }, + } + cmd.Flags().StringVar(&cfg.Profile, "profile", "", "Installation profile (minimal|demo)") + _ = cmd.RegisterFlagCompletionFunc("profile", func(_ *cobra.Command, _ []string, _ string) ([]string, cobra.ShellCompDirective) { + return profiles.Profiles, cobra.ShellCompDirectiveNoFileComp + }) + return cmd +} + +// NewUninstallCmd constructs the kagent uninstall command. +func NewUninstallCmd() *cobra.Command { + return &cobra.Command{ + Use: "uninstall", + Short: "Uninstall kagent", + Long: `Uninstall kagent`, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + runUninstall(cmd.Context(), options.Namespace) + return nil + }, + } +} diff --git a/go/core/cli/internal/cli/mcp/add_tool.go b/go/core/cli/internal/commands/mcp/add_tool.go similarity index 100% rename from go/core/cli/internal/cli/mcp/add_tool.go rename to go/core/cli/internal/commands/mcp/add_tool.go diff --git a/go/core/cli/internal/cli/mcp/build.go b/go/core/cli/internal/commands/mcp/build.go similarity index 100% rename from go/core/cli/internal/cli/mcp/build.go rename to go/core/cli/internal/commands/mcp/build.go diff --git a/go/core/cli/internal/cli/mcp/build_test.go b/go/core/cli/internal/commands/mcp/build_test.go similarity index 100% rename from go/core/cli/internal/cli/mcp/build_test.go rename to go/core/cli/internal/commands/mcp/build_test.go diff --git a/go/core/cli/internal/cli/mcp/deploy.go b/go/core/cli/internal/commands/mcp/deploy.go similarity index 100% rename from go/core/cli/internal/cli/mcp/deploy.go rename to go/core/cli/internal/commands/mcp/deploy.go diff --git a/go/core/cli/internal/cli/mcp/init.go b/go/core/cli/internal/commands/mcp/init.go similarity index 100% rename from go/core/cli/internal/cli/mcp/init.go rename to go/core/cli/internal/commands/mcp/init.go diff --git a/go/core/cli/internal/cli/mcp/init_test.go b/go/core/cli/internal/commands/mcp/init_test.go similarity index 100% rename from go/core/cli/internal/cli/mcp/init_test.go rename to go/core/cli/internal/commands/mcp/init_test.go diff --git a/go/core/cli/internal/cli/mcp/inspector.go b/go/core/cli/internal/commands/mcp/inspector.go similarity index 100% rename from go/core/cli/internal/cli/mcp/inspector.go rename to go/core/cli/internal/commands/mcp/inspector.go diff --git a/go/core/cli/internal/cli/mcp/integration_test.go b/go/core/cli/internal/commands/mcp/integration_test.go similarity index 100% rename from go/core/cli/internal/cli/mcp/integration_test.go rename to go/core/cli/internal/commands/mcp/integration_test.go diff --git a/go/core/cli/internal/cli/mcp/root.go b/go/core/cli/internal/commands/mcp/root.go similarity index 100% rename from go/core/cli/internal/cli/mcp/root.go rename to go/core/cli/internal/commands/mcp/root.go diff --git a/go/core/cli/internal/cli/mcp/run.go b/go/core/cli/internal/commands/mcp/run.go similarity index 100% rename from go/core/cli/internal/cli/mcp/run.go rename to go/core/cli/internal/commands/mcp/run.go diff --git a/go/core/cli/internal/cli/mcp/run_test.go b/go/core/cli/internal/commands/mcp/run_test.go similarity index 100% rename from go/core/cli/internal/cli/mcp/run_test.go rename to go/core/cli/internal/commands/mcp/run_test.go diff --git a/go/core/cli/internal/cli/mcp/secrets.go b/go/core/cli/internal/commands/mcp/secrets.go similarity index 100% rename from go/core/cli/internal/cli/mcp/secrets.go rename to go/core/cli/internal/commands/mcp/secrets.go diff --git a/go/core/cli/internal/commands/version.go b/go/core/cli/internal/commands/version.go new file mode 100644 index 000000000..78eedc0b4 --- /dev/null +++ b/go/core/cli/internal/commands/version.go @@ -0,0 +1,58 @@ +package commands + +import ( + "context" + "encoding/json" + "os" + "time" + + "github.com/kagent-dev/kagent/go/api/client" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + "github.com/kagent-dev/kagent/go/core/internal/version" + "github.com/spf13/cobra" +) + +func runVersion(clientSet *client.ClientSet) { + versionInfo := map[string]string{ + "kagent_version": version.Version, + "git_commit": version.GitCommit, + "build_date": version.BuildDate, + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second*5) + defer cancel() + serverVersion, err := clientSet.Version.GetVersion(ctx) + if err != nil { + versionInfo["backend_version"] = "unknown" + } else { + versionInfo["backend_version"] = serverVersion.KAgentVersion + } + + json.NewEncoder(os.Stdout).Encode(versionInfo) //nolint:errcheck +} + +// NewVersionCmd constructs the kagent version command. +func NewVersionCmd() *cobra.Command { + return &cobra.Command{ + Use: "version", + Short: "Print the kagent version", + Long: `Print the kagent version`, + RunE: func(cmd *cobra.Command, _ []string) error { + options, err := connection.OptionsFromCommand(cmd) + if err != nil { + return err + } + // The CLI version prints whether or not the server answers; an + // unreachable server reports its version as "unknown". + session, _ := connection.Open(cmd.Context(), options) + if session == nil { + clientSet := options.Client() + defer clientSet.Close() //nolint:errcheck + runVersion(clientSet) + return nil + } + defer session.Close() //nolint:errcheck + runVersion(session.Client) + return nil + }, + } +} diff --git a/go/core/cli/internal/cli/connection/connection.go b/go/core/cli/internal/connection/connection.go similarity index 97% rename from go/core/cli/internal/cli/connection/connection.go rename to go/core/cli/internal/connection/connection.go index 2ce0f049d..70cf3a361 100644 --- a/go/core/cli/internal/cli/connection/connection.go +++ b/go/core/cli/internal/connection/connection.go @@ -30,7 +30,8 @@ const ( kubectlErrorLimit = 8 << 10 ) -// Options contains only the settings needed to connect to kagent. +// Options is how the CLI reaches kagent: where to dial, who to dial as, the +// namespace to port-forward into, and whether to narrate the attempt. type Options struct { KAgentURL string KAgentGRPCURL string diff --git a/go/core/cli/internal/cli/connection/connection_test.go b/go/core/cli/internal/connection/connection_test.go similarity index 100% rename from go/core/cli/internal/cli/connection/connection_test.go rename to go/core/cli/internal/connection/connection_test.go diff --git a/go/core/cli/internal/connection/flags.go b/go/core/cli/internal/connection/flags.go new file mode 100644 index 000000000..2d618b0dd --- /dev/null +++ b/go/core/cli/internal/connection/flags.go @@ -0,0 +1,70 @@ +package connection + +import ( + "github.com/spf13/cobra" + "github.com/spf13/pflag" +) + +// Flag names are unexported so RegisterFlags and OptionsFromCommand are the +// only things that can disagree about them, and they cannot. +const ( + flagKAgentURL = "kagent-url" + flagKAgentGRPCURL = "grpc-url" + flagKAgentGRPCTLS = "grpc-tls" + flagKAgentGRPCCAFile = "grpc-ca-file" + flagKAgentGRPCServerName = "grpc-server-name" + flagNamespace = "namespace" + flagVerbose = "verbose" + flagTimeout = "timeout" + flagUserID = "user-id" +) + +// RegisterFlags declares the CLI-wide connection flags, defaulted from DefaultOptions. +func RegisterFlags(flags *pflag.FlagSet) { + defaults := DefaultOptions() + flags.String(flagKAgentURL, defaults.KAgentURL, "KAgent REST URL") + flags.String(flagKAgentGRPCURL, defaults.KAgentGRPCURL, "KAgent gRPC target") + flags.Bool(flagKAgentGRPCTLS, defaults.KAgentGRPCTLS, "Use TLS for KAgent gRPC") + flags.String(flagKAgentGRPCCAFile, defaults.KAgentGRPCCAFile, "CA certificate file for KAgent gRPC") + flags.String(flagKAgentGRPCServerName, defaults.KAgentGRPCServerName, "TLS server name for KAgent gRPC") + flags.StringP(flagNamespace, "n", defaults.Namespace, "Namespace") + flags.BoolP(flagVerbose, "v", defaults.Verbose, "Verbose output") + flags.Duration(flagTimeout, defaults.Timeout, "Timeout") + flags.String(flagUserID, defaults.UserID, "Caller identity used to select the server-side data partition") +} + +// OptionsFromCommand resolves connection options from the flags a command was +// invoked with, which include the root's persistent flags. +func OptionsFromCommand(cmd *cobra.Command) (Options, error) { + flags := cmd.Flags() + var options Options + var err error + if options.KAgentURL, err = flags.GetString(flagKAgentURL); err != nil { + return Options{}, err + } + if options.KAgentGRPCURL, err = flags.GetString(flagKAgentGRPCURL); err != nil { + return Options{}, err + } + if options.KAgentGRPCTLS, err = flags.GetBool(flagKAgentGRPCTLS); err != nil { + return Options{}, err + } + if options.KAgentGRPCCAFile, err = flags.GetString(flagKAgentGRPCCAFile); err != nil { + return Options{}, err + } + if options.KAgentGRPCServerName, err = flags.GetString(flagKAgentGRPCServerName); err != nil { + return Options{}, err + } + if options.Namespace, err = flags.GetString(flagNamespace); err != nil { + return Options{}, err + } + if options.Verbose, err = flags.GetBool(flagVerbose); err != nil { + return Options{}, err + } + if options.Timeout, err = flags.GetDuration(flagTimeout); err != nil { + return Options{}, err + } + if options.UserID, err = flags.GetString(flagUserID); err != nil { + return Options{}, err + } + return options, nil +} diff --git a/go/core/cli/internal/connection/flags_test.go b/go/core/cli/internal/connection/flags_test.go new file mode 100644 index 000000000..79bbfef79 --- /dev/null +++ b/go/core/cli/internal/connection/flags_test.go @@ -0,0 +1,49 @@ +package connection + +import ( + "testing" + "time" + + "github.com/spf13/cobra" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestOptionsFromCommandReadsInheritedFlags(t *testing.T) { + var got Options + root := &cobra.Command{Use: "root"} + RegisterFlags(root.PersistentFlags()) + root.AddCommand(&cobra.Command{ + Use: "child", + RunE: func(cmd *cobra.Command, _ []string) error { + var err error + got, err = OptionsFromCommand(cmd) + return err + }, + }) + root.SetArgs([]string{ + "child", + "--kagent-url", "https://api.example.test", + "--grpc-url", "grpc.example.test:443", + "--grpc-tls", + "--grpc-ca-file", "/tmp/ca.pem", + "--grpc-server-name", "grpc.example.test", + "--namespace", "agents", + "--verbose", + "--timeout", "12s", + "--user-id", "reviewer@example.test", + }) + + require.NoError(t, root.ExecuteContext(t.Context())) + assert.Equal(t, Options{ + KAgentURL: "https://api.example.test", + KAgentGRPCURL: "grpc.example.test:443", + KAgentGRPCTLS: true, + KAgentGRPCCAFile: "/tmp/ca.pem", + KAgentGRPCServerName: "grpc.example.test", + Namespace: "agents", + Verbose: true, + Timeout: 12 * time.Second, + UserID: "reviewer@example.test", + }, got) +} diff --git a/go/core/cli/internal/connection/session.go b/go/core/cli/internal/connection/session.go new file mode 100644 index 000000000..561118156 --- /dev/null +++ b/go/core/cli/internal/connection/session.go @@ -0,0 +1,41 @@ +package connection + +import ( + "context" + "fmt" + + "github.com/kagent-dev/kagent/go/api/client" +) + +// Session is a connected kagent client for one command invocation, together +// with the namespace the command is scoped to. +type Session struct { + Client *client.ClientSet + Namespace string + + portForward *PortForward +} + +// Open reaches the server, starting a port-forward when the default local +// endpoint is unreachable. The caller must Close the returned session. +func Open(ctx context.Context, options Options) (*Session, error) { + portForward, err := Connect(ctx, &options) + if err != nil { + return nil, fmt.Errorf("connect to kagent: %w", err) + } + return &Session{ + Client: options.Client(), + Namespace: options.Namespace, + portForward: portForward, + }, nil +} + +// Close releases the client before tearing down the port-forward it rode on. +func (s *Session) Close() error { + if s == nil { + return nil + } + err := s.Client.Close() + s.portForward.Stop() + return err +} diff --git a/go/core/cli/internal/cli/output/output.go b/go/core/cli/internal/output/output.go similarity index 81% rename from go/core/cli/internal/cli/output/output.go rename to go/core/cli/internal/output/output.go index 3b327168a..8f0e5fe7e 100644 --- a/go/core/cli/internal/cli/output/output.go +++ b/go/core/cli/internal/output/output.go @@ -6,10 +6,19 @@ import ( "fmt" "io" + "github.com/spf13/cobra" "google.golang.org/protobuf/encoding/protojson" "google.golang.org/protobuf/proto" ) +// FlagName is the root persistent flag that selects the output format. +const FlagName = "output-format" + +// FromCommand reads the output format a command was invoked with. +func FromCommand(cmd *cobra.Command) (string, error) { + return cmd.Flags().GetString(FlagName) +} + // Format selects the CLI payload encoding. type Format string diff --git a/go/core/cli/internal/cli/output/output_test.go b/go/core/cli/internal/output/output_test.go similarity index 100% rename from go/core/cli/internal/cli/output/output_test.go rename to go/core/cli/internal/output/output_test.go diff --git a/go/core/cli/internal/tui/workspace_test.go b/go/core/cli/internal/tui/workspace_test.go index 00f231b62..c1f84f8d6 100644 --- a/go/core/cli/internal/tui/workspace_test.go +++ b/go/core/cli/internal/tui/workspace_test.go @@ -10,7 +10,7 @@ import ( tea "github.com/charmbracelet/bubbletea" apiv1alpha1 "github.com/kagent-dev/kagent/go/api/gen/kagent/api/v1alpha1" clia2a "github.com/kagent-dev/kagent/go/core/cli/internal/a2a" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "google.golang.org/protobuf/types/known/timestamppb" diff --git a/go/core/cli/root.go b/go/core/cli/root.go new file mode 100644 index 000000000..8803cdfe1 --- /dev/null +++ b/go/core/cli/root.go @@ -0,0 +1,70 @@ +package cli + +import ( + "fmt" + "strings" + + "github.com/kagent-dev/kagent/go/core/cli/internal/commands" + agentinstancecli "github.com/kagent-dev/kagent/go/core/cli/internal/commands/agentinstance" + dbcli "github.com/kagent-dev/kagent/go/core/cli/internal/commands/db" + "github.com/kagent-dev/kagent/go/core/cli/internal/commands/mcp" + "github.com/kagent-dev/kagent/go/core/cli/internal/connection" + clioutput "github.com/kagent-dev/kagent/go/core/cli/internal/output" + "github.com/spf13/cobra" +) + +// Root creates a fresh kagent command tree. +func Root() *cobra.Command { + rootCmd := &cobra.Command{ + Use: "kagent", + Short: "kagent is a CLI for kagent", + Long: "kagent is a CLI for kagent", + SilenceErrors: true, + SilenceUsage: true, + RunE: runInteractive, + } + connection.RegisterFlags(rootCmd.PersistentFlags()) + rootCmd.PersistentFlags().StringP(clioutput.FlagName, "o", string(clioutput.FormatTable), "Output format") + + getCmd := newResourceGroupCmd("get", "Get a kagent resource") + createCmd := newResourceGroupCmd("create", "Create a kagent resource") + deleteCmd := newResourceGroupCmd("delete", "Delete a kagent resource") + + getCmd.AddCommand(agentinstancecli.NewGetCmd()) + getCmd.AddCommand(commands.NewGetAgentTemplateCmd()) + createCmd.AddCommand(agentinstancecli.NewCreateCmd()) + deleteCmd.AddCommand(agentinstancecli.NewDeleteCmd()) + + rootCmd.AddCommand( + getCmd, + createCmd, + deleteCmd, + agentinstancecli.NewInvokeCmd(), + commands.NewInstallCmd(), + commands.NewUninstallCmd(), + commands.NewBugReportCmd(), + commands.NewVersionCmd(), + commands.NewDashboardCmd(), + mcp.NewMCPCmd(), + commands.NewEnvCmd(), + dbcli.NewDBCmd(), + ) + return rootCmd +} + +// newResourceGroupCmd builds a parent command that only routes to resource subcommands. +func newResourceGroupCmd(use, short string) *cobra.Command { + return &cobra.Command{ + Use: use, + Short: short, + Long: short, + Args: cobra.NoArgs, + RunE: func(cmd *cobra.Command, _ []string) error { + resourceTypes := make([]string, 0, len(cmd.Commands())) + for _, child := range cmd.Commands() { + resourceTypes = append(resourceTypes, child.Name()) + } + return fmt.Errorf("resource type is required; available resource types: %s", strings.Join(resourceTypes, ", ")) + }, + } +} diff --git a/go/core/cli/cmd/kagent/main_test.go b/go/core/cli/root_test.go similarity index 50% rename from go/core/cli/cmd/kagent/main_test.go rename to go/core/cli/root_test.go index 292d454c0..d2012bc69 100644 --- a/go/core/cli/cmd/kagent/main_test.go +++ b/go/core/cli/root_test.go @@ -1,60 +1,31 @@ -package main +package cli_test import ( "bytes" - "context" "testing" - "time" - "github.com/kagent-dev/kagent/go/core/cli/internal/cli/connection" + "github.com/kagent-dev/kagent/go/core/cli" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) -func TestRootCommandUsesOptionValuesAsFlagDefaults(t *testing.T) { - opts := &rootOptions{ - Connection: connection.Options{ - KAgentURL: "http://kagent.example.test", - KAgentGRPCURL: "grpc.kagent.example.test:443", - KAgentGRPCTLS: true, - KAgentGRPCCAFile: "/tmp/kagent-ca.pem", - KAgentGRPCServerName: "grpc.kagent.example.test", - Namespace: "configured-ns", - Verbose: true, - Timeout: 45 * time.Second, - UserID: "configured-user", - }, - OutputFormat: "json", - } - - rootCmd := newRootCommand(context.Background(), opts) - - assert.Equal(t, "http://kagent.example.test", rootCmd.PersistentFlags().Lookup("kagent-url").DefValue) - assert.Equal(t, "grpc.kagent.example.test:443", rootCmd.PersistentFlags().Lookup("grpc-url").DefValue) - assert.Equal(t, "true", rootCmd.PersistentFlags().Lookup("grpc-tls").DefValue) - assert.Equal(t, "/tmp/kagent-ca.pem", rootCmd.PersistentFlags().Lookup("grpc-ca-file").DefValue) - assert.Equal(t, "grpc.kagent.example.test", rootCmd.PersistentFlags().Lookup("grpc-server-name").DefValue) - assert.Equal(t, "configured-ns", rootCmd.PersistentFlags().Lookup("namespace").DefValue) - assert.Equal(t, "json", rootCmd.PersistentFlags().Lookup("output-format").DefValue) - assert.Equal(t, "true", rootCmd.PersistentFlags().Lookup("verbose").DefValue) - assert.Equal(t, "45s", rootCmd.PersistentFlags().Lookup("timeout").DefValue) - assert.Equal(t, "configured-user", rootCmd.PersistentFlags().Lookup("user-id").DefValue) - - assert.Equal(t, "configured-ns", opts.Connection.Namespace) +func TestRootCommandUsesDefaultFlagValues(t *testing.T) { + rootCmd := cli.Root() + + assert.Equal(t, "http://localhost:8083", rootCmd.PersistentFlags().Lookup("kagent-url").DefValue) + assert.Equal(t, "localhost:8084", rootCmd.PersistentFlags().Lookup("grpc-url").DefValue) + assert.Equal(t, "false", rootCmd.PersistentFlags().Lookup("grpc-tls").DefValue) + assert.Empty(t, rootCmd.PersistentFlags().Lookup("grpc-ca-file").DefValue) + assert.Empty(t, rootCmd.PersistentFlags().Lookup("grpc-server-name").DefValue) + assert.Equal(t, "kagent", rootCmd.PersistentFlags().Lookup("namespace").DefValue) + assert.Equal(t, "table", rootCmd.PersistentFlags().Lookup("output-format").DefValue) + assert.Equal(t, "false", rootCmd.PersistentFlags().Lookup("verbose").DefValue) + assert.Equal(t, "5m0s", rootCmd.PersistentFlags().Lookup("timeout").DefValue) + assert.Equal(t, "admin@kagent.dev", rootCmd.PersistentFlags().Lookup("user-id").DefValue) } func TestRootCommandFlagsOverrideOptionValues(t *testing.T) { - opts := &rootOptions{ - Connection: connection.Options{ - KAgentURL: "http://kagent.example.test", - KAgentGRPCURL: "grpc.kagent.example.test:443", - Namespace: "configured-ns", - Timeout: 45 * time.Second, - }, - OutputFormat: "json", - } - - rootCmd := newRootCommand(context.Background(), opts) + rootCmd := cli.Root() require.NoError(t, rootCmd.ParseFlags([]string{ "--kagent-url", "http://flag.example.test", "--grpc-url", "grpc.flag.example.test:8443", @@ -68,20 +39,42 @@ func TestRootCommandFlagsOverrideOptionValues(t *testing.T) { "--user-id", "flag-user", })) - assert.Equal(t, "http://flag.example.test", opts.Connection.KAgentURL) - assert.Equal(t, "grpc.flag.example.test:8443", opts.Connection.KAgentGRPCURL) - assert.True(t, opts.Connection.KAgentGRPCTLS) - assert.Equal(t, "/tmp/flag-ca.pem", opts.Connection.KAgentGRPCCAFile) - assert.Equal(t, "grpc.flag.example.test", opts.Connection.KAgentGRPCServerName) - assert.Equal(t, "flag-ns", opts.Connection.Namespace) - assert.Equal(t, "yaml", opts.OutputFormat) - assert.True(t, opts.Connection.Verbose) - assert.Equal(t, 10*time.Second, opts.Connection.Timeout) - assert.Equal(t, "flag-user", opts.Connection.UserID) + want := map[string]string{ + "kagent-url": "http://flag.example.test", + "grpc-url": "grpc.flag.example.test:8443", + "grpc-tls": "true", + "grpc-ca-file": "/tmp/flag-ca.pem", + "grpc-server-name": "grpc.flag.example.test", + "namespace": "flag-ns", + "output-format": "yaml", + "verbose": "true", + "timeout": "10s", + "user-id": "flag-user", + } + for name, value := range want { + assert.Equal(t, value, rootCmd.PersistentFlags().Lookup(name).Value.String()) + } +} + +func TestRootCommandAllowsNoTimeout(t *testing.T) { + rootCmd := cli.Root() + + require.NoError(t, rootCmd.ParseFlags([]string{"--timeout", "0"})) + assert.Equal(t, "0s", rootCmd.PersistentFlags().Lookup("timeout").Value.String()) +} + +func TestRootCommandsOwnIndependentFlagState(t *testing.T) { + first := cli.Root() + second := cli.Root() + + require.NoError(t, first.ParseFlags([]string{"--namespace", "first"})) + + assert.Equal(t, "first", first.PersistentFlags().Lookup("namespace").Value.String()) + assert.Equal(t, "kagent", second.PersistentFlags().Lookup("namespace").Value.String()) } func TestRootCommandDoesNotValidateClientFlagsForIndependentCommand(t *testing.T) { - rootCmd := newRootCommand(t.Context(), defaultRootOptions()) + rootCmd := cli.Root() rootCmd.SetArgs([]string{"--output-format", "yaml", "--user-id", "invalid user", "env"}) rootCmd.SetOut(&bytes.Buffer{}) @@ -89,7 +82,7 @@ func TestRootCommandDoesNotValidateClientFlagsForIndependentCommand(t *testing.T } func TestRootCommandInvokeContract(t *testing.T) { - rootCmd := newRootCommand(t.Context(), defaultRootOptions()) + rootCmd := cli.Root() assert.True(t, rootCmd.SilenceErrors) assert.True(t, rootCmd.SilenceUsage) @@ -111,7 +104,7 @@ func TestRootCommandInvokeContract(t *testing.T) { } func TestRootCommandV2CatalogAndLifecycleContract(t *testing.T) { - rootCmd := newRootCommand(t.Context(), defaultRootOptions()) + rootCmd := cli.Root() getTemplateCmd, _, err := rootCmd.Find([]string{"get", "agent-template"}) require.NoError(t, err) @@ -138,7 +131,7 @@ func TestRootCommandV2CatalogAndLifecycleContract(t *testing.T) { } func TestRootCommandRemovesLegacyPaths(t *testing.T) { - rootCmd := newRootCommand(t.Context(), defaultRootOptions()) + rootCmd := cli.Root() rootCommands := make([]string, 0, len(rootCmd.Commands())) for _, command := range rootCmd.Commands() { @@ -161,7 +154,7 @@ func TestRootCommandRemovesLegacyPaths(t *testing.T) { } func TestRootCommandRequiresTerminalForInteractiveUse(t *testing.T) { - rootCmd := newRootCommand(t.Context(), defaultRootOptions()) + rootCmd := cli.Root() rootCmd.SetArgs(nil) rootCmd.SetIn(&bytes.Buffer{}) rootCmd.SetOut(&bytes.Buffer{}) @@ -172,3 +165,45 @@ func TestRootCommandRequiresTerminalForInteractiveUse(t *testing.T) { assert.Contains(t, err.Error(), "kagent requires a terminal") assert.Contains(t, err.Error(), "kagent invoke") } + +func TestRootCommandOutputFormatReachesResourceCommands(t *testing.T) { + // An unparseable format is rejected before any command connects, so this + // reaches the run function without touching the network or a cluster. + for name, args := range map[string][]string{ + "get agent-instance": {"get", "agent-instance"}, + "get agent-template": {"get", "agent-template"}, + "create agent-instance": {"create", "agent-instance", "--harness", "kagent", "--agent-template", "example"}, + "delete agent-instance": {"delete", "agent-instance", "8bd650a8-9775-488f-8bc1-0d52bf7bdcab"}, + "invoke": {"invoke", "--agent-instance", "8bd650a8-9775-488f-8bc1-0d52bf7bdcab", "--task", "hello"}, + } { + t.Run(name, func(t *testing.T) { + rootCmd := cli.Root() + rootCmd.SetArgs(append(args, "--output-format", "bogus")) + rootCmd.SetOut(&bytes.Buffer{}) + rootCmd.SetErr(&bytes.Buffer{}) + + err := rootCmd.ExecuteContext(t.Context()) + + require.Error(t, err) + assert.Contains(t, err.Error(), `unsupported output format "bogus"`) + }) + } +} + +func TestRootResourceGroupsNameAvailableTypes(t *testing.T) { + for name, want := range map[string]string{ + "get": "agent-instance, agent-template", + "create": "agent-instance", + "delete": "agent-instance", + } { + t.Run(name, func(t *testing.T) { + rootCmd := cli.Root() + rootCmd.SetArgs([]string{name}) + + err := rootCmd.ExecuteContext(t.Context()) + + require.Error(t, err) + assert.Contains(t, err.Error(), "available resource types: "+want) + }) + } +} diff --git a/go/core/v2/agentplugins/materialize_test.go b/go/core/v2/agentplugins/materialize_test.go index c76dbd927..52ac0e43a 100644 --- a/go/core/v2/agentplugins/materialize_test.go +++ b/go/core/v2/agentplugins/materialize_test.go @@ -70,8 +70,12 @@ func TestFetchSourceReusesExistingMaterialization(t *testing.T) { if err != nil { t.Fatalf("fetchSource() redownloaded existing materialization: %v", err) } - if root != destination { - t.Fatalf("fetchSource() root = %q, want %q", root, destination) + wantRoot, err := filepath.EvalSymlinks(destination) + if err != nil { + t.Fatal(err) + } + if root != wantRoot { + t.Fatalf("fetchSource() root = %q, want %q", root, wantRoot) } }