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
3 changes: 3 additions & 0 deletions cmd/root.go
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,9 @@ func rootCmdRun(cmd *cobra.Command, _ []string) {
log.WithField("error", err).Fatal("failed to create pelican system user")
return
}
if err := config.ConfigurePasswd(); err != nil {
log.WithField("error", err).Fatal("failed to create passwd files for pelican")
}
log.WithFields(log.Fields{
"username": config.Get().System.Username,
"uid": config.Get().System.User.Uid,
Expand Down
82 changes: 59 additions & 23 deletions config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,11 @@ import (
"text/template"
"time"

"github.com/gbrlsnchs/jwt/v3"

"emperror.dev/errors"
"github.com/acobaugh/osrelease"
"github.com/apex/log"
"github.com/creasty/defaults"
"github.com/gbrlsnchs/jwt/v3"
"golang.org/x/sys/unix"
"gopkg.in/yaml.v2"

Expand Down Expand Up @@ -129,9 +128,9 @@ type RemoteQueryConfiguration struct {
// be less likely to cause performance issues on the Panel.
BootServersPerPage int `default:"50" yaml:"boot_servers_per_page"`

//When using services like Cloudflare Access to manage access to
//a specific system via an external authentication system,
//it is possible to add special headers to bypass authentication.
//When using services like Cloudflare Access to manage access to
//a specific system via an external authentication system,
//it is possible to add special headers to bypass authentication.
//The mentioned headers can be appended to queries sent from Wings to the panel.
CustomHeaders map[string]string `yaml:"custom_headers"`
}
Expand Down Expand Up @@ -187,11 +186,23 @@ type SystemConfiguration struct {
Uid int `yaml:"uid"`
Gid int `yaml:"gid"`

// Passwd controls weather a passwd file is mounted in the container
// at /etc/passwd to resolve missing user issues
Passwd bool `json:"mount_passwd" yaml:"mount_passwd" default:"true"`
PasswdFile string `json:"passwd_file" yaml:"passwd_file" default:"/etc/pelican/passwd"`
} `yaml:"user"`
// Passwd controls weather a passwd and group file is mounted in the container
// at /etc/passwd to resolve missing user/group issues inside the container
Passwd struct {
Enable bool `json:"enable" yaml:"enable" default:"true"`
Directory string `json:"directory" yaml:"directory" default:"/etc/pelican"`
} `json:"passwd" yaml:"passwd"`
} `json:"user" yaml:"user"`

// MachineID manages the mounting of a 'machine-id' file for containers as required for
// some game servers. I.E. Hytale
MachineID struct {
// Enable controls if the machine-id file is generated and mounted into the server container
// This is enabled by default
Enable bool `json:"enable" yaml:"enable" default:"true"`
// FilePath is the full path to the machine-id file that will be generated and mounted
Directory string `json:"directory" yaml:"directory" default:"/etc/pelican/machine-id"`
} `json:"machine_id" yaml:"machine_id"`

// The amount of time in seconds that can elapse before a server's disk space calculation is
// considered stale and a re-check should occur. DANGER: setting this value too low can seriously
Expand Down Expand Up @@ -604,19 +615,6 @@ func ConfigureDirectories() error {
return err
}

log.WithField("filepath", _config.System.User.PasswdFile).Debug("ensuring passwd file exists")
if passwd, err := os.Create(_config.System.User.PasswdFile); err != nil {
return err
} else {
// the WriteFile method returns an error if unsuccessful
err := os.WriteFile(passwd.Name(), []byte(fmt.Sprintf("container:x:%d:%d::/home/container:/usr/sbin/nologin", _config.System.User.Uid, _config.System.User.Gid)), 0644)
// handle this error
if err != nil {
// print it out
fmt.Println(err)
}
}

// There are a non-trivial number of users out there whose data directories are actually a
// symlink to another location on the disk. If we do not resolve that final destination at this
// point things will appear to work, but endless errors will be encountered when we try to
Expand All @@ -638,6 +636,11 @@ func ConfigureDirectories() error {
return err
}

log.WithField("path", _config.System.TmpDirectory).Debug("ensuring temporary data directory exists")
if err := os.MkdirAll(_config.System.TmpDirectory, 0o700); err != nil {
return err
}

log.WithField("path", _config.System.ArchiveDirectory).Debug("ensuring archive data directory exists")
if err := os.MkdirAll(_config.System.ArchiveDirectory, 0o700); err != nil {
return err
Expand All @@ -648,9 +651,42 @@ func ConfigureDirectories() error {
return err
}

log.WithField("path", _config.System.User.Passwd.Directory).Debug("ensuring passwd directory exists")
if err := os.MkdirAll(_config.System.User.Passwd.Directory, 0o700); err != nil {
return err
}

log.WithField("path", _config.System.MachineID.Directory).Debug("ensuring machine-id directory exists")
if err := os.MkdirAll(_config.System.MachineID.Directory, 0o700); err != nil {
return err
}
return nil
}

// ConfigurePasswd generates the passwd and group files to be used by
// this looks cleaner than the previous way and is similar to pterodactyl
func ConfigurePasswd() (err error) {
if !_config.System.User.Passwd.Enable {
return
}
log.WithField("filepath", filepath.Join(_config.System.User.Passwd.Directory, "passwd")).
Debug("ensuring passwd file exists")
if err = os.WriteFile(filepath.Join(_config.System.User.Passwd.Directory, "passwd"),
[]byte(fmt.Sprintf("container:x:%d:%d::/home/container:/usr/sbin/nologin",
_config.System.User.Uid, _config.System.User.Gid)), 0644); err != nil {
return fmt.Errorf("could not write passwd file: %w", err)
}

log.WithField("filepath", filepath.Join(_config.System.User.Passwd.Directory, "group")).
Debug("ensuring group file exists")
if err = os.WriteFile(filepath.Join(_config.System.User.Passwd.Directory, "group"),
[]byte(fmt.Sprintf("container:x:%d:container",
_config.System.User.Gid)), 0644); err != nil {
return fmt.Errorf("could not write group file: %w", err)
}
return
}

// EnableLogRotation writes a logrotate file for wings to the system logrotate
// configuration directory if one exists and a logrotate file is not found. This
// allows us to basically automate away the log rotation for most installs, but
Expand Down
1 change: 1 addition & 0 deletions router/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ func Configure(m *wserver.Manager, client remote.Client) *gin.Engine {
protected.GET("/api/servers", getAllServers)
protected.POST("/api/servers", postCreateServer)
protected.DELETE("/api/transfers/:server", deleteTransfer)
protected.POST("/api/deauthorize-user", postDeauthorizeUser)

// These are server specific routes, and require that the request be authorized, and
// that the server exist on the Daemon.
Expand Down
14 changes: 12 additions & 2 deletions router/router_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,6 @@ func getServerInstallLogs(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"data": output})
}


// Handles a request to control the power state of a server. If the action being passed
// through is invalid a 404 is returned. Otherwise, a HTTP/202 Accepted response is returned
// and the actual power action is run asynchronously so that we don't have to block the
Expand Down Expand Up @@ -281,7 +280,16 @@ func deleteServer(c *gin.Context) {
p := fs.Path()
_ = fs.UnixFS().Close()
if err := os.RemoveAll(p); err != nil {
log.WithFields(log.Fields{"path": p, "error": err}).Warn("failed to remove server files during deletion process")
log.WithFields(log.Fields{"path": p, "error": err}).
Warn("failed to remove server files during deletion process")
}
}(s)

// remove hanging machine-id file for the server when removing
go func(s *server.Server) {
if err := os.Remove(filepath.Join(config.Get().System.MachineID.Directory, s.ID())); err != nil {
log.WithFields(log.Fields{"server_id": s.ID(), "error": err}).
Warn("failed to remove machine-id file for server")
}
}(s)

Expand All @@ -295,6 +303,8 @@ func deleteServer(c *gin.Context) {
// Adds any of the JTIs passed through in the body to the deny list for the websocket
// preventing any JWT generated before the current time from being used to connect to
// the socket or send along commands.
//
// deprecated: prefer /api/deauthorize-user
func postServerDenyWSTokens(c *gin.Context) {
var data struct {
JTIs []string `json:"jtis"`
Expand Down
104 changes: 88 additions & 16 deletions router/router_server_ws.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,17 @@ package router

import (
"context"
"encoding/json"
"net/http"
"time"

"emperror.dev/errors"
"github.com/gin-gonic/gin"
"github.com/goccy/go-json"
ws "github.com/gorilla/websocket"

"github.com/pelican-dev/wings/router/middleware"
"github.com/pelican-dev/wings/router/websocket"
"github.com/pelican-dev/wings/server"
"golang.org/x/time/rate"
)

var expectedCloseCodes = []int{
Expand All @@ -25,6 +28,27 @@ func getServerWebsocket(c *gin.Context) {
manager := middleware.ExtractManager(c)
s, _ := manager.Get(c.Param("server"))

// Limit the total number of websockets that can be opened at any one time for
// a server instance. This applies across all users connected to the server, and
// is not applied on a per-user basis.
//
// todo: it would be great to make this per-user instead, but we need to modify
// how we even request this endpoint in order for that to be possible. Some type
// of signed identifier in the URL that is verified on this end and set by the
// panel using a shared secret is likely the easiest option. The benefit of that
// is that we can both scope things to the user before authentication, and also
// verify that the JWT provided by the panel is assigned to the same user.
if s.Websockets().Len() >= 30 {
c.AbortWithStatusJSON(http.StatusBadRequest, gin.H{
"error": "Too many open websocket connections.",
})

return
}

c.Header("Content-Security-Policy", "default-src 'self'")
c.Header("X-Frame-Options", "DENY")

// Create a context that can be canceled when the user disconnects from this
// socket that will also cancel listeners running in separate threads. If the
// connection itself is terminated listeners using this context will also be
Expand All @@ -37,53 +61,101 @@ func getServerWebsocket(c *gin.Context) {
middleware.CaptureAndAbort(c, err)
return
}
defer handler.Connection.Close()

// Track this open connection on the server so that we can close them all programmatically
// if the server is deleted.
s.Websockets().Push(handler.Uuid(), &cancel)
handler.Logger().Debug("opening connection to server websocket")
defer s.Websockets().Remove(handler.Uuid())

defer func() {
s.Websockets().Remove(handler.Uuid())
handler.Logger().Debug("closing connection to server websocket")
go func() {
select {
// When the main context is canceled (through disconnect, server deletion, or server
// suspension) close the connection itself.
case <-ctx.Done():
handler.Logger().Debug("closing connection to server websocket")
if err := handler.Connection.Close(); err != nil {
handler.Logger().WithError(err).Error("failed to close websocket connection")
}
break
}
}()

// If the server is deleted we need to send a close message to the connected client
// so that they disconnect since there will be no more events sent along. Listen for
// the request context being closed to break this loop, otherwise this routine will
// be left hanging in the background.
go func() {
select {
case <-ctx.Done():
break
return
// If the server is deleted we need to send a close message to the connected client
// so that they disconnect since there will be no more events sent along. Listen for
// the request context being closed to break this loop, otherwise this routine will
//be left hanging in the background.
case <-s.Context().Done():
_ = handler.Connection.WriteControl(ws.CloseMessage, ws.FormatCloseMessage(ws.CloseGoingAway, "server deleted"), time.Now().Add(time.Second*5))
cancel()
break
}
}()

for {
j := websocket.Message{}
// Due to how websockets are handled we need to connect to the socket
// and _then_ abort it if the server is suspended. You cannot capture
// the HTTP response in the websocket client, thus we connect and then
// immediately close with failure.
if s.IsSuspended() {
_ = handler.Connection.WriteMessage(ws.CloseMessage, ws.FormatCloseMessage(4409, "server is suspended"))

_, p, err := handler.Connection.ReadMessage()
return
}

// There is a separate rate limiter that applies to individual message types
// within the actual websocket logic handler. _This_ rate limiter just exists
// to avoid enormous floods of data through the socket since we need to parse
// JSON each time. This rate limit realistically should never be hit since this
// would require sending 50+ messages a second over the websocket (no more than
// 10 per 200ms).
var throttled bool
rl := rate.NewLimiter(rate.Every(time.Millisecond*200), 10)

for {
t, p, err := handler.Connection.ReadMessage()
if err != nil {
if ws.IsUnexpectedCloseError(err, expectedCloseCodes...) {
handler.Logger().WithField("error", err).Warn("error handling websocket message for server")
}
break
}

if !rl.Allow() {
if !throttled {
throttled = true
_ = handler.Connection.WriteJSON(websocket.Message{Event: websocket.ThrottledEvent, Args: []string{"global"}})
}
continue
}

throttled = false

// If the message isn't a format we expect, or the length of the message is far larger
// than we'd ever expect, drop it. The websocket upgrader logic does enforce a maximum
// _compressed_ message size of 4Kb but that could decompress to a much larger amount
// of data.
if t != ws.TextMessage || len(p) > 32_768 {
continue
}

// Discard and JSON parse errors into the void and don't continue processing this
// specific socket request. If we did a break here the client would get disconnected
// from the socket, which is NOT what we want to do.
var j websocket.Message
if err := json.Unmarshal(p, &j); err != nil {
continue
}

go func(msg websocket.Message) {
if err := handler.HandleInbound(ctx, msg); err != nil {
_ = handler.SendErrorJson(msg, err)
if errors.Is(err, server.ErrSuspended) {
cancel()
} else {
_ = handler.SendErrorJson(msg, err)
}
}
}(j)
}
Expand Down
31 changes: 31 additions & 0 deletions router/router_system.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
"github.com/pelican-dev/wings/config"
"github.com/pelican-dev/wings/internal/diagnostics"
"github.com/pelican-dev/wings/router/middleware"
"github.com/pelican-dev/wings/router/tokens"
"github.com/pelican-dev/wings/server"
"github.com/pelican-dev/wings/server/installer"
"github.com/pelican-dev/wings/system"
Expand Down Expand Up @@ -256,3 +257,33 @@ func postUpdateConfiguration(c *gin.Context) {
Applied: true,
})
}

func postDeauthorizeUser(c *gin.Context) {
var data struct {
User string `json:"user"`
Servers []string `json:"servers"`
}

if err := c.BindJSON(&data); err != nil {
return
}
Comment thread
parkervcp marked this conversation as resolved.

// todo: disconnect websockets more gracefully
m := middleware.ExtractManager(c)
if len(data.Servers) > 0 {
for _, uuid := range data.Servers {
if s, ok := m.Get(uuid); ok {
s.Websockets().CancelAll()
s.Sftp().Cancel(data.User)
tokens.DenyForServer(s.ID(), data.User)
}
}
} else {
for _, s := range m.All() {
s.Websockets().CancelAll()
s.Sftp().Cancel(data.User)
}
}
Comment thread
parkervcp marked this conversation as resolved.

c.Status(http.StatusNoContent)
}
Loading
Loading