Loading internal/api/console.go +1 −1 Changes for internal/api/console.go: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -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 } Loading internal/rcon/rcon.go +84 −18 Changes for internal/rcon/rcon.go: 84 added lines, 18 removed lines. Original line number Diff line number Diff line Loading @@ -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" Loading @@ -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") Loading @@ -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 Loading Loading @@ -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() if err := writePacket(c.conn, id, typeExecCommand, cmd); err != nil { 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 "", c.ctxErr(ctx, err) } var reply []byte sentEnd := false for { respID, _, body, err := readPacket(c.conn) if err != nil { return "", err return "", c.ctxErr(ctx, err) } if respID != id { return "", fmt.Errorf("rcon: response id mismatch: got %d want %d", respID, id) 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. } return body, nil } } // 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 err } // auth performs the SERVERDATA_AUTH handshake. Some servers emit an empty Loading Loading @@ -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 } Loading internal/rcon/rcon_test.go +114 −0 Changes for internal/rcon/rcon_test.go: 114 added lines, 0 removed lines. Original line number Diff line number Diff line package rcon_test import ( "context" "encoding/binary" "errors" "io" "net" "strings" "sync" "testing" "time" Loading @@ -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 { Loading Loading @@ -72,13 +77,42 @@ 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: 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) Loading Loading @@ -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) } } Loading
internal/api/console.go +1 −1 Changes for internal/api/console.go: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -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 } Loading
internal/rcon/rcon.go +84 −18 Changes for internal/rcon/rcon.go: 84 added lines, 18 removed lines. Original line number Diff line number Diff line Loading @@ -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" Loading @@ -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") Loading @@ -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 Loading Loading @@ -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() if err := writePacket(c.conn, id, typeExecCommand, cmd); err != nil { 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 "", c.ctxErr(ctx, err) } var reply []byte sentEnd := false for { respID, _, body, err := readPacket(c.conn) if err != nil { return "", err return "", c.ctxErr(ctx, err) } if respID != id { return "", fmt.Errorf("rcon: response id mismatch: got %d want %d", respID, id) 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. } return body, nil } } // 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 err } // auth performs the SERVERDATA_AUTH handshake. Some servers emit an empty Loading Loading @@ -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 } Loading
internal/rcon/rcon_test.go +114 −0 Changes for internal/rcon/rcon_test.go: 114 added lines, 0 removed lines. Original line number Diff line number Diff line package rcon_test import ( "context" "encoding/binary" "errors" "io" "net" "strings" "sync" "testing" "time" Loading @@ -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 { Loading Loading @@ -72,13 +77,42 @@ 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: 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) Loading Loading @@ -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) } }