Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions internal/devserver/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down
30 changes: 30 additions & 0 deletions internal/temporalcli/commands.gen.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
55 changes: 55 additions & 0 deletions internal/temporalcli/commands.nexus_operation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down
210 changes: 209 additions & 1 deletion internal/temporalcli/commands.nexus_operation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package temporalcli_test
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"testing"
Expand All @@ -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"
Expand Down Expand Up @@ -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, &notFound)
}

func (s *SharedServerSuite) TestNexusOperationList() {
endpointName, w := s.setupNexusEndpointAndWorker(s.T())
defer w.Stop()
Expand Down Expand Up @@ -793,4 +1002,3 @@ func (s *SharedServerSuite) TestNexusOperationStart_InvalidSearchAttribute() {
s.Error(res.Err)
s.ErrorContains(res.Err, "invalid search attribute")
}

Loading
Loading