summaryrefslogtreecommitdiff
path: root/internal
diff options
context:
space:
mode:
authorPaul Bütow <pbuetow@mimecast.com>2020-01-26 11:26:53 +0000
committerPaul Bütow <pbuetow@mimecast.com>2020-02-07 13:31:15 +0000
commit0945da8dfefcbb723eecea0e5f4eafff63398253 (patch)
treef06dab4d2bf21d25d176b23d5baeca588d27f5d7 /internal
parent2a8e5de265a0e0a31a5834909d6879f5c9941467 (diff)
Introduce drun command, refactor code to use context package
Diffstat (limited to 'internal')
-rw-r--r--internal/clients/args.go3
-rw-r--r--internal/clients/baseclient.go130
-rw-r--r--internal/clients/catclient.go20
-rw-r--r--internal/clients/client.go5
-rw-r--r--internal/clients/connectionmaker.go12
-rw-r--r--internal/clients/execclient.go48
-rw-r--r--internal/clients/grepclient.go20
-rw-r--r--internal/clients/handlers/basehandler.go84
-rw-r--r--internal/clients/handlers/clienthandler.go11
-rw-r--r--internal/clients/handlers/handler.go12
-rw-r--r--internal/clients/handlers/healthhandler.go21
-rw-r--r--internal/clients/handlers/maprhandler.go21
-rw-r--r--internal/clients/handlers/withcancel.go24
-rw-r--r--internal/clients/healthclient.go7
-rw-r--r--internal/clients/maker.go8
-rw-r--r--internal/clients/maprclient.go52
-rw-r--r--internal/clients/remote/connection.go116
-rw-r--r--internal/clients/runclient.go40
-rw-r--r--internal/clients/stats.go8
-rw-r--r--internal/clients/tailclient.go21
-rw-r--r--internal/discovery/comma.go2
-rw-r--r--internal/discovery/discovery.go21
-rw-r--r--internal/discovery/file.go2
-rw-r--r--internal/io/fs/catfile.go (renamed from internal/fs/catfile.go)6
-rw-r--r--internal/io/fs/filereader.go (renamed from internal/fs/filereader.go)9
-rw-r--r--internal/io/fs/permissions/permission.go (renamed from internal/fs/permissions/permission.go)2
-rw-r--r--internal/io/fs/permissions/permission_linux.c (renamed from internal/fs/permissions/permission_linux.c)0
-rw-r--r--internal/io/fs/permissions/permission_linux.go (renamed from internal/fs/permissions/permission_linux.go)0
-rw-r--r--internal/io/fs/permissions/permission_linux.h (renamed from internal/fs/permissions/permission_linux.h)0
-rw-r--r--internal/io/fs/permissions/permission_test.go (renamed from internal/fs/permissions/permission_test.go)0
-rw-r--r--internal/io/fs/readfile.go (renamed from internal/fs/readfile.go)73
-rw-r--r--internal/io/fs/stats.go (renamed from internal/fs/stats.go)0
-rw-r--r--internal/io/fs/tailfile.go (renamed from internal/fs/tailfile.go)6
-rw-r--r--internal/io/line/line.go (renamed from internal/fs/lineread.go)14
-rw-r--r--internal/io/logger/logger.go (renamed from internal/logger/logger.go)56
-rw-r--r--internal/io/run/run.go104
-rw-r--r--internal/mapr/aggregateset.go5
-rw-r--r--internal/mapr/client/aggregate.go25
-rw-r--r--internal/mapr/groupset.go5
-rw-r--r--internal/mapr/logformat/parser.go2
-rw-r--r--internal/mapr/query.go2
-rw-r--r--internal/mapr/server/aggregate.go141
-rw-r--r--internal/mapr/wherecondition.go2
-rw-r--r--internal/omode/mode.go6
-rw-r--r--internal/pprof/pprof.go3
-rw-r--r--internal/prompt/prompt.go2
-rw-r--r--internal/server/handlers/controlhandler.go42
-rw-r--r--internal/server/handlers/handler.go2
-rw-r--r--internal/server/handlers/mapcommand.go35
-rw-r--r--internal/server/handlers/readcommand.go158
-rw-r--r--internal/server/handlers/runcommand.go73
-rw-r--r--internal/server/handlers/serverhandler.go521
-rw-r--r--internal/server/server.go70
-rw-r--r--internal/server/stats.go10
-rw-r--r--internal/ssh/client/authmethods.go2
-rw-r--r--internal/ssh/client/hostkeycallback.go10
-rw-r--r--internal/ssh/server/hostkey.go2
-rw-r--r--internal/ssh/server/publickeycallback.go2
-rw-r--r--internal/ssh/ssh.go2
-rw-r--r--internal/user/name.go15
-rw-r--r--internal/user/server/user.go44
-rw-r--r--internal/version/version.go22
62 files changed, 1161 insertions, 1000 deletions
diff --git a/internal/clients/args.go b/internal/clients/args.go
index 5fe0a72..dea5a9e 100644
--- a/internal/clients/args.go
+++ b/internal/clients/args.go
@@ -9,10 +9,9 @@ type Args struct {
Mode omode.Mode
ServersStr string
UserName string
- Files string
+ What string
Regex string
TrustAllHosts bool
Discovery string
ConnectionsPerCPU int
- PingTimeout int
}
diff --git a/internal/clients/baseclient.go b/internal/clients/baseclient.go
index 574ae94..b1540ea 100644
--- a/internal/clients/baseclient.go
+++ b/internal/clients/baseclient.go
@@ -1,13 +1,14 @@
package clients
import (
+ "context"
"regexp"
"sync"
"time"
"github.com/mimecast/dtail/internal/clients/remote"
"github.com/mimecast/dtail/internal/discovery"
- "github.com/mimecast/dtail/internal/logger"
+ "github.com/mimecast/dtail/internal/io/logger"
"github.com/mimecast/dtail/internal/omode"
"github.com/mimecast/dtail/internal/ssh/client"
@@ -27,111 +28,110 @@ type baseClient struct {
sshAuthMethods []gossh.AuthMethod
// To deal with SSH host keys
hostKeyCallback *client.HostKeyCallback
- // To stop the client.
- stop chan struct{}
- // To indicate that the client has stopped.
- stopped chan struct{}
// Throttle how fast we initiate SSH connections concurrently
throttleCh chan struct{}
// Retry connection upon failure?
retry bool
- // Connection helper.
- maker connectionMaker
+ // Connection maker helper.
+ maker maker
}
-func (c *baseClient) init(maker connectionMaker) {
+func (c *baseClient) init(maker maker) {
logger.Info("Initiating base client")
c.maker = maker
- //c.connections = make(map[string]*remote.Connection)
c.sshAuthMethods, c.hostKeyCallback = client.InitSSHAuthMethods(c.TrustAllHosts, c.throttleCh)
+ discoveryService := discovery.New(c.Discovery, c.ServersStr, discovery.Shuffle)
- // Retrieve a shuffled list of remote dtail servers.
- shuffleServers := true
- discoveryService := discovery.New(c.Discovery, c.ServersStr, shuffleServers)
for _, server := range discoveryService.ServerList() {
- c.connections = append(c.connections, c.maker.makeConnection(server, c.sshAuthMethods, c.hostKeyCallback))
+ c.connections = append(c.connections, c.makeConnection(server, c.sshAuthMethods, c.hostKeyCallback))
}
if _, err := regexp.Compile(c.Regex); err != nil {
logger.FatalExit(c.Regex, "Can't test compile regex", err)
}
- // Periodically check for unknown hosts, and ask the user whether to trust them or not.
- go c.hostKeyCallback.PromptAddHosts(c.stop)
-
- // Periodically print out connection stats to the client.
c.stats = newTailStats(len(c.connections))
- go c.stats.periodicLogStats(c.throttleCh, c.stop)
}
-func (c *baseClient) Start() (status int) {
+func (c *baseClient) Start(ctx context.Context) (status int) {
+ // Periodically check for unknown hosts, and ask the user whether to trust them or not.
+ go c.hostKeyCallback.PromptAddHosts(ctx)
+ // Periodically print out connection stats to the client.
+ go c.stats.periodicLogStats(ctx, c.throttleCh)
+ // Keep count of active connections
active := make(chan struct{}, len(c.connections))
- var wg sync.WaitGroup
- wg.Add(len(c.connections))
-
+ var mutex sync.Mutex
for i, conn := range c.connections {
go func(i int, conn *remote.Connection) {
- active <- struct{}{}
- defer func() {
- logger.Debug(conn.Server, "Disconnected completely...")
- <-active
- }()
- wg.Done()
-
- for {
- conn.Start(c.throttleCh, c.stats.connectionsEstCh)
- if !c.retry {
- return
- }
- time.Sleep(time.Second * 2)
- logger.Debug(conn.Server, "Reconencting")
- conn = c.maker.makeConnection(conn.Server, c.sshAuthMethods, c.hostKeyCallback)
- c.connections[i] = conn
+ connStatus := c.start(ctx, active, i, conn)
+
+ // Update global status.
+ mutex.Lock()
+ defer mutex.Unlock()
+ if connStatus > status {
+ status = connStatus
}
}(i, conn)
}
- wg.Wait()
- c.waitUntilDone(active)
-
+ c.waitUntilDone(ctx, active)
return
}
-func (c *baseClient) waitUntilDone(active chan struct{}) {
- defer close(c.stopped)
+func (c *baseClient) start(ctx context.Context, active chan struct{}, i int, conn *remote.Connection) (status int) {
+ // Increment connection count
+ active <- struct{}{}
+ // Derement connection count
+ defer func() { <-active }()
- if c.Mode != omode.TailClient {
- c.waitUntilZero(active)
- logger.Info("All connections stopped")
- return
- }
+ for {
+ connCtx, cancel := conn.Handler.WithCancel(ctx)
+ defer cancel()
- <-c.stop
- logger.Info("Stopping client")
- for _, conn := range c.connections {
- conn.Stop()
+ conn.Start(connCtx, cancel, c.throttleCh, c.stats.connectionsEstCh)
+ // Retrieve status code from handler (dtail client will exit with that status)
+ status = conn.Handler.Status()
+
+ if !c.retry {
+ return
+ }
+
+ time.Sleep(time.Second * 2)
+ logger.Debug(conn.Server, "Reconnecting")
+
+ conn = c.makeConnection(conn.Server, c.sshAuthMethods, c.hostKeyCallback)
+ c.connections[i] = conn
}
+}
- c.waitUntilZero(active)
+func (c *baseClient) makeConnection(server string, sshAuthMethods []gossh.AuthMethod, hostKeyCallback *client.HostKeyCallback) *remote.Connection {
+ conn := remote.NewConnection(server, c.UserName, sshAuthMethods, hostKeyCallback)
+ conn.Handler = c.maker.makeHandler(server)
+ conn.Commands = c.maker.makeCommands()
+
+ return conn
}
-func (c *baseClient) waitUntilZero(active chan struct{}) {
+func (c *baseClient) waitUntilDone(ctx context.Context, active chan struct{}) {
+ defer logger.Info("Terminated connection")
+
+ // We want to have at least one active connection
+ <-active
+ // Put it back on the channel
+ active <- struct{}{}
+
+ if c.Mode == omode.TailClient {
+ <-ctx.Done()
+ }
+
for {
- logger.Debug("Active connections", len(active))
- if len(active) == 0 {
+ numActive := len(active)
+ if numActive == 0 {
return
}
+ logger.Debug("Active connections", numActive)
time.Sleep(time.Second)
}
}
-
-func (c *baseClient) Stop() {
- close(c.stop)
- <-c.WaitC()
-}
-
-func (c *baseClient) WaitC() <-chan struct{} {
- return c.stopped
-}
diff --git a/internal/clients/catclient.go b/internal/clients/catclient.go
index 5ea701d..7fd6bdc 100644
--- a/internal/clients/catclient.go
+++ b/internal/clients/catclient.go
@@ -7,11 +7,7 @@ import (
"strings"
"github.com/mimecast/dtail/internal/clients/handlers"
- "github.com/mimecast/dtail/internal/clients/remote"
"github.com/mimecast/dtail/internal/omode"
- "github.com/mimecast/dtail/internal/ssh/client"
-
- gossh "golang.org/x/crypto/ssh"
)
// CatClient is a client for returning a whole file from the beginning to the end.
@@ -31,8 +27,6 @@ func NewCatClient(args Args) (*CatClient, error) {
c := CatClient{
baseClient: baseClient{
Args: args,
- stop: make(chan struct{}),
- stopped: make(chan struct{}),
throttleCh: make(chan struct{}, args.ConnectionsPerCPU*runtime.NumCPU()),
retry: false,
},
@@ -43,11 +37,13 @@ func NewCatClient(args Args) (*CatClient, error) {
return &c, nil
}
-func (c CatClient) makeConnection(server string, sshAuthMethods []gossh.AuthMethod, hostKeyCallback *client.HostKeyCallback) *remote.Connection {
- conn := remote.NewConnection(server, c.UserName, sshAuthMethods, hostKeyCallback)
- conn.Handler = handlers.NewClientHandler(server, c.PingTimeout)
- for _, file := range strings.Split(c.Files, ",") {
- conn.Commands = append(conn.Commands, fmt.Sprintf("%s %s regex %s", c.Mode.String(), file, c.Regex))
+func (c CatClient) makeHandler(server string) handlers.Handler {
+ return handlers.NewClientHandler(server)
+}
+
+func (c CatClient) makeCommands() (commands []string) {
+ for _, file := range strings.Split(c.What, ",") {
+ commands = append(commands, fmt.Sprintf("%s %s regex %s", c.Mode.String(), file, c.Regex))
}
- return conn
+ return
}
diff --git a/internal/clients/client.go b/internal/clients/client.go
index 85d1aae..1fc5e23 100644
--- a/internal/clients/client.go
+++ b/internal/clients/client.go
@@ -1,7 +1,8 @@
package clients
+import "context"
+
// Client is the interface for the end user command line client.
type Client interface {
- Start() int
- Stop()
+ Start(ctx context.Context) int
}
diff --git a/internal/clients/connectionmaker.go b/internal/clients/connectionmaker.go
deleted file mode 100644
index 0617992..0000000
--- a/internal/clients/connectionmaker.go
+++ /dev/null
@@ -1,12 +0,0 @@
-package clients
-
-import (
- "github.com/mimecast/dtail/internal/clients/remote"
- "github.com/mimecast/dtail/internal/ssh/client"
-
- gossh "golang.org/x/crypto/ssh"
-)
-
-type connectionMaker interface {
- makeConnection(server string, sshAuthMethods []gossh.AuthMethod, hostKeyCallback *client.HostKeyCallback) *remote.Connection
-}
diff --git a/internal/clients/execclient.go b/internal/clients/execclient.go
deleted file mode 100644
index 10bd081..0000000
--- a/internal/clients/execclient.go
+++ /dev/null
@@ -1,48 +0,0 @@
-package clients
-
-import (
- "fmt"
- "runtime"
- "strings"
-
- "github.com/mimecast/dtail/internal/clients/handlers"
- "github.com/mimecast/dtail/internal/clients/remote"
- "github.com/mimecast/dtail/internal/omode"
- "github.com/mimecast/dtail/internal/ssh/client"
-
- gossh "golang.org/x/crypto/ssh"
-)
-
-// ExecClient is a client for execute various commands on the server.
-type ExecClient struct {
- baseClient
-}
-
-// NewExecClient returns a new cat client.
-func NewExecClient(args Args) (*ExecClient, error) {
- args.Regex = "."
- args.Mode = omode.ExecClient
-
- c := ExecClient{
- baseClient: baseClient{
- Args: args,
- stop: make(chan struct{}),
- stopped: make(chan struct{}),
- throttleCh: make(chan struct{}, args.ConnectionsPerCPU*runtime.NumCPU()),
- retry: false,
- },
- }
-
- c.init(c)
-
- return &c, nil
-}
-
-func (c ExecClient) makeConnection(server string, sshAuthMethods []gossh.AuthMethod, hostKeyCallback *client.HostKeyCallback) *remote.Connection {
- conn := remote.NewConnection(server, c.UserName, sshAuthMethods, hostKeyCallback)
- conn.Handler = handlers.NewClientHandler(server, c.PingTimeout)
- for _, file := range strings.Split(c.Files, ";") {
- conn.Commands = append(conn.Commands, fmt.Sprintf("%s %s", c.Mode.String(), file))
- }
- return conn
-}
diff --git a/internal/clients/grepclient.go b/internal/clients/grepclient.go
index c568f63..8d11458 100644
--- a/internal/clients/grepclient.go
+++ b/internal/clients/grepclient.go
@@ -7,11 +7,7 @@ import (
"strings"
"github.com/mimecast/dtail/internal/clients/handlers"
- "github.com/mimecast/dtail/internal/clients/remote"
"github.com/mimecast/dtail/internal/omode"
- "github.com/mimecast/dtail/internal/ssh/client"
-
- gossh "golang.org/x/crypto/ssh"
)
// GrepClient searches a remote file for all lines matching a regular expression. Only the matching lines are displayed.
@@ -29,8 +25,6 @@ func NewGrepClient(args Args) (*GrepClient, error) {
c := GrepClient{
baseClient: baseClient{
Args: args,
- stop: make(chan struct{}),
- stopped: make(chan struct{}),
throttleCh: make(chan struct{}, args.ConnectionsPerCPU*runtime.NumCPU()),
retry: false,
},
@@ -41,13 +35,13 @@ func NewGrepClient(args Args) (*GrepClient, error) {
return &c, nil
}
-func (c GrepClient) makeConnection(server string, sshAuthMethods []gossh.AuthMethod, hostKeyCallback *client.HostKeyCallback) *remote.Connection {
- conn := remote.NewConnection(server, c.UserName, sshAuthMethods, hostKeyCallback)
- conn.Handler = handlers.NewClientHandler(server, c.PingTimeout)
+func (c GrepClient) makeHandler(server string) handlers.Handler {
+ return handlers.NewClientHandler(server)
+}
- for _, file := range strings.Split(c.Files, ",") {
- conn.Commands = append(conn.Commands, fmt.Sprintf("%s %s regex %s", c.Mode.String(), file, c.Regex))
+func (c GrepClient) makeCommands() (commands []string) {
+ for _, file := range strings.Split(c.What, ",") {
+ commands = append(commands, fmt.Sprintf("%s %s regex %s", c.Mode.String(), file, c.Regex))
}
-
- return conn
+ return
}
diff --git a/internal/clients/handlers/basehandler.go b/internal/clients/handlers/basehandler.go
index 19246f9..68b8ddc 100644
--- a/internal/clients/handlers/basehandler.go
+++ b/internal/clients/handlers/basehandler.go
@@ -1,60 +1,44 @@
package handlers
import (
- "github.com/mimecast/dtail/internal/logger"
- "errors"
+ "encoding/base64"
"fmt"
"io"
+ "strconv"
"strings"
"time"
+
+ "github.com/mimecast/dtail/internal/io/logger"
+ "github.com/mimecast/dtail/internal/version"
)
type baseHandler struct {
+ withCancel
server string
shellStarted bool
commands chan string
- pong chan struct{}
receiveBuf []byte
- stop chan struct{}
- pingTimeout int
+ status int
}
func (h *baseHandler) Server() string {
return h.server
}
-// Used to determine whether server is still responding to requests or not.
-func (h *baseHandler) Ping() error {
- if h.pingTimeout == 0 {
- // Server ping disabled
- return nil
- }
-
- if err := h.SendCommand("ping"); err != nil {
- return err
- }
-
- select {
- case <-h.pong:
- return nil
- case <-time.After(time.Duration(h.pingTimeout) * time.Second):
- }
-
- return errors.New("Didn't receive any server pongs (ping replies)")
+func (h *baseHandler) Status() int {
+ return h.status
}
-func (h *baseHandler) SendCommand(command string) error {
- if command == "ping" {
- logger.Trace("Sending command", h.server, command)
- } else {
- logger.Debug("Sending command", h.server, command)
- }
+// SendMessage to the server.
+func (h *baseHandler) SendMessage(command string) error {
+ encoded := base64.StdEncoding.EncodeToString([]byte(command))
+ logger.Debug("Sending command", h.server, command, encoded)
select {
- case h.commands <- fmt.Sprintf("%s;", command):
+ case h.commands <- fmt.Sprintf("protocol %s base64 %v;", version.ProtocolCompat, encoded):
case <-time.After(time.Second * 5):
- return errors.New("Timed out sending command " + command)
- case <-h.stop:
+ return fmt.Errorf("Timed out sending command '%s' (base64: '%s')", command, encoded)
+ case <-h.ctx.Done():
}
return nil
@@ -81,7 +65,7 @@ func (h *baseHandler) Read(p []byte) (n int, err error) {
select {
case command := <-h.commands:
n = copy(p, []byte(command))
- case <-h.stop:
+ case <-h.ctx.Done():
return 0, io.EOF
}
return
@@ -92,6 +76,7 @@ func (h *baseHandler) handleMessageType(message string) {
if len(h.receiveBuf) == 0 {
return
}
+
// Hidden server commands starti with a dot "."
if h.receiveBuf[0] == '.' {
h.handleHiddenMessage(message)
@@ -108,6 +93,7 @@ func (h *baseHandler) handleMessageType(message string) {
h.receiveBuf = h.receiveBuf[:0]
return
}
+
logger.Raw(message)
h.receiveBuf = h.receiveBuf[:0]
}
@@ -116,19 +102,27 @@ func (h *baseHandler) handleMessageType(message string) {
// to the end user.
func (h *baseHandler) handleHiddenMessage(message string) {
switch {
- case strings.HasPrefix(message, ".pong"):
- h.pong <- struct{}{}
case strings.HasPrefix(message, ".syn close connection"):
- h.SendCommand("ack close connection")
- }
-}
+ h.SendMessage(".ack close connection")
+ select {
+ case <-time.After(time.Second * 1):
+ logger.Debug("Shutting down client after timeout and sending ack to server")
+ h.withCancel.shutdown()
+ case <-h.ctx.Done():
+ }
-// Stop the handler.
-func (h *baseHandler) Stop() {
- select {
- case <-h.stop:
- default:
- logger.Debug("Stopping base handler", h.server)
- close(h.stop)
+ case strings.HasPrefix(message, ".run exitstatus"):
+ splitted := strings.Split(strings.TrimSuffix(message, "\n"), " ")
+ if len(splitted) != 3 {
+ logger.Error("Unable to retrieve exitstatus", message)
+ return
+ }
+ i, err := strconv.Atoi(splitted[2])
+ if err != nil {
+ logger.Error("Unable to retrieve exitstatus", message, err)
+ return
+ }
+ h.status = i
+ logger.Debug("Retrieved exitstatus", h.status)
}
}
diff --git a/internal/clients/handlers/clienthandler.go b/internal/clients/handlers/clienthandler.go
index 4738cd3..fcd8052 100644
--- a/internal/clients/handlers/clienthandler.go
+++ b/internal/clients/handlers/clienthandler.go
@@ -1,7 +1,7 @@
package handlers
import (
- "github.com/mimecast/dtail/internal/logger"
+ "github.com/mimecast/dtail/internal/io/logger"
)
// ClientHandler is the basic client handler interface.
@@ -10,7 +10,7 @@ type ClientHandler struct {
}
// NewClientHandler creates a new client handler.
-func NewClientHandler(server string, pingTimeout int) *ClientHandler {
+func NewClientHandler(server string) *ClientHandler {
logger.Debug(server, "Creating new client handler")
return &ClientHandler{
@@ -18,9 +18,10 @@ func NewClientHandler(server string, pingTimeout int) *ClientHandler {
server: server,
shellStarted: false,
commands: make(chan string),
- pong: make(chan struct{}, 1),
- stop: make(chan struct{}),
- pingTimeout: pingTimeout,
+ status: -1,
+ withCancel: withCancel{
+ done: make(chan struct{}),
+ },
},
}
}
diff --git a/internal/clients/handlers/handler.go b/internal/clients/handlers/handler.go
index 2013be0..c53ca34 100644
--- a/internal/clients/handlers/handler.go
+++ b/internal/clients/handlers/handler.go
@@ -1,12 +1,16 @@
package handlers
-import "io"
+import (
+ "context"
+ "io"
+)
// Handler provides all methods which can be run on any client handler.