From 234498b85a2338c9c369b9890da745bc4928e601 Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Fri, 25 Sep 2026 18:50:01 +0800 Subject: [PATCH] =?UTF-8?q?fix(rcon):=20=E5=A4=9A=E5=8C=85=E5=9B=9E?= =?UTF-8?q?=E5=A4=8D=E5=9C=A8=E9=A6=96=E5=8C=85=E5=88=B0=E8=BE=BE=E5=90=8E?= =?UTF-8?q?=E5=8F=91=E7=BB=93=E6=9D=9F=E6=A0=87=E8=AE=B0=E9=87=8D=E7=BB=84?= =?UTF-8?q?=EF=BC=8C=E5=8D=95=E5=8C=85=E4=B8=8A=E9=99=90=E6=94=BE=E5=AE=BD?= =?UTF-8?q?=E5=88=B0=204106=EF=BC=8C=E5=91=BD=E4=BB=A4=E5=B8=A6=20ctx=20?= =?UTF-8?q?=E4=B8=8E=E9=BB=98=E8=AE=A4=2010s=20=E8=B6=85=E6=97=B6=EF=BC=8C?= =?UTF-8?q?console=20=E8=B5=B0=20ExecuteContext?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/api/console.go | 2 +- internal/rcon/rcon.go | 106 ++++++++++++++++++++++++++------- internal/rcon/rcon_test.go | 116 ++++++++++++++++++++++++++++++++++++- 3 files changed, 202 insertions(+), 22 deletions(-) diff --git a/internal/api/console.go b/internal/api/console.go index 6722732..429e3f5 100644 --- a/internal/api/console.go +++ b/internal/api/console.go @@ -101,7 +101,7 @@ func (k *K8sConsole) RunCommand(ctx context.Context, name, command string) (stri } defer conn.Close() - out, err := conn.Execute(command) + out, err := conn.ExecuteContext(ctx, command) if err != nil { return "", ErrConsoleUnavailable } diff --git a/internal/rcon/rcon.go b/internal/rcon/rcon.go index 1c771de..e7490ef 100644 --- a/internal/rcon/rcon.go +++ b/internal/rcon/rcon.go @@ -6,13 +6,15 @@ // - Graceful shutdown: Execute("save-all flush") right before the operator // scales a server to zero. // -// Multi-packet responses (a single command whose reply exceeds one ~4 KiB -// packet) are not reassembled; Phase-1 commands ("list", "save-all", "stop") -// always fit in one packet. This is intentional and documented rather than -// silently truncating large replies. +// A reply longer than one packet (Minecraft splits a command's output into +// 4096-byte bodies: banlist, whitelist list, any console command) is +// reassembled: once the reply begins, Execute sends an empty RESPONSE_VALUE +// packet, which the server answers only after the whole reply ("Unknown +// request 0" on vanilla and Paper), so that answer marks the end. package rcon import ( + "context" "encoding/binary" "errors" "fmt" @@ -36,8 +38,18 @@ const authFailedID int32 = -1 const ( minPacketLen = 10 maxPacketLen = 4096 + // maxReplyPacketLen admits a full 4096-byte Minecraft reply body plus the + // id, type and two terminators. + maxReplyPacketLen = maxPacketLen + minPacketLen + // maxReplyBytes bounds one reassembled reply, so a server that never sends + // the end marker cannot grow it without limit. + maxReplyBytes = 1 << 20 ) +// DefaultCommandTimeout bounds a command when neither the caller's context nor +// SetDeadline gives one, so a hung server cannot hold a console request open. +const DefaultCommandTimeout = 10 * time.Second + // ErrAuthFailed is returned by Dial when the RCON password is rejected. var ErrAuthFailed = errors.New("rcon: authentication failed") @@ -48,6 +60,9 @@ const DefaultPort = 25575 type Conn struct { conn net.Conn reqID int32 + // deadline is the caller's SetDeadline, which a command honours in place + // of DefaultCommandTimeout. + deadline time.Time } // Dial opens a TCP connection to addr and authenticates with password. The @@ -87,25 +102,76 @@ func Dial(addr, password string, timeout time.Duration) (*Conn, error) { } // SetDeadline sets an absolute deadline for subsequent Execute calls. -func (c *Conn) SetDeadline(t time.Time) error { return c.conn.SetDeadline(t) } +func (c *Conn) SetDeadline(t time.Time) error { + c.deadline = t + return c.conn.SetDeadline(t) +} // Close closes the underlying connection. func (c *Conn) Close() error { return c.conn.Close() } -// Execute runs a single command and returns the server's reply body. +// Execute runs a single command and returns the server's whole reply. func (c *Conn) Execute(cmd string) (string, error) { - id := c.nextID() + return c.ExecuteContext(context.Background(), cmd) +} + +// ExecuteContext runs a single command and returns the server's whole reply, +// reassembled across packets. It gives up when ctx is done (cancelled or past +// its deadline), and otherwise at the SetDeadline deadline, or after +// DefaultCommandTimeout when none was set. +func (c *Conn) ExecuteContext(ctx context.Context, cmd string) (string, error) { + dl := c.deadline + if dl.IsZero() { + dl = time.Now().Add(DefaultCommandTimeout) + } + if err := c.conn.SetDeadline(dl); err != nil { + return "", err + } + stop := context.AfterFunc(ctx, func() { _ = c.conn.SetDeadline(time.Unix(1, 0)) }) + defer stop() + + id, end := c.nextID(), c.nextID() if err := writePacket(c.conn, id, typeExecCommand, cmd); err != nil { - return "", err + return "", c.ctxErr(ctx, err) } - respID, _, body, err := readPacket(c.conn) - if err != nil { - return "", err + var reply []byte + sentEnd := false + for { + respID, _, body, err := readPacket(c.conn) + if err != nil { + return "", c.ctxErr(ctx, err) + } + switch { + case respID == id: + if len(reply)+len(body) > maxReplyBytes { + return "", fmt.Errorf("rcon: reply exceeds %d bytes", maxReplyBytes) + } + reply = append(reply, body...) + // The end marker goes out only once the reply has begun: Minecraft + // drops a connection whose read holds more than one packet, and it + // sends every fragment before it reads again, so the marker's answer + // follows the last one. + if !sentEnd { + if err := writePacket(c.conn, end, typeResponseValue, ""); err != nil { + return "", c.ctxErr(ctx, err) + } + sentEnd = true + } + case respID == end: + return string(reply), nil + default: + // A leftover from an earlier exchange (a server that answers the end + // marker twice); it belongs to no live request. + } } - if respID != id { - return "", fmt.Errorf("rcon: response id mismatch: got %d want %d", respID, id) +} + +// ctxErr prefers ctx's error when a cancellation is what cut the exchange. +func (c *Conn) ctxErr(ctx context.Context, err error) error { + if ctxErr := ctx.Err(); ctxErr != nil { + return ctxErr } - return body, nil + return err } // auth performs the SERVERDATA_AUTH handshake. Some servers emit an empty @@ -161,23 +227,23 @@ func writePacket(w io.Writer, id, typ int32, body string) error { } // readPacket decodes one RCON packet. -func readPacket(r io.Reader) (id, typ int32, body string, err error) { +func readPacket(r io.Reader) (id, typ int32, body []byte, err error) { var lenBuf [4]byte if _, err = io.ReadFull(r, lenBuf[:]); err != nil { - return 0, 0, "", err + return 0, 0, nil, err } length := int32(binary.LittleEndian.Uint32(lenBuf[:])) - if length < minPacketLen || length > maxPacketLen { - return 0, 0, "", fmt.Errorf("rcon: invalid packet length %d", length) + if length < minPacketLen || length > maxReplyPacketLen { + return 0, 0, nil, fmt.Errorf("rcon: invalid packet length %d", length) } payload := make([]byte, length) if _, err = io.ReadFull(r, payload); err != nil { - return 0, 0, "", err + return 0, 0, nil, err } id = int32(binary.LittleEndian.Uint32(payload[0:4])) typ = int32(binary.LittleEndian.Uint32(payload[4:8])) // Strip the two trailing null bytes from the body. - body = string(payload[8 : length-2]) + body = payload[8 : length-2] return id, typ, body, nil } diff --git a/internal/rcon/rcon_test.go b/internal/rcon/rcon_test.go index 33424a3..ca4954a 100644 --- a/internal/rcon/rcon_test.go +++ b/internal/rcon/rcon_test.go @@ -1,10 +1,12 @@ package rcon_test import ( + "context" "encoding/binary" "errors" "io" "net" + "strings" "sync" "testing" "time" @@ -19,6 +21,9 @@ type fakeRCON struct { password string replies map[string]string wg sync.WaitGroup + // hang leaves every command unanswered; doubleEnd answers the end marker + // twice, as a Source-engine server does. + hang, doubleEnd bool } func startFakeRCON(t *testing.T, password string, replies map[string]string) *fakeRCON { @@ -72,14 +77,43 @@ func (f *fakeRCON) handle(conn net.Conn) { _ = writeFramePacket(conn, -1, 0, "") continue } + if f.hang { + continue + } + // Minecraft drops a connection whose read holds more than the one + // packet: a client that pipelines its next packet is cut off. + if pipelined(conn) { + return + } + // Minecraft splits a reply into 4096-byte bodies, even mid-rune. reply := f.replies[body] + for len(reply) > 4096 { + _ = writeFramePacket(conn, id, 0, reply[:4096]) + reply = reply[4096:] + } _ = writeFramePacket(conn, id, 0, reply) default: - _ = writeFramePacket(conn, id, 0, "") + if f.hang { + continue + } + _ = writeFramePacket(conn, id, 0, "Unknown request 0") + if f.doubleEnd { + _ = writeFramePacket(conn, id, 0, "") + } } } } +// pipelined reports whether the client already sent more bytes behind the +// packet just read. +func pipelined(conn net.Conn) bool { + _ = conn.SetReadDeadline(time.Now().Add(30 * time.Millisecond)) + defer conn.SetReadDeadline(time.Time{}) + var b [1]byte + n, _ := conn.Read(b[:]) + return n > 0 +} + func writeFramePacket(w io.Writer, id, typ int32, body string) error { b := []byte(body) length := int32(4 + 4 + len(b) + 2) @@ -178,3 +212,83 @@ func TestExecuteGracefulShutdownSequence(t *testing.T) { t.Fatalf("stop = %q, %v", out, err) } } + +func TestExecuteReassemblesALongReply(t *testing.T) { + // 3000 three-byte runes: 9000 bytes over three packets, split mid-rune. + long := strings.Repeat("封", 3000) + f := startFakeRCON(t, "pw", map[string]string{"banlist": long}) + defer f.stop() + c, err := rcon.Dial(f.addr(), "pw", time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer c.Close() + got, err := c.Execute("banlist") + if err != nil { + t.Fatalf("execute: %v", err) + } + if len(got) != 9000 || got != long { + t.Fatalf("reply = %d bytes, want the 9000-byte original", len(got)) + } +} + +func TestExecuteSkipsALeftoverEndMarker(t *testing.T) { + f := startFakeRCON(t, "pw", map[string]string{"list": "There are 0 of a max of 20 players online", "seed": "Seed: [42]"}) + f.doubleEnd = true + defer f.stop() + c, err := rcon.Dial(f.addr(), "pw", time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer c.Close() + if got, err := c.Execute("list"); err != nil || got != "There are 0 of a max of 20 players online" { + t.Fatalf("first = %q, %v", got, err) + } + if got, err := c.Execute("seed"); err != nil || got != "Seed: [42]" { + t.Fatalf("second = %q, %v", got, err) + } +} + +func TestExecuteContextGivesUpOnAHungServer(t *testing.T) { + f := startFakeRCON(t, "pw", nil) + f.hang = true + defer f.stop() + c, err := rcon.Dial(f.addr(), "pw", time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer c.Close() + + ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond) + defer cancel() + start := time.Now() + if _, err := c.ExecuteContext(ctx, "banlist"); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("deadline: err = %v, want context.DeadlineExceeded", err) + } + if took := time.Since(start); took > 2*time.Second { + t.Fatalf("deadline: took %v", took) + } + + ctx, cancel = context.WithCancel(context.Background()) + time.AfterFunc(100*time.Millisecond, cancel) + start = time.Now() + if _, err := c.ExecuteContext(ctx, "banlist"); !errors.Is(err, context.Canceled) { + t.Fatalf("cancel: err = %v, want context.Canceled", err) + } + if took := time.Since(start); took > 2*time.Second { + t.Fatalf("cancel: took %v", took) + } + + // A caller's own SetDeadline bounds a context-free Execute. + if err := c.SetDeadline(time.Now().Add(150 * time.Millisecond)); err != nil { + t.Fatal(err) + } + start = time.Now() + var ne net.Error + if _, err := c.Execute("banlist"); !errors.As(err, &ne) || !ne.Timeout() { + t.Fatalf("set deadline: err = %v, want a timeout", err) + } + if took := time.Since(start); took > 2*time.Second { + t.Fatalf("set deadline: took %v", took) + } +}