diff --git a/cmd/config.go b/cmd/config.go index 4d85863..f7fb8c7 100644 --- a/cmd/config.go +++ b/cmd/config.go @@ -3,7 +3,6 @@ package cmd import ( "io" - "github.com/zrepl/zrepl/rpc" "github.com/zrepl/zrepl/zfs" ) @@ -20,8 +19,18 @@ type Global struct { } } -type RPCConnecter interface { - Connect() (rpc.RPCClient, error) +type JobDebugSettings struct { + Conn struct { + ReadDump string `mapstructure:"read_dump"` + WriteDump string `mapstructure:"write_dump"` + } + RPC struct { + Log bool + } +} + +type RWCConnecter interface { + Connect() (io.ReadWriteCloser, error) } type AuthenticatedChannelListenerFactory interface { Listen() (AuthenticatedChannelListener, error) diff --git a/cmd/config_connect.go b/cmd/config_connect.go index ee98154..35be638 100644 --- a/cmd/config_connect.go +++ b/cmd/config_connect.go @@ -7,9 +7,7 @@ import ( "github.com/jinzhu/copier" "github.com/mitchellh/mapstructure" "github.com/pkg/errors" - "github.com/zrepl/zrepl/rpc" "github.com/zrepl/zrepl/sshbytestream" - "github.com/zrepl/zrepl/util" ) type SSHStdinserverConnecter struct { @@ -20,8 +18,6 @@ type SSHStdinserverConnecter struct { TransportOpenCommand []string `mapstructure:"transport_open_command"` SSHCommand string `mapstructure:"ssh_command"` Options []string - ConnLogReadFile string `mapstructure:"connlog_read_file"` - ConnLogWriteFile string `mapstructure:"connlog_write_file"` } func parseSSHStdinserverConnecter(i map[string]interface{}) (c *SSHStdinserverConnecter, err error) { @@ -37,20 +33,14 @@ func parseSSHStdinserverConnecter(i map[string]interface{}) (c *SSHStdinserverCo } -func (c *SSHStdinserverConnecter) Connect() (client rpc.RPCClient, err error) { - var stream io.ReadWriteCloser +func (c *SSHStdinserverConnecter) Connect() (rwc io.ReadWriteCloser, err error) { var rpcTransport sshbytestream.SSHTransport if err = copier.Copy(&rpcTransport, c); err != nil { return } - if stream, err = sshbytestream.Outgoing(rpcTransport); err != nil { + if rwc, err = sshbytestream.Outgoing(rpcTransport); err != nil { + err = errors.WithStack(err) return } - stream, err = util.NewReadWriteCloserLogger(stream, c.ConnLogReadFile, c.ConnLogWriteFile) - if err != nil { - return - } - client = rpc.NewClient(stream) - return client, nil - + return } diff --git a/cmd/config_job_local.go b/cmd/config_job_local.go index 5805540..ce85a1c 100644 --- a/cmd/config_job_local.go +++ b/cmd/config_job_local.go @@ -3,30 +3,32 @@ package cmd import ( "time" - "github.com/zrepl/zrepl/rpc" "github.com/mitchellh/mapstructure" "github.com/pkg/errors" + "github.com/zrepl/zrepl/rpc" ) type LocalJob struct { Name string Mapping *DatasetMapFilter - SnapshotFilter *PrefixSnapshotFilter + SnapshotFilter *PrefixSnapshotFilter Interval time.Duration InitialReplPolicy InitialReplPolicy PruneLHS PrunePolicy PruneRHS PrunePolicy + Debug JobDebugSettings } func parseLocalJob(name string, i map[string]interface{}) (j *LocalJob, err error) { var asMap struct { - Mapping map[string]string - SnapshotPrefix string `mapstructure:"snapshot_prefix"` - Interval string - InitialReplPolicy string `mapstructure:"initial_repl_policy"` - PruneLHS map[string]interface{} `mapstructure:"prune_lhs"` - PruneRHS map[string]interface{} `mapstructure:"prune_rhs"` + Mapping map[string]string + SnapshotPrefix string `mapstructure:"snapshot_prefix"` + Interval string + InitialReplPolicy string `mapstructure:"initial_repl_policy"` + PruneLHS map[string]interface{} `mapstructure:"prune_lhs"` + PruneRHS map[string]interface{} `mapstructure:"prune_rhs"` + Debug map[string]interface{} } if err = mapstructure.Decode(i, &asMap); err != nil { @@ -62,6 +64,11 @@ func parseLocalJob(name string, i map[string]interface{}) (j *LocalJob, err erro return } + if err = mapstructure.Decode(asMap.Debug, &j.Debug); err != nil { + err = errors.Wrap(err, "cannot parse 'debug'") + return + } + return } diff --git a/cmd/config_job_pull.go b/cmd/config_job_pull.go index 75a7ee6..ea2501d 100644 --- a/cmd/config_job_pull.go +++ b/cmd/config_job_pull.go @@ -5,15 +5,18 @@ import ( "github.com/mitchellh/mapstructure" "github.com/pkg/errors" + "github.com/zrepl/zrepl/rpc" + "github.com/zrepl/zrepl/util" ) type PullJob struct { Name string - Connect RPCConnecter + Connect RWCConnecter Mapping *DatasetMapFilter SnapshotFilter *PrefixSnapshotFilter InitialReplPolicy InitialReplPolicy Prune PrunePolicy + Debug JobDebugSettings } func parsePullJob(name string, i map[string]interface{}) (j *PullJob, err error) { @@ -24,6 +27,7 @@ func parsePullJob(name string, i map[string]interface{}) (j *PullJob, err error) InitialReplPolicy string `mapstructure:"initial_repl_policy"` Prune map[string]interface{} SnapshotPrefix string `mapstructure:"snapshot_prefix"` + Debug map[string]interface{} } if err = mapstructure.Decode(i, &asMap); err != nil { @@ -60,6 +64,11 @@ func parsePullJob(name string, i map[string]interface{}) (j *PullJob, err error) return } + if err = mapstructure.Decode(asMap.Debug, &j.Debug); err != nil { + err = errors.Wrap(err, "cannot parse 'debug'") + return + } + return } @@ -68,11 +77,23 @@ func (j *PullJob) JobName() string { } func (j *PullJob) JobDo(log Logger) (err error) { - client, err := j.Connect.Connect() + + rwc, err := j.Connect.Connect() if err != nil { log.Printf("error connect: %s", err) return err } + + rwc, err = util.NewReadWriteCloserLogger(rwc, j.Debug.Conn.ReadDump, j.Debug.Conn.WriteDump) + if err != nil { + return + } + + client := rpc.NewClient(rwc) + if j.Debug.RPC.Log { + client.SetLogger(log, true) + } + defer closeRPCWithTimeout(log, client, time.Second*10, "") return doPull(PullContext{client, log, j.Mapping, j.InitialReplPolicy}) } diff --git a/cmd/config_job_source.go b/cmd/config_job_source.go index 257e8ee..f25415a 100644 --- a/cmd/config_job_source.go +++ b/cmd/config_job_source.go @@ -1,15 +1,15 @@ package cmd import ( - "time" - mapstructure "github.com/mitchellh/mapstructure" "github.com/pkg/errors" "github.com/zrepl/zrepl/rpc" + "github.com/zrepl/zrepl/util" "io" "os" "os/signal" "syscall" + "time" ) type SourceJob struct { @@ -19,6 +19,7 @@ type SourceJob struct { SnapshotFilter *PrefixSnapshotFilter Interval time.Duration Prune PrunePolicy + Debug JobDebugSettings } func parseSourceJob(name string, i map[string]interface{}) (j *SourceJob, err error) { @@ -29,6 +30,7 @@ func parseSourceJob(name string, i map[string]interface{}) (j *SourceJob, err er SnapshotPrefix string `mapstructure:"snapshot_prefix"` Interval string Prune map[string]interface{} + Debug map[string]interface{} } if err = mapstructure.Decode(i, &asMap); err != nil { @@ -59,6 +61,11 @@ func parseSourceJob(name string, i map[string]interface{}) (j *SourceJob, err er return } + if err = mapstructure.Decode(asMap.Debug, &j.Debug); err != nil { + err = errors.Wrap(err, "cannot parse 'debug'") + return + } + return } @@ -100,6 +107,11 @@ outer: break outer // closed because of accept error } + rwc, err := util.NewReadWriteCloserLogger(rwc, j.Debug.Conn.ReadDump, j.Debug.Conn.WriteDump) + if err != nil { + panic(err) + } + // construct connection handler handler := Handler{ Logger: log, @@ -108,6 +120,10 @@ outer: // handle connection rpcServer := rpc.NewServer(rwc) + if j.Debug.RPC.Log { + rpclog := util.NewPrefixLogger(log, "rpc") + rpcServer.SetLogger(rpclog, true) + } registerEndpoints(rpcServer, handler) if err = rpcServer.Serve(); err != nil { log.Printf("error serving connection: %s", err) diff --git a/cmd/config_parse.go b/cmd/config_parse.go index 0aa1930..47524ab 100644 --- a/cmd/config_parse.go +++ b/cmd/config_parse.go @@ -115,7 +115,7 @@ func parseJob(i map[string]interface{}) (j Job, err error) { } -func parseConnect(i map[string]interface{}) (c RPCConnecter, err error) { +func parseConnect(i map[string]interface{}) (c RWCConnecter, err error) { t, err := extractStringField(i, "type", true) if err != nil { diff --git a/cmd/config_serve_stdinserver.go b/cmd/config_serve_stdinserver.go index e8f6f88..3f411b5 100644 --- a/cmd/config_serve_stdinserver.go +++ b/cmd/config_serve_stdinserver.go @@ -4,7 +4,6 @@ import ( "github.com/ftrvxmtrx/fd" "github.com/mitchellh/mapstructure" "github.com/pkg/errors" - "github.com/zrepl/zrepl/util" "io" "net" "os" @@ -13,9 +12,7 @@ import ( ) type StdinserverListenerFactory struct { - ClientIdentity string `mapstructure:"client_identity"` - ConnLogReadFile string `mapstructure:"connlog_read_file"` - ConnLogWriteFile string `mapstructure:"connlog_write_file"` + ClientIdentity string `mapstructure:"client_identity"` } func parseStdinserverListenerFactory(i map[string]interface{}) (f *StdinserverListenerFactory, err error) { @@ -63,15 +60,13 @@ func (f *StdinserverListenerFactory) Listen() (al AuthenticatedChannelListener, return nil, errors.Wrapf(err, "cannot listen on unix socket %s", unixaddr) } - l := &StdinserverListener{ul, f.ConnLogReadFile, f.ConnLogWriteFile} + l := &StdinserverListener{ul} return l, nil } type StdinserverListener struct { - l *net.UnixListener - ConnLogReadFile string `mapstructure:"connlog_read_file"` - ConnLogWriteFile string `mapstructure:"connlog_write_file"` + l *net.UnixListener } type fdRWC struct { @@ -109,12 +104,7 @@ func (l *StdinserverListener) Accept() (ch io.ReadWriteCloser, err error) { rwc := fdRWC{files[0], files[1], c.(*net.UnixConn)} - rwclog, err := util.NewReadWriteCloserLogger(rwc, l.ConnLogReadFile, l.ConnLogWriteFile) - if err != nil { - panic(err) - } - - return rwclog, nil + return rwc, nil } diff --git a/cmd/config_test.go b/cmd/config_test.go index a775831..09b3826 100644 --- a/cmd/config_test.go +++ b/cmd/config_test.go @@ -11,10 +11,11 @@ import ( func TestSampleConfigsAreParsedWithoutErrors(t *testing.T) { - paths:= []string{ + paths := []string{ "./sampleconf/localbackup/host1.yml", "./sampleconf/pullbackup/backuphost.yml", "./sampleconf/pullbackup/productionhost.yml", + "./sampleconf/random/debugging.yml", } for _, p := range paths { diff --git a/cmd/sampleconf/random/debugging.yml b/cmd/sampleconf/random/debugging.yml new file mode 100644 index 0000000..4030db8 --- /dev/null +++ b/cmd/sampleconf/random/debugging.yml @@ -0,0 +1,33 @@ +global: + serve: + stdinserver: + sockdir: /var/run/zrepl/stdinserver + +jobs: + +- name: debian2_pull + # JOB DEBUGGING OPTIONS + # should be equal for all job types, but each job implements the debugging itself + # => consult job documentation for supported options + debug: + conn: # debug the io.ReadWriteCloser connection + read_dump: /tmp/connlog_read # dump results of Read() invocations to this file + write_dump: /tmp/connlog_write # dump results of Write() invocations to this file + rpc: # debug the RPC protocol implementation + log: true # log output from rpc layer to the job log + + # ... just to make the unit tests pass. + # check other examples, e.g. localbackup or pullbackup for what the sutff below means + type: source + serve: + type: stdinserver + client_identity: debian2 + datasets: { + "pool1/db<": ok + } + snapshot_prefix: zrepl_ + interval: 1s + prune: + policy: grid + grid: 1x10s(keep=all) +