diff --git a/internal/devserver/server.go b/internal/devserver/server.go index 09d47854f..dcd09031f 100644 --- a/internal/devserver/server.go +++ b/internal/devserver/server.go @@ -40,6 +40,7 @@ import ( uiserveroptions "github.com/temporalio/ui-server/v2/server/server_options" "go.temporal.io/api/enums/v1" "go.temporal.io/server/chasm/lib/activity" + "go.temporal.io/server/chasm/lib/nexusoperation" "go.temporal.io/server/common/authorization" "go.temporal.io/server/common/cluster" "go.temporal.io/server/common/config" @@ -250,6 +251,9 @@ func (s *StartOptions) buildServerOptions() ([]temporal.ServerOption, *slog.Leve dynConf[activity.EnableStandaloneActivityOperatorCommands.Key()] = true dynConf[dynamicconfig.FrontendEnableBatchOperationsForStandaloneActivities.Key()] = true + // Enable SANO; this can be removed when SANO is on by default. + dynConf[nexusoperation.Enabled.Key()] = true + // Dynamic config if set for k, v := range s.DynamicConfigValues { dynConf[dynamicconfig.MakeKey(k)] = v diff --git a/internal/temporalcli/commands.gen.go b/internal/temporalcli/commands.gen.go index a48152714..3f625ab05 100644 --- a/internal/temporalcli/commands.gen.go +++ b/internal/temporalcli/commands.gen.go @@ -1523,6 +1523,7 @@ func NewTemporalNexusOperationCommand(cctx *CommandContext, parent *TemporalNexu s.Command.Args = cobra.NoArgs s.Command.AddCommand(&NewTemporalNexusOperationCancelCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalNexusOperationCountCommand(cctx, &s).Command) + s.Command.AddCommand(&NewTemporalNexusOperationDeleteCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalNexusOperationDescribeCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalNexusOperationExecuteCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalNexusOperationListCommand(cctx, &s).Command) @@ -1588,6 +1589,35 @@ func NewTemporalNexusOperationCountCommand(cctx *CommandContext, parent *Tempora return &s } +type TemporalNexusOperationDeleteCommand struct { + Parent *TemporalNexusOperationCommand + Command cobra.Command + NexusOperationReferenceOptions + Yes bool +} + +func NewTemporalNexusOperationDeleteCommand(cctx *CommandContext, parent *TemporalNexusOperationCommand) *TemporalNexusOperationDeleteCommand { + var s TemporalNexusOperationDeleteCommand + s.Parent = parent + s.Command.DisableFlagsInUseLine = true + s.Command.Use = "delete [flags]" + s.Command.Short = "Remove a Nexus Operation Execution (Experimental)" + if hasHighlighting { + s.Command.Long = "Delete a Nexus Operation Execution and its history. The request runs\nasynchronously. If the Operation is running, the Service terminates it\nbefore deletion. Without a Run ID, the latest run is deleted.\n\n\x1b[1mtemporal nexus operation delete \\\n --operation-id YourOperationId \\\n --yes\x1b[0m\n\nUse \"--run-id\" to delete a specific run. Confirmation is required unless\n\"--yes\" is passed." + } else { + s.Command.Long = "Delete a Nexus Operation Execution and its history. The request runs\nasynchronously. If the Operation is running, the Service terminates it\nbefore deletion. Without a Run ID, the latest run is deleted.\n\n```\ntemporal nexus operation delete \\\n --operation-id YourOperationId \\\n --yes\n```\n\nUse \"--run-id\" to delete a specific run. Confirmation is required unless\n\"--yes\" is passed." + } + s.Command.Args = cobra.NoArgs + s.Command.Flags().BoolVarP(&s.Yes, "yes", "y", false, "Don't prompt to confirm deletion.") + s.NexusOperationReferenceOptions.BuildFlags(s.Command.Flags()) + s.Command.Run = func(c *cobra.Command, args []string) { + if err := s.run(cctx, args); err != nil { + cctx.Options.Fail(err) + } + } + return &s +} + type TemporalNexusOperationDescribeCommand struct { Parent *TemporalNexusOperationCommand Command cobra.Command diff --git a/internal/temporalcli/commands.nexus_operation.go b/internal/temporalcli/commands.nexus_operation.go index b289e4456..c97a39703 100644 --- a/internal/temporalcli/commands.nexus_operation.go +++ b/internal/temporalcli/commands.nexus_operation.go @@ -18,6 +18,8 @@ import ( "go.temporal.io/sdk/converter" ) +const nexusOperationDeleteWarning = "WARNING: Deleting Nexus Operation Executions in a global Namespace removes them from all replicas. Requests sent to a passive cluster are forwarded to the active cluster by default; to target the passive cluster directly, specify `--grpc-meta xdc-redirection=false`." + func (c *TemporalNexusOperationStartCommand) run(cctx *CommandContext, args []string) error { cl, err := dialClient(cctx, &c.Parent.Parent.ClientOptions) if err != nil { @@ -352,6 +354,59 @@ func (c *TemporalNexusOperationTerminateCommand) run(cctx *CommandContext, args return nil } +func (c *TemporalNexusOperationDeleteCommand) run(cctx *CommandContext, _ []string) error { + cl, err := dialClient(cctx, &c.Parent.Parent.ClientOptions) + if err != nil { + return err + } + defer cl.Close() + + // Only warn when the namespace is global, or can't get the namespace info + nsResp, nsErr := cl.WorkflowService().DescribeNamespace(cctx, &workflowservice.DescribeNamespaceRequest{ + Namespace: c.Parent.Parent.Namespace, + }) + if nsErr != nil || nsResp.GetIsGlobalNamespace() { + fmt.Fprintln(cctx.Options.Stderr, nexusOperationDeleteWarning) + } + + yes, err := cctx.promptYes(nexusOperationDeleteConfirmationMessage(c.OperationId, c.RunId), c.Yes) + if err != nil { + return err + } else if !yes { + return fmt.Errorf("user denied confirmation") + } + + _, err = cl.WorkflowService().DeleteNexusOperationExecution(cctx, &workflowservice.DeleteNexusOperationExecutionRequest{ + Namespace: c.Parent.Parent.Namespace, + OperationId: c.OperationId, + RunId: c.RunId, + }) + if err != nil { + return fmt.Errorf("failed to delete nexus operation: %w", err) + } + if cctx.JSONOutput { + return cctx.Printer.PrintStructured(struct { + OperationId string `json:"operationId"` + RunId string `json:"runId,omitempty"` + Status string `json:"status"` + }{ + OperationId: c.OperationId, + RunId: c.RunId, + Status: "DELETE_REQUESTED", + }, printer.StructuredOptions{}) + } + cctx.Printer.Println("Nexus Operation deletion requested") + return nil +} + +func nexusOperationDeleteConfirmationMessage(operationID, runID string) string { + action := fmt.Sprintf("Delete Nexus Operation %q", operationID) + if runID != "" { + action += fmt.Sprintf(" with Run ID %q", runID) + } + return fmt.Sprintf("%s? y/N", action) +} + func (c *TemporalNexusOperationListCommand) run(cctx *CommandContext, _ []string) error { cl, err := dialClient(cctx, &c.Parent.Parent.ClientOptions) if err != nil { diff --git a/internal/temporalcli/commands.nexus_operation_test.go b/internal/temporalcli/commands.nexus_operation_test.go index 0418c22f4..71fc7e2db 100644 --- a/internal/temporalcli/commands.nexus_operation_test.go +++ b/internal/temporalcli/commands.nexus_operation_test.go @@ -3,6 +3,7 @@ package temporalcli_test import ( "context" "encoding/json" + "errors" "fmt" "strings" "testing" @@ -13,6 +14,7 @@ import ( "github.com/stretchr/testify/require" nexuspb "go.temporal.io/api/nexus/v1" "go.temporal.io/api/operatorservice/v1" + "go.temporal.io/api/serviceerror" "go.temporal.io/sdk/client" "go.temporal.io/sdk/temporalnexus" "go.temporal.io/sdk/workflow" @@ -270,6 +272,213 @@ func (s *SharedServerSuite) TestNexusOperationTerminate() { s.Contains(res.Stdout.String(), "Nexus Operation terminated") } +func (s *SharedServerSuite) TestNexusOperationDelete() { + endpointName, w := s.setupNexusEndpointAndWorker(s.T()) + defer w.Stop() + + opID := "delete-op-" + uuid.NewString()[:8] + res := s.Execute( + "nexus", "operation", "start", + "--address", s.Address(), + "--endpoint", endpointName, + "--service", "test-service", + "--operation", "test-op", + "--operation-id", opID, + "--input", `"hello"`, + ) + s.NoError(res.Err) + + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + ) + return res.Err == nil + }, 30*time.Second, 500*time.Millisecond) + + res = s.Execute( + "nexus", "operation", "delete", + "--address", s.Address(), + "--operation-id", opID, + ) + s.EqualError(res.Err, "user denied confirmation") + s.Contains(res.Stdout.String(), opID) + + res = s.Execute( + "nexus", "operation", "delete", + "--address", s.Address(), + "--operation-id", opID, + "--yes", + ) + s.NoError(res.Err) + s.Contains(res.Stdout.String(), "Nexus Operation deletion requested") + + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + ) + return isNotFoundErr(res.Err) + }, 30*time.Second, 500*time.Millisecond) +} + +func (s *SharedServerSuite) TestNexusOperationDelete_WithoutRunIDDeletesLatestRun() { + endpointName, w := s.setupNexusEndpointAndWorker(s.T()) + defer w.Stop() + + opID := "delete-latest-op-" + uuid.NewString()[:8] + startOperation := func(extraArgs ...string) string { + args := []string{ + "nexus", "operation", "start", + "--address", s.Address(), + "--endpoint", endpointName, + "--service", "test-service", + "--operation", "test-op", + "--operation-id", opID, + "--input", `"hello"`, + "--output", "json", + } + res := s.Execute(append(args, extraArgs...)...) + s.NoError(res.Err) + var started struct { + RunId string `json:"runId"` + } + s.NoError(json.Unmarshal(res.Stdout.Bytes(), &started)) + s.NotEmpty(started.RunId) + return started.RunId + } + + firstRunID := startOperation() + + // Wait for the first run to close so the same Operation ID can be reused. + res := s.Execute( + "nexus", "operation", "result", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", firstRunID, + "--output", "json", + ) + s.NoError(res.Err) + + latestRunID := startOperation("--id-reuse-policy", "AllowDuplicate") + s.NotEqual(firstRunID, latestRunID) + + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", latestRunID, + ) + return res.Err == nil + }, 30*time.Second, 500*time.Millisecond) + + res = s.Execute( + "nexus", "operation", "delete", + "--address", s.Address(), + "--operation-id", opID, + "--yes", + ) + s.NoError(res.Err) + s.Contains(res.Stdout.String(), "Nexus Operation deletion requested") + + // Omitting --run-id deletes the latest run, leaving the earlier run intact. + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", latestRunID, + ) + return isNotFoundErr(res.Err) + }, 30*time.Second, 500*time.Millisecond) + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", firstRunID, + ) + return res.Err == nil + }, 30*time.Second, 500*time.Millisecond) +} + +func (s *SharedServerSuite) TestNexusOperationDelete_RunID_JSON() { + endpointName, w := s.setupNexusEndpointAndWorker(s.T()) + defer w.Stop() + + opID := "delete-run-op-" + uuid.NewString()[:8] + res := s.Execute( + "nexus", "operation", "start", + "--address", s.Address(), + "--endpoint", endpointName, + "--service", "test-service", + "--operation", "test-op", + "--operation-id", opID, + "--input", `"hello"`, + "--output", "json", + ) + s.NoError(res.Err) + var started struct { + RunId string `json:"runId"` + } + s.NoError(json.Unmarshal(res.Stdout.Bytes(), &started)) + s.NotEmpty(started.RunId) + + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", started.RunId, + ) + return res.Err == nil + }, 30*time.Second, 500*time.Millisecond) + + res = s.Execute( + "nexus", "operation", "delete", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", started.RunId, + "--yes", + "--output", "json", + ) + s.NoError(res.Err) + var deleted struct { + OperationId string `json:"operationId"` + RunId string `json:"runId"` + Status string `json:"status"` + } + s.NoError(json.Unmarshal(res.Stdout.Bytes(), &deleted)) + s.Equal(opID, deleted.OperationId) + s.Equal(started.RunId, deleted.RunId) + s.Equal("DELETE_REQUESTED", deleted.Status) + + s.Eventually(func() bool { + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + "--run-id", started.RunId, + ) + return isNotFoundErr(res.Err) + }, 30*time.Second, 500*time.Millisecond) +} + +func (s *SharedServerSuite) TestNexusOperationDelete_MissingOperationID() { + res := s.Execute("nexus", "operation", "delete", "--address", s.Address(), "--yes") + s.ErrorContains(res.Err, "operation-id") +} + +// isNotFoundErr reports whether err is the server's NotFound, so deletion waits +// don't treat a transient RPC failure as proof the execution is gone. +func isNotFoundErr(err error) bool { + var notFound *serviceerror.NotFound + return errors.As(err, ¬Found) +} + func (s *SharedServerSuite) TestNexusOperationList() { endpointName, w := s.setupNexusEndpointAndWorker(s.T()) defer w.Stop() @@ -793,4 +1002,3 @@ func (s *SharedServerSuite) TestNexusOperationStart_InvalidSearchAttribute() { s.Error(res.Err) s.ErrorContains(res.Err, "invalid search attribute") } - diff --git a/internal/temporalcli/commands.yaml b/internal/temporalcli/commands.yaml index 3a94ce74b..fd74ba530 100644 --- a/internal/temporalcli/commands.yaml +++ b/internal/temporalcli/commands.yaml @@ -2200,6 +2200,7 @@ commands: - nexus operation describe - nexus operation cancel - nexus operation terminate + - nexus operation delete - nexus operation list - nexus operation count - cli reference @@ -2314,6 +2315,29 @@ commands: Reason for termination. Defaults to a message with the current user's name. + - name: temporal nexus operation delete + summary: Remove a Nexus Operation Execution (Experimental) + description: | + Delete a Nexus Operation Execution and its history. The request runs + asynchronously. If the Operation is running, the Service terminates it + before deletion. Without a Run ID, the latest run is deleted. + + ``` + temporal nexus operation delete \ + --operation-id YourOperationId \ + --yes + ``` + + Use "--run-id" to delete a specific run. Confirmation is required unless + "--yes" is passed. + option-sets: + - nexus-operation-reference + options: + - name: yes + type: bool + short: y + description: Don't prompt to confirm deletion. + - name: temporal nexus operation list summary: List Nexus Operations matching a query (Experimental) description: | diff --git a/internal/temporalcli/commands_test.go b/internal/temporalcli/commands_test.go index 6b621ef18..af9ef6ebc 100644 --- a/internal/temporalcli/commands_test.go +++ b/internal/temporalcli/commands_test.go @@ -406,6 +406,7 @@ func StartDevServer(t *testing.T, options DevServerOptions) *DevServer { d.Options.DynamicConfigValues = map[string]any{} } d.Options.DynamicConfigValues["system.forceSearchAttributesCacheRefreshOnRead"] = true + d.Options.DynamicConfigValues["system.forceNexusEndpointRefreshOnRead"] = true d.Options.DynamicConfigValues["frontend.workerVersioningRuleAPIs"] = true d.Options.DynamicConfigValues["frontend.workerVersioningDataAPIs"] = true d.Options.DynamicConfigValues["frontend.workerVersioningWorkflowAPIs"] = true