summaryrefslogtreecommitdiff
path: root/internal
diff options
context:
space:
mode:
Diffstat (limited to 'internal')
-rw-r--r--internal/clients/args.go25
-rw-r--r--internal/clients/baseclient.go97
-rw-r--r--internal/clients/catclient.go13
-rw-r--r--internal/clients/connectors/connector.go17
-rw-r--r--internal/clients/connectors/serverconnection.go203
-rw-r--r--internal/clients/connectors/serverless.go116
-rw-r--r--internal/clients/grepclient.go15
-rw-r--r--internal/clients/handlers/basehandler.go97
-rw-r--r--internal/clients/handlers/clienthandler.go4
-rw-r--r--internal/clients/handlers/healthhandler.go105
-rw-r--r--internal/clients/handlers/maprhandler.go56
-rw-r--r--internal/clients/healthclient.go114
-rw-r--r--internal/clients/maker.go3
-rw-r--r--internal/clients/maprclient.go84
-rw-r--r--internal/clients/remote/connection.go212
-rw-r--r--internal/clients/stats.go83
-rw-r--r--internal/clients/tailclient.go16
-rw-r--r--internal/color/brush/brush.go194
-rw-r--r--internal/color/color.go174
-rw-r--r--internal/color/color_test.go53
-rw-r--r--internal/color/colorfy.go58
-rw-r--r--internal/color/paint.go91
-rw-r--r--internal/color/table.go53
-rw-r--r--internal/config/args.go162
-rw-r--r--internal/config/client.go191
-rw-r--r--internal/config/common.go28
-rw-r--r--internal/config/config.go71
-rw-r--r--internal/config/initializer.go171
-rw-r--r--internal/config/read.go37
-rw-r--r--internal/config/server.go10
-rw-r--r--internal/discovery/comma.go5
-rw-r--r--internal/discovery/discovery.go31
-rw-r--r--internal/discovery/file.go9
-rw-r--r--internal/done.go10
-rw-r--r--internal/io/dlog/dlog.go272
-rw-r--r--internal/io/dlog/level.go84
-rw-r--r--internal/io/dlog/loggers/factory.go54
-rw-r--r--internal/io/dlog/loggers/file.go165
-rw-r--r--internal/io/dlog/loggers/fout.go46
-rw-r--r--internal/io/dlog/loggers/logger.go19
-rw-r--r--internal/io/dlog/loggers/none.go21
-rw-r--r--internal/io/dlog/loggers/stdout.go54
-rw-r--r--internal/io/dlog/loggers/strategy.go34
-rw-r--r--internal/io/dlog/rotation.go27
-rw-r--r--internal/io/fs/catfile.go4
-rw-r--r--internal/io/fs/filereader.go7
-rw-r--r--internal/io/fs/permissions/permission.go6
-rw-r--r--internal/io/fs/permissions/permission_linuxacl.c (renamed from internal/io/fs/permissions/permission_linux.c)4
-rw-r--r--internal/io/fs/permissions/permission_linuxacl.go (renamed from internal/io/fs/permissions/permission_linux.go)6
-rw-r--r--internal/io/fs/permissions/permission_linuxacl.h (renamed from internal/io/fs/permissions/permission_linux.h)2
-rw-r--r--internal/io/fs/permissions/permission_test.go2
-rw-r--r--internal/io/fs/readfile.go284
-rw-r--r--internal/io/fs/tailfile.go4
-rw-r--r--internal/io/line/line.go11
-rw-r--r--internal/io/logger/logger.go397
-rw-r--r--internal/io/logger/modes.go11
-rw-r--r--internal/io/logger/strategy.go22
-rw-r--r--internal/io/pool/builder.go21
-rw-r--r--internal/io/pool/bytesbuffer.go22
-rw-r--r--internal/io/prompt/prompt.go13
-rw-r--r--internal/io/signal/signal.go8
-rw-r--r--internal/lcontext/lcontext.go22
-rw-r--r--internal/mapr/aggregateset.go29
-rw-r--r--internal/mapr/client/aggregate.go29
-rw-r--r--internal/mapr/funcs/function.go12
-rw-r--r--internal/mapr/funcs/function_test.go21
-rw-r--r--internal/mapr/funcs/maskdigits.go2
-rw-r--r--internal/mapr/globalgroupset.go11
-rw-r--r--internal/mapr/groupset.go291
-rw-r--r--internal/mapr/logformat/default.go41
-rw-r--r--internal/mapr/logformat/default_test.go88
-rw-r--r--internal/mapr/logformat/generickv.go31
-rw-r--r--internal/mapr/logformat/parser.go15
-rw-r--r--internal/mapr/query.go27
-rw-r--r--internal/mapr/query_test.go125
-rw-r--r--internal/mapr/selectcondition.go11
-rw-r--r--internal/mapr/server/aggregate.go155
-rw-r--r--internal/mapr/setclause.go2
-rw-r--r--internal/mapr/setcondition.go15
-rw-r--r--internal/mapr/token.go18
-rw-r--r--internal/mapr/whereclause.go16
-rw-r--r--internal/mapr/wherecondition.go50
-rw-r--r--internal/omode/mode.go3
-rw-r--r--internal/protocol/protocol.go18
-rw-r--r--internal/regex/regex.go20
-rw-r--r--internal/regex/regex_test.go33
-rw-r--r--internal/server/continuous.go36
-rw-r--r--internal/server/handlers/basehandler.go320
-rw-r--r--internal/server/handlers/controlhandler.go100
-rw-r--r--internal/server/handlers/healthhandler.go58
-rw-r--r--internal/server/handlers/mapcommand.go7
-rw-r--r--internal/server/handlers/readcommand.go97
-rw-r--r--internal/server/handlers/serverhandler.go373
-rw-r--r--internal/server/scheduler.go37
-rw-r--r--internal/server/server.go131
-rw-r--r--internal/server/stats.go21
-rw-r--r--internal/source/source.go30
-rw-r--r--internal/ssh/client/authmethods.go53
-rw-r--r--internal/ssh/client/customkeycallback.go3
-rw-r--r--internal/ssh/client/knownhostscallback.go33
-rw-r--r--internal/ssh/server/hostkey.go18
-rw-r--r--internal/ssh/server/publickeycallback.go38
-rw-r--r--internal/ssh/ssh.go11
-rw-r--r--internal/user/name.go3
-rw-r--r--internal/user/server/user.go57
-rw-r--r--internal/version/version.go34
106 files changed, 4292 insertions, 2461 deletions
diff --git a/internal/clients/args.go b/internal/clients/args.go
deleted file mode 100644
index 34fcfa2..0000000
--- a/internal/clients/args.go
+++ /dev/null
@@ -1,25 +0,0 @@
-package clients
-
-import (
- "github.com/mimecast/dtail/internal/omode"
-
- gossh "golang.org/x/crypto/ssh"
-)
-
-// Args is a helper struct to summarize common client arguments.
-type Args struct {
- Mode omode.Mode
- ServersStr string
- UserName string
- What string
- Arguments []string
- RegexStr string
- RegexInvert bool
- TrustAllHosts bool
- Discovery string
- ConnectionsPerCPU int
- Timeout int
- SSHAuthMethods []gossh.AuthMethod
- SSHHostKeyCallback gossh.HostKeyCallback
- PrivateKeyPathFile string
-}
diff --git a/internal/clients/baseclient.go b/internal/clients/baseclient.go
index 69055a3..4a7bd84 100644
--- a/internal/clients/baseclient.go
+++ b/internal/clients/baseclient.go
@@ -5,10 +5,10 @@ import (
"sync"
"time"
- "github.com/mimecast/dtail/internal/clients/remote"
+ "github.com/mimecast/dtail/internal/clients/connectors"
+ "github.com/mimecast/dtail/internal/config"
"github.com/mimecast/dtail/internal/discovery"
- "github.com/mimecast/dtail/internal/io/logger"
- "github.com/mimecast/dtail/internal/omode"
+ "github.com/mimecast/dtail/internal/io/dlog"
"github.com/mimecast/dtail/internal/regex"
"github.com/mimecast/dtail/internal/ssh/client"
@@ -17,13 +17,13 @@ import (
// This is the main client data structure.
type baseClient struct {
- Args
+ config.Args
// To display client side stats
stats *stats
// List of remote servers to connect to.
servers []string
// We have one connection per remote server.
- connections []*remote.Connection
+ connections []connectors.Connector
// SSH auth methods to use to connect to the remote servers.
sshAuthMethods []gossh.AuthMethod
// To deal with SSH host keys
@@ -39,7 +39,7 @@ type baseClient struct {
}
func (c *baseClient) init() {
- logger.Info("Initiating base client")
+ dlog.Client.Debug("Initiating base client", c.Args.String())
flag := regex.Default
if c.Args.RegexInvert {
@@ -47,12 +47,16 @@ func (c *baseClient) init() {
}
regex, err := regex.New(c.Args.RegexStr, flag)
if err != nil {
- logger.FatalExit(c.Regex, "invalid regex!", err, regex)
+ dlog.Client.FatalPanic(c.Regex, "Invalid regex!", err, regex)
}
c.Regex = regex
- logger.Debug("Regex", c.Regex)
- c.sshAuthMethods, c.hostKeyCallback = client.InitSSHAuthMethods(c.Args.SSHAuthMethods, c.Args.SSHHostKeyCallback, c.Args.TrustAllHosts, c.throttleCh, c.Args.PrivateKeyPathFile)
+ if c.Args.Serverless {
+ return
+ }
+ c.sshAuthMethods, c.hostKeyCallback = client.InitSSHAuthMethods(
+ c.Args.SSHAuthMethods, c.Args.SSHHostKeyCallback, c.Args.TrustAllHosts,
+ c.throttleCh, c.Args.PrivateKeyPathFile)
}
func (c *baseClient) makeConnections(maker maker) {
@@ -60,26 +64,31 @@ func (c *baseClient) makeConnections(maker maker) {
discoveryService := discovery.New(c.Discovery, c.ServersStr, discovery.Shuffle)
for _, server := range discoveryService.ServerList() {
- c.connections = append(c.connections, c.makeConnection(server, c.sshAuthMethods, c.hostKeyCallback))
+ c.connections = append(c.connections, c.makeConnection(server,
+ c.sshAuthMethods, c.hostKeyCallback))
}
c.stats = newTailStats(len(c.connections))
}
func (c *baseClient) Start(ctx context.Context, statsCh <-chan string) (status int) {
- // Periodically check for unknown hosts, and ask the user whether to trust them or not.
- go c.hostKeyCallback.PromptAddHosts(ctx)
+ dlog.Client.Trace("Starting base client")
+ // Can be nil when serverless.
+ if c.hostKeyCallback != nil {
+ // Periodically check for unknown hosts, and ask the user whether to trust them or not.
+ go c.hostKeyCallback.PromptAddHosts(ctx)
+ }
// Print client stats every time something on statsCh is recieved.
- go c.stats.Start(ctx, c.throttleCh, statsCh)
- // Keep count of active connections
- active := make(chan struct{}, len(c.connections))
+ go c.stats.Start(ctx, c.throttleCh, statsCh, c.Args.Quiet)
+ 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) {
- connStatus := c.start(ctx, active, i, conn)
- // Update global status.
+ for i, conn := range c.connections {
+ go func(i int, conn connectors.Connector) {
+ defer wg.Done()
+ connStatus := c.startConnection(ctx, i, conn)
mutex.Lock()
defer mutex.Unlock()
if connStatus > status {
@@ -88,15 +97,12 @@ func (c *baseClient) Start(ctx context.Context, statsCh <-chan string) (status i
}(i, conn)
}
- c.waitUntilDone(ctx, active)
+ wg.Wait()
return
}
-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 }()
+func (c *baseClient) startConnection(ctx context.Context, i int,
+ conn connectors.Connector) (status int) {
for {
connCtx, cancel := context.WithCancel(ctx)
@@ -104,46 +110,25 @@ func (c *baseClient) start(ctx context.Context, active chan struct{}, i int, con
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()
+ 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)
+ dlog.Client.Debug(conn.Server(), "Reconnecting")
+ conn = c.makeConnection(conn.Server(), c.sshAuthMethods, c.hostKeyCallback)
c.connections[i] = conn
}
}
-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) 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 && c.retry {
- <-ctx.Done()
- }
-
- for {
- numActive := len(active)
- if numActive == 0 {
- return
- }
- logger.Debug("Active connections", numActive)
- time.Sleep(time.Second)
+func (c *baseClient) makeConnection(server string, sshAuthMethods []gossh.AuthMethod,
+ hostKeyCallback client.HostKeyCallback) connectors.Connector {
+ if c.Args.Serverless {
+ return connectors.NewServerless(c.UserName, c.maker.makeHandler(server),
+ c.maker.makeCommands())
}
+ return connectors.NewServerConnection(server, c.UserName, sshAuthMethods,
+ hostKeyCallback, c.maker.makeHandler(server), c.maker.makeCommands())
}
diff --git a/internal/clients/catclient.go b/internal/clients/catclient.go
index d8e9196..bd65560 100644
--- a/internal/clients/catclient.go
+++ b/internal/clients/catclient.go
@@ -7,6 +7,8 @@ import (
"strings"
"github.com/mimecast/dtail/internal/clients/handlers"
+ "github.com/mimecast/dtail/internal/config"
+ "github.com/mimecast/dtail/internal/io/dlog"
"github.com/mimecast/dtail/internal/omode"
)
@@ -16,11 +18,10 @@ type CatClient struct {
}
// NewCatClient returns a new cat client.
-func NewCatClient(args Args) (*CatClient, error) {
+func NewCatClient(args config.Args) (*CatClient, error) {
if args.RegexStr != "" {
return nil, errors.New("Can't use regex with 'cat' operating mode")
}
-
args.Mode = omode.CatClient
c := CatClient{
@@ -33,7 +34,6 @@ func NewCatClient(args Args) (*CatClient, error) {
c.init()
c.makeConnections(c)
-
return &c, nil
}
@@ -42,8 +42,13 @@ func (c CatClient) makeHandler(server string) handlers.Handler {
}
func (c CatClient) makeCommands() (commands []string) {
+ regex, err := c.Regex.Serialize()
+ if err != nil {
+ dlog.Client.FatalPanic(err)
+ }
for _, file := range strings.Split(c.What, ",") {
- commands = append(commands, fmt.Sprintf("%s %s %s", c.Mode.String(), file, c.Regex.Serialize()))
+ commands = append(commands, fmt.Sprintf("%s:%s %s %s",
+ c.Mode.String(), c.Args.SerializeOptions(), file, regex))
}
return
}
diff --git a/internal/clients/connectors/connector.go b/internal/clients/connectors/connector.go
new file mode 100644
index 0000000..3ab6a08
--- /dev/null
+++ b/internal/clients/connectors/connector.go
@@ -0,0 +1,17 @@
+package connectors
+
+import (
+ "context"
+
+ "github.com/mimecast/dtail/internal/clients/handlers"
+)
+
+// Connector interface.
+type Connector interface {
+ // Start the connection.
+ Start(ctx context.Context, cancel context.CancelFunc, throttleCh, statsCh chan struct{})
+ // Server hostname.
+ Server() string
+ // Handler for the connection.
+ Handler() handlers.Handler
+}
diff --git a/internal/clients/connectors/serverconnection.go b/internal/clients/connectors/serverconnection.go
new file mode 100644
index 0000000..1df4d73
--- /dev/null
+++ b/internal/clients/connectors/serverconnection.go
@@ -0,0 +1,203 @@
+package connectors
+
+import (
+ "context"
+ "fmt"
+ "io"