This is an automated email from the ASF dual-hosted git repository.

lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git


The following commit(s) were added to refs/heads/rocketmq-studio by this push:
     new 7de310121 fix(rmqctl): handle Ctrl-C as a user interrupt instead of a 
server failure (#4944)
7de310121 is described below

commit 7de310121c50688abd82422880ef4c3817cc4a35
Author: Apulupie <[email protected]>
AuthorDate: Thu Oct 1 17:26:04 2026 +0800

    fix(rmqctl): handle Ctrl-C as a user interrupt instead of a server failure 
(#4944)
    
    * fix(rmqctl): report a canceled context as CANCELED instead of UNAVAILABLE
    
    * fix(rmqctl): abort the confirmation prompt when the command is interrupted
---
 rmqctl/cmd/app.go                  |  14 +++--
 rmqctl/cmd/catalog.go              |  72 ++++++++++++++++-------
 rmqctl/cmd/catalog_confirm_test.go | 114 +++++++++++++++++++++++++++++++++++++
 rmqctl/cmd/error.go                |  15 +++++
 rmqctl/cmd/error_test.go           |  89 +++++++++++++++++++++++++++++
 rmqctl/cmd/testhelpers_test.go     |   5 +-
 rmqctl/internal/types/types.go     |   1 +
 7 files changed, 283 insertions(+), 27 deletions(-)

diff --git a/rmqctl/cmd/app.go b/rmqctl/cmd/app.go
index 50968a202..169e2018f 100644
--- a/rmqctl/cmd/app.go
+++ b/rmqctl/cmd/app.go
@@ -44,12 +44,14 @@ type App struct {
 }
 
 // confirmFunc is the interactive confirmation hook for dangerous operations.
-// It receives the tool command path, risk level and target server, writes the
-// prompt to out, reads one line from in, and returns nil only when the user
-// explicitly approves. A non-nil error aborts the command before any HTTP
-// request is sent. Production code leaves App.confirm nil so the default
-// TTY-aware implementation in catalog.go is used; tests inject a stub.
-type confirmFunc func(in io.Reader, out io.Writer, commandPath, riskLevel, 
server string) error
+// It receives the command context, the tool command path, risk level and 
target
+// server, writes the prompt to out, reads one line from in, and returns nil
+// only when the user explicitly approves. A non-nil error aborts the command
+// before any HTTP request is sent. The context lets the implementation abort
+// the prompt when the command is canceled (Ctrl-C) while it waits for input.
+// Production code leaves App.confirm nil so the default TTY-aware
+// implementation in catalog.go is used; tests inject a stub.
+type confirmFunc func(ctx context.Context, in io.Reader, out io.Writer, 
commandPath, riskLevel, server string) error
 
 func NewApp(out io.Writer, err io.Writer) *App {
        return &App{
diff --git a/rmqctl/cmd/catalog.go b/rmqctl/cmd/catalog.go
index e29069c69..0f230cb51 100644
--- a/rmqctl/cmd/catalog.go
+++ b/rmqctl/cmd/catalog.go
@@ -18,6 +18,7 @@ package cmd
 
 import (
        "bufio"
+       "context"
        "fmt"
        "io"
        "os"
@@ -318,18 +319,23 @@ func confirmRisk(cmd *cobra.Command, runtime 
commandRuntime, tool toolcatalog.To
        if confirm == nil {
                confirm = defaultConfirm
        }
-       return confirm(cmd.InOrStdin(), cmd.ErrOrStderr(), tool.CommandPath(), 
tool.RiskLevel, runtime.targetServer())
+       return confirm(cmd.Context(), cmd.InOrStdin(), cmd.ErrOrStderr(), 
tool.CommandPath(), tool.RiskLevel, runtime.targetServer())
 }
 
 // defaultConfirm is the production confirmation implementation. It only
 // prompts when stdin is a character device (interactive terminal); pipes and
-// redirected input are rejected so scripts must pass --yes explicitly.
-func defaultConfirm(in io.Reader, out io.Writer, commandPath, riskLevel, 
server string) error {
-       file, ok := in.(*os.File)
-       if !ok || file == nil {
+// redirected input are rejected so scripts must pass --yes explicitly. The
+// command context is watched while waiting for the answer so a Ctrl-C aborts
+// the prompt instead of leaving it open for a "yes" that would run the
+// mutation with an already-canceled context.
+func defaultConfirm(ctx context.Context, in io.Reader, out io.Writer, 
commandPath, riskLevel, server string) error {
+       // Anything that can report its own mode qualifies as a potential 
terminal: production passes
+       // *os.File, and the character-device check below is what actually 
gates interactive use.
+       statter, ok := in.(interface{ Stat() (os.FileInfo, error) })
+       if !ok || statter == nil {
                return riskConfirmationRequired(commandPath, riskLevel)
        }
-       stat, err := file.Stat()
+       stat, err := statter.Stat()
        if err != nil {
                return riskConfirmationRequired(commandPath, riskLevel)
        }
@@ -338,21 +344,49 @@ func defaultConfirm(in io.Reader, out io.Writer, 
commandPath, riskLevel, server
        }
        fmt.Fprintf(out, "WARNING: %q is a %s operation.\n", commandPath, 
riskLevel)
        fmt.Fprintf(out, "Arguments will be sent to %s. Type \"yes\" to 
continue: ", server)
-       reader := bufio.NewReader(in)
-       answer, err := reader.ReadString('\n')
-       if err != nil {
-               return types.NewCLIError(
-                       types.CodeCommandFailed,
-                       fmt.Sprintf("confirmation for %q failed: %v", 
commandPath, err),
-                       "Re-run the command and confirm interactively, or pass 
--yes to skip the prompt.")
+
+       // The stdin read itself cannot be interrupted, so run it alongside the 
context: a canceled
+       // context always wins over an answer that may already have been typed, 
and the buffered
+       // channel lets the reader goroutine retire after a canceled prompt 
without leaking.
+       reads := make(chan confirmRead, 1)
+       go func() {
+               reader := bufio.NewReader(in)
+               answer, err := reader.ReadString('\n')
+               reads <- confirmRead{answer: answer, err: err}
+       }()
+       if ctx.Err() != nil {
+               return confirmInterrupted(commandPath)
        }
-       if !isAffirmative(strings.TrimSpace(answer)) {
-               return types.NewCLIError(
-                       types.CodeCommandFailed,
-                       fmt.Sprintf("execution of %q cancelled", commandPath),
-                       "Re-run the command and type yes, or pass --yes to skip 
the prompt.")
+       select {
+       case <-ctx.Done():
+               return confirmInterrupted(commandPath)
+       case read := <-reads:
+               if read.err != nil {
+                       return types.NewCLIError(
+                               types.CodeCommandFailed,
+                               fmt.Sprintf("confirmation for %q failed: %v", 
commandPath, read.err),
+                               "Re-run the command and confirm interactively, 
or pass --yes to skip the prompt.")
+               }
+               if !isAffirmative(strings.TrimSpace(read.answer)) {
+                       return types.NewCLIError(
+                               types.CodeCommandFailed,
+                               fmt.Sprintf("execution of %q cancelled", 
commandPath),
+                               "Re-run the command and type yes, or pass --yes 
to skip the prompt.")
+               }
+               return nil
        }
-       return nil
+}
+
+type confirmRead struct {
+       answer string
+       err    error
+}
+
+func confirmInterrupted(commandPath string) error {
+       return types.NewCLIError(
+               types.CodeCanceled,
+               fmt.Sprintf("confirmation for %q was interrupted", commandPath),
+               "The command was canceled (for example with Ctrl-C); rerun it 
if that was unintended.")
 }
 
 func riskConfirmationRequired(commandPath, riskLevel string) error {
diff --git a/rmqctl/cmd/catalog_confirm_test.go 
b/rmqctl/cmd/catalog_confirm_test.go
new file mode 100644
index 000000000..170522f11
--- /dev/null
+++ b/rmqctl/cmd/catalog_confirm_test.go
@@ -0,0 +1,114 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package cmd
+
+import (
+       "context"
+       "errors"
+       "io"
+       "os"
+       "strings"
+       "testing"
+       "time"
+
+       "github.com/apache/rocketmq-dashboard/rmqctl/internal/types"
+)
+
+// charDeviceReader wraps a pipe with a character-device mode bit so 
defaultConfirm treats it as
+// an interactive prompt without touching the test process's real stdin.
+type charDeviceReader struct{ r io.Reader }
+
+func (c charDeviceReader) Read(p []byte) (int, error) { return c.r.Read(p) }
+
+// Stat pretends to be a character device so the TTY gate passes.
+func (c charDeviceReader) Stat() (os.FileInfo, error) {
+       return charDeviceFileInfo{}, nil
+}
+
+type charDeviceFileInfo struct{}
+
+func (charDeviceFileInfo) Name() string       { return "stdin" }
+func (charDeviceFileInfo) Size() int64        { return 0 }
+func (charDeviceFileInfo) Mode() os.FileMode  { return os.ModeCharDevice | 
0o600 }
+func (charDeviceFileInfo) ModTime() time.Time { return time.Time{} }
+func (charDeviceFileInfo) IsDir() bool        { return false }
+func (charDeviceFileInfo) Sys() any           { return nil }
+
+func TestDefaultConfirmCanceledContextWinsOverTypedYes(t *testing.T) {
+       pipeR, pipeW, err := os.Pipe()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer pipeR.Close()
+       defer pipeW.Close()
+       // The operator's "yes" is already in the pipe: an interrupt must still 
abort the prompt,
+       // because running the mutation with a canceled context would only fail 
later and mislead.
+       if _, err := pipeW.WriteString("yes\n"); err != nil {
+               t.Fatal(err)
+       }
+       ctx, cancel := context.WithCancel(context.Background())
+       cancel()
+       err = defaultConfirm(ctx, charDeviceReader{r: pipeR}, io.Discard, 
"rmq.consumer.delete", "L3", "http://studio";)
+       var cliErr *types.CLIError
+       if !errors.As(err, &cliErr) || cliErr.Code != types.CodeCanceled {
+               t.Fatalf("defaultConfirm with canceled context: got %v, want 
CANCELED", err)
+       }
+}
+
+func TestDefaultConfirmInterruptDuringReadAborts(t *testing.T) {
+       pipeR, pipeW, err := os.Pipe()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer pipeR.Close()
+       defer pipeW.Close()
+       ctx, cancel := context.WithCancel(context.Background())
+       go func() {
+               // Simulate Ctrl-C arriving while the prompt waits on an empty 
pipe.
+               time.Sleep(50 * time.Millisecond)
+               cancel()
+       }()
+       err = defaultConfirm(ctx, charDeviceReader{r: pipeR}, io.Discard, 
"rmq.consumer.delete", "L3", "http://studio";)
+       cancel()
+       var cliErr *types.CLIError
+       if !errors.As(err, &cliErr) || cliErr.Code != types.CodeCanceled {
+               t.Fatalf("defaultConfirm interrupted mid-read: got %v, want 
CANCELED", err)
+       }
+}
+
+func TestDefaultConfirmAffirmativeProceeds(t *testing.T) {
+       pipeR, pipeW, err := os.Pipe()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer pipeR.Close()
+       defer pipeW.Close()
+       if _, err := pipeW.WriteString("yes\n"); err != nil {
+               t.Fatal(err)
+       }
+       if err := defaultConfirm(context.Background(), charDeviceReader{r: 
pipeR}, io.Discard, "rmq.consumer.delete", "L3", "http://studio";); err != nil {
+               t.Fatalf("defaultConfirm with yes: got %v, want nil", err)
+       }
+}
+
+func TestDefaultConfirmNonTTYRejected(t *testing.T) {
+       err := defaultConfirm(context.Background(), strings.NewReader("yes\n"), 
io.Discard, "rmq.consumer.delete", "L3", "http://studio";)
+       var cliErr *types.CLIError
+       if !errors.As(err, &cliErr) || cliErr.Code != types.CodeCommandFailed {
+               t.Fatalf("defaultConfirm on non-chardevice stdin: got %v, want 
COMMAND_FAILED rejection", err)
+       }
+}
diff --git a/rmqctl/cmd/error.go b/rmqctl/cmd/error.go
index 9f9821bff..de145bf40 100644
--- a/rmqctl/cmd/error.go
+++ b/rmqctl/cmd/error.go
@@ -62,6 +62,7 @@ var errorMatchers = []errorMatcher{
        matchCLIError,
        matchAPIError,
        matchContextDeadline,
+       matchContextCanceled,
        matchNetError,
 }
 
@@ -94,6 +95,20 @@ func matchContextDeadline(err error) (*types.CLIError, bool) 
{
                "Increase --timeout or check Studio Server availability."), true
 }
 
+// matchContextCanceled reports a canceled context as a user interrupt, not a 
server problem:
+// Ctrl-C (signal.NotifyContext) cancels the command context while a request 
is in flight, and a
+// wrapped *url.Error carrying context.Canceled would otherwise fall through 
to the net.Error
+// matcher and surface as UNAVAILABLE, sending the operator to debug a healthy 
server.
+func matchContextCanceled(err error) (*types.CLIError, bool) {
+       if !errors.Is(err, context.Canceled) {
+               return nil, false
+       }
+       return types.NewCLIError(
+               types.CodeCanceled,
+               "command was interrupted before it finished",
+               "The command was canceled (for example with Ctrl-C); rerun it 
if that was unintended."), true
+}
+
 func matchNetError(err error) (*types.CLIError, bool) {
        if _, ok := errors.AsType[net.Error](err); !ok {
                return nil, false
diff --git a/rmqctl/cmd/error_test.go b/rmqctl/cmd/error_test.go
new file mode 100644
index 000000000..32b147f2c
--- /dev/null
+++ b/rmqctl/cmd/error_test.go
@@ -0,0 +1,89 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package cmd
+
+import (
+       "context"
+       "errors"
+       "net"
+       "net/url"
+       "testing"
+
+       "github.com/apache/rocketmq-dashboard/rmqctl/internal/types"
+)
+
+func TestNormalizeCLIError(t *testing.T) {
+       tests := []struct {
+               name string
+               err  error
+               want string
+       }{
+               {
+                       // Ctrl-C cancels the command context mid-request; the 
transport wraps it in a
+                       // *url.Error, which is a net.Error. It must surface as 
CANCELED, not UNAVAILABLE.
+                       name: "interrupted request reports canceled",
+                       err: &url.Error{
+                               Op:  "Post",
+                               URL: "http://studio.example/api";,
+                               Err: context.Canceled,
+                       },
+                       want: types.CodeCanceled,
+               },
+               {
+                       name: "bare context canceled reports canceled",
+                       err:  context.Canceled,
+                       want: types.CodeCanceled,
+               },
+               {
+                       name: "deadline exceeded reports timeout",
+                       err: &url.Error{
+                               Op:  "Get",
+                               URL: "http://studio.example/api";,
+                               Err: context.DeadlineExceeded,
+                       },
+                       want: types.CodeTimeout,
+               },
+               {
+                       name: "connection refused reports unavailable",
+                       err: &url.Error{
+                               Op:  "Get",
+                               URL: "http://studio.example/api";,
+                               Err: &net.OpError{Op: "dial", Err: 
errors.New("connection refused")},
+                       },
+                       want: types.CodeUnavailable,
+               },
+               {
+                       name: "plain error reports command failed",
+                       err:  errors.New("something else failed"),
+                       want: types.CodeCommandFailed,
+               },
+       }
+
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       got := normalizeCLIError(tt.err)
+                       if got.Code != tt.want {
+                               t.Fatalf("normalizeCLIError(%v).Code = %q, want 
%q", tt.err, got.Code, tt.want)
+                       }
+                       if got.Code == types.CodeCanceled {
+                               if got.Hint == "" || got.Hint == "Check 
--server, network connectivity, and Studio Server status." {
+                                       t.Fatalf("canceled hint must point at 
the interrupt, got %q", got.Hint)
+                               }
+                       }
+               })
+       }
+}
diff --git a/rmqctl/cmd/testhelpers_test.go b/rmqctl/cmd/testhelpers_test.go
index 2631529d7..dd96322e3 100644
--- a/rmqctl/cmd/testhelpers_test.go
+++ b/rmqctl/cmd/testhelpers_test.go
@@ -19,6 +19,7 @@ package cmd
 import (
        "bufio"
        "bytes"
+       "context"
        "encoding/json"
        "fmt"
        "io"
@@ -132,7 +133,7 @@ func executeTestAppWithStdin(t *testing.T, client 
*http.Client, serverURL, insta
        app := NewApp(stdout, stderr)
        app.HTTP = client
        app.Store.Getenv = testEnv
-       app.confirm = func(in io.Reader, out io.Writer, commandPath, riskLevel, 
server string) error {
+       app.confirm = func(_ context.Context, in io.Reader, out io.Writer, 
commandPath, riskLevel, server string) error {
                fmt.Fprintf(out, "WARNING: %q is a %s operation.\n", 
commandPath, riskLevel)
                fmt.Fprintf(out, "Arguments will be sent to %s. Type \"yes\" to 
continue: ", server)
                reader := bufio.NewReader(in)
@@ -164,7 +165,7 @@ func executeTestAppWithStdin(t *testing.T, client 
*http.Client, serverURL, insta
 
 // stubConfirmReject mimics the production non-TTY rejection: it always returns
 // the "requires interactive confirmation" error without reading stdin.
-func stubConfirmReject(in io.Reader, out io.Writer, commandPath, riskLevel, 
server string) error {
+func stubConfirmReject(_ context.Context, in io.Reader, out io.Writer, 
commandPath, riskLevel, server string) error {
        return types.NewCLIError(
                types.CodeCommandFailed,
                fmt.Sprintf("%q is a %s operation and requires interactive 
confirmation", commandPath, riskLevel),
diff --git a/rmqctl/internal/types/types.go b/rmqctl/internal/types/types.go
index 5f6f59578..3ce36e7df 100644
--- a/rmqctl/internal/types/types.go
+++ b/rmqctl/internal/types/types.go
@@ -64,6 +64,7 @@ const (
        CodeTimeout         = "TIMEOUT"
        CodeUnavailable     = "UNAVAILABLE"
        CodeCommandFailed   = "COMMAND_FAILED"
+       CodeCanceled        = "CANCELED"
 )
 
 type CLIError struct {

Reply via email to