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 {