diff --git a/cmd/config.go b/cmd/config.go index 5b2bcdd..1c0786c 100644 --- a/cmd/config.go +++ b/cmd/config.go @@ -15,10 +15,6 @@ import ( "strings" ) -const LOCAL_TRANSPORT_IDENTITY string = "local" - -const DEFAULT_INITIAL_REPL_POLICY = InitialReplPolicyMostRecent - type Pool struct { Name string Transport Transport @@ -42,22 +38,15 @@ type SSHTransport struct { ConnLogWriteFile string `mapstructure:"connlog_write_file"` } -type InitialReplPolicy string - -const ( - InitialReplPolicyMostRecent InitialReplPolicy = "most_recent" - InitialReplPolicyAll InitialReplPolicy = "all" -) - type Push struct { To *Pool Datasets []zfs.DatasetPath - InitialReplPolicy InitialReplPolicy + InitialReplPolicy rpc.InitialReplPolicy } type Pull struct { From *Pool Mapping zfs.DatasetMapping - InitialReplPolicy InitialReplPolicy + InitialReplPolicy rpc.InitialReplPolicy } type ClientMapping struct { From string @@ -131,8 +120,8 @@ func parsePools(v interface{}) (pools []Pool, err error) { pools = make([]Pool, len(asList)) for i, p := range asList { - if p.Name == LOCAL_TRANSPORT_IDENTITY { - err = errors.New(fmt.Sprintf("pool name '%s' reserved for local pulls", LOCAL_TRANSPORT_IDENTITY)) + if p.Name == rpc.LOCAL_TRANSPORT_IDENTITY { + err = errors.New(fmt.Sprintf("pool name '%s' reserved for local pulls", rpc.LOCAL_TRANSPORT_IDENTITY)) return } @@ -207,7 +196,7 @@ func parsePushs(v interface{}, pl poolLookup) (p []Push, err error) { } } - if push.InitialReplPolicy, err = parseInitialReplPolicy(e.InitialReplPolicy, DEFAULT_INITIAL_REPL_POLICY); err != nil { + if push.InitialReplPolicy, err = parseInitialReplPolicy(e.InitialReplPolicy, rpc.DEFAULT_INITIAL_REPL_POLICY); err != nil { return } @@ -235,7 +224,7 @@ func parsePulls(v interface{}, pl poolLookup) (p []Pull, err error) { var fromPool *Pool - if e.From == LOCAL_TRANSPORT_IDENTITY { + if e.From == rpc.LOCAL_TRANSPORT_IDENTITY { fromPool = &Pool{ Name: "local", Transport: LocalTransport{}, @@ -252,7 +241,7 @@ func parsePulls(v interface{}, pl poolLookup) (p []Pull, err error) { if pull.Mapping, err = parseComboMapping(e.Mapping); err != nil { return } - if pull.InitialReplPolicy, err = parseInitialReplPolicy(e.InitialReplPolicy, DEFAULT_INITIAL_REPL_POLICY); err != nil { + if pull.InitialReplPolicy, err = parseInitialReplPolicy(e.InitialReplPolicy, rpc.DEFAULT_INITIAL_REPL_POLICY); err != nil { return } @@ -262,7 +251,7 @@ func parsePulls(v interface{}, pl poolLookup) (p []Pull, err error) { return } -func parseInitialReplPolicy(v interface{}, defaultPolicy InitialReplPolicy) (p InitialReplPolicy, err error) { +func parseInitialReplPolicy(v interface{}, defaultPolicy rpc.InitialReplPolicy) (p rpc.InitialReplPolicy, err error) { s, ok := v.(string) if !ok { goto err @@ -272,9 +261,9 @@ func parseInitialReplPolicy(v interface{}, defaultPolicy InitialReplPolicy) (p I case s == "": p = defaultPolicy case s == "most_recent": - p = InitialReplPolicyMostRecent + p = rpc.InitialReplPolicyMostRecent case s == "all": - p = InitialReplPolicyAll + p = rpc.InitialReplPolicyAll default: goto err } diff --git a/cmd/main.go b/cmd/main.go index a8ad382..915e8a3 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -158,7 +158,7 @@ func cmdRun(c *cli.Context) error { Repeats: true, RunFunc: func(log jobrun.Logger) error { log.Printf("doing pull: %v", pull) - return doPull(pull, c, log) + return jobPull(pull, c, log) }, } @@ -193,33 +193,7 @@ func cmdRun(c *cli.Context) error { return nil } -func closeRPCWithTimeout(log Logger, remote rpc.RPCRequester, timeout time.Duration, goodbye string) { - log.Printf("closing rpc connection") - - ch := make(chan error) - go func() { - ch <- remote.CloseRequest(rpc.CloseRequest{goodbye}) - }() - - var err error - select { - case <-time.After(timeout): - err = fmt.Errorf("timeout exceeded (%s)", timeout) - case closeRequestErr := <-ch: - err = closeRequestErr - } - - if err != nil { - log.Printf("error closing connection: %s", err) - err = remote.ForceClose() - if err != nil { - log.Printf("error force-closing connection: %s", err) - } - } - return -} - -func doPull(pull Pull, c *cli.Context, log jobrun.Logger) (err error) { +func jobPull(pull Pull, c *cli.Context, log jobrun.Logger) (err error) { if lt, ok := pull.From.Transport.(LocalTransport); ok { lt.SetHandler(Handler{ @@ -238,185 +212,5 @@ func doPull(pull Pull, c *cli.Context, log jobrun.Logger) (err error) { defer closeRPCWithTimeout(log, remote, time.Second*10, "") - fsr := rpc.FilesystemRequest{ - Direction: rpc.DirectionPull, - } - var remoteFilesystems []zfs.DatasetPath - if remoteFilesystems, err = remote.FilesystemRequest(fsr); err != nil { - return - } - - type RemoteLocalMapping struct { - Remote zfs.DatasetPath - Local zfs.DatasetPath - LocalExists bool - } - replMapping := make(map[string]RemoteLocalMapping, len(remoteFilesystems)) - localTraversal := zfs.NewDatasetPathForest() - localExists, err := zfs.ZFSListFilesystemExists() - if err != nil { - log.Printf("cannot get local filesystems map: %s", err) - return err - } - - { - - log.Printf("mapping using %#v\n", pull.Mapping) - for fs := range remoteFilesystems { - var err error - var localFs zfs.DatasetPath - localFs, err = pull.Mapping.Map(remoteFilesystems[fs]) - if err != nil { - if err != zfs.NoMatchError { - log.Printf("error mapping %s: %#v\n", remoteFilesystems[fs], err) - return err - } - continue - } - m := RemoteLocalMapping{remoteFilesystems[fs], localFs, localExists(localFs)} - replMapping[m.Local.ToString()] = m - localTraversal.Add(m.Local) - } - - } - - log.Printf("remoteFilesystems: %#v\nreplMapping: %#v\n", remoteFilesystems, replMapping) - - // per fs sync, assume sorted in top-down order TODO - - localTraversal.WalkTopDown(func(v zfs.DatasetPathVisit) bool { - - if v.FilledIn { - if localExists(v.Path) { - return true - } - log.Printf("creating fill-in dataset %s", v.Path) - return false - } - - m, ok := replMapping[v.Path.ToString()] - if !ok { - panic("internal inconsistency: replMapping should contain mapping for any path that was not filled in by WalkTopDown()") - } - - log := func(format string, args ...interface{}) { - log.Printf("[%s => %s]: %s", m.Remote.ToString(), m.Local.ToString(), fmt.Sprintf(format, args...)) - } - - log("mapping: %#v\n", m) - - var versions []zfs.FilesystemVersion - if m.LocalExists { - if versions, err = zfs.ZFSListFilesystemVersions(m.Local); err != nil { - log("cannot get filesystem versions, stopping...: %v\n", m.Local.ToString(), m, err) - return false - } - } - - var theirVersions []zfs.FilesystemVersion - theirVersions, err = remote.FilesystemVersionsRequest(rpc.FilesystemVersionsRequest{ - Filesystem: m.Remote, - }) - if err != nil { - log("cannot fetch remote filesystem versions, stopping: %s", err) - return false - } - - diff := zfs.MakeFilesystemDiff(versions, theirVersions) - log("diff: %#v\n", diff) - - if diff.IncrementalPath == nil { - log("performing initial sync, following policy: %#v", pull.InitialReplPolicy) - - if pull.InitialReplPolicy != InitialReplPolicyMostRecent { - panic(fmt.Sprintf("policy %#v not implemented", pull.InitialReplPolicy)) - } - - snapsOnly := make([]zfs.FilesystemVersion, 0, len(diff.MRCAPathRight)) - for s := range diff.MRCAPathRight { - if diff.MRCAPathRight[s].Type == zfs.Snapshot { - snapsOnly = append(snapsOnly, diff.MRCAPathRight[s]) - } - } - - if len(snapsOnly) < 1 { - log("cannot perform initial sync: no remote snapshots. stopping...") - return false - } - - r := rpc.InitialTransferRequest{ - Filesystem: m.Remote, - FilesystemVersion: snapsOnly[len(snapsOnly)-1], - } - - log("requesting initial transfer") - - var stream io.Reader - if stream, err = remote.InitialTransferRequest(r); err != nil { - log("error initial transfer request, stopping...: %s", err) - return false - } - - log("received initial transfer request response. zfs recv...") - - if err = zfs.ZFSRecv(m.Local, stream, "-u"); err != nil { - log("error receiving stream, stopping...: %s", err) - return false - } - - log("configuring properties of received filesystem") - - if err = zfs.ZFSSet(m.Local, "readonly", "on"); err != nil { - - } - - log("finished initial transfer") - - } else if len(diff.IncrementalPath) < 2 { - log("remote and local are in sync") - } else { - - log("incremental transfers using path: %#v", diff.IncrementalPath) - - for i := 0; i < len(diff.IncrementalPath)-1; i++ { - - from, to := diff.IncrementalPath[i], diff.IncrementalPath[i+1] - - log := func(format string, args ...interface{}) { - log("[%s => %s]: %s", from.Name, to.Name, fmt.Sprintf(format, args...)) - } - - r := rpc.IncrementalTransferRequest{ - Filesystem: m.Remote, - From: from, - To: to, - } - log("requesting incremental transfer: %#v", r) - - var stream io.Reader - if stream, err = remote.IncrementalTransferRequest(r); err != nil { - log("error requesting incremental transfer, stopping...: %s", err.Error()) - return false - } - - log("receving incremental transfer") - - if err = zfs.ZFSRecv(m.Local, stream); err != nil { - log("error receiving stream, stopping...: %s", err) - return false - } - - log("finished incremental transfer") - - } - - log("finished incremental transfer path") - - } - - return true - - }) - - return nil + return doPull(PullContext{remote, log, pull.Mapping, pull.InitialReplPolicy}) } diff --git a/cmd/replication.go b/cmd/replication.go new file mode 100644 index 0000000..25bb7c2 --- /dev/null +++ b/cmd/replication.go @@ -0,0 +1,232 @@ +package main + +import ( + "fmt" + "github.com/zrepl/zrepl/rpc" + "github.com/zrepl/zrepl/zfs" + "io" + "time" +) + +func closeRPCWithTimeout(log Logger, remote rpc.RPCRequester, timeout time.Duration, goodbye string) { + log.Printf("closing rpc connection") + + ch := make(chan error) + go func() { + ch <- remote.CloseRequest(rpc.CloseRequest{goodbye}) + }() + + var err error + select { + case <-time.After(timeout): + err = fmt.Errorf("timeout exceeded (%s)", timeout) + case closeRequestErr := <-ch: + err = closeRequestErr + } + + if err != nil { + log.Printf("error closing connection: %s", err) + err = remote.ForceClose() + if err != nil { + log.Printf("error force-closing connection: %s", err) + } + } + return +} + +type PullContext struct { + Remote rpc.RPCRequester + Log Logger + Mapping zfs.DatasetMapping + InitialReplPolicy rpc.InitialReplPolicy +} + +func doPull(pull PullContext) (err error) { + + remote := pull.Remote + log := pull.Log + + fsr := rpc.FilesystemRequest{ + Direction: rpc.DirectionPull, + } + var remoteFilesystems []zfs.DatasetPath + if remoteFilesystems, err = remote.FilesystemRequest(fsr); err != nil { + return + } + + type RemoteLocalMapping struct { + Remote zfs.DatasetPath + Local zfs.DatasetPath + LocalExists bool + } + replMapping := make(map[string]RemoteLocalMapping, len(remoteFilesystems)) + localTraversal := zfs.NewDatasetPathForest() + localExists, err := zfs.ZFSListFilesystemExists() + if err != nil { + log.Printf("cannot get local filesystems map: %s", err) + return err + } + + { + + log.Printf("mapping using %#v\n", pull.Mapping) + for fs := range remoteFilesystems { + var err error + var localFs zfs.DatasetPath + localFs, err = pull.Mapping.Map(remoteFilesystems[fs]) + if err != nil { + if err != zfs.NoMatchError { + log.Printf("error mapping %s: %#v\n", remoteFilesystems[fs], err) + return err + } + continue + } + m := RemoteLocalMapping{remoteFilesystems[fs], localFs, localExists(localFs)} + replMapping[m.Local.ToString()] = m + localTraversal.Add(m.Local) + } + + } + + log.Printf("remoteFilesystems: %#v\nreplMapping: %#v\n", remoteFilesystems, replMapping) + + // per fs sync, assume sorted in top-down order TODO + + localTraversal.WalkTopDown(func(v zfs.DatasetPathVisit) bool { + + if v.FilledIn { + if localExists(v.Path) { + return true + } + log.Printf("aborting, don't know how to create fill-in dataset %s", v.Path) + err = fmt.Errorf("aborting, don't know how to create fill-in dataset: %s", v.Path) + return false + } + + m, ok := replMapping[v.Path.ToString()] + if !ok { + panic("internal inconsistency: replMapping should contain mapping for any path that was not filled in by WalkTopDown()") + } + + log := func(format string, args ...interface{}) { + log.Printf("[%s => %s]: %s", m.Remote.ToString(), m.Local.ToString(), fmt.Sprintf(format, args...)) + } + + log("mapping: %#v\n", m) + + var versions []zfs.FilesystemVersion + if m.LocalExists { + if versions, err = zfs.ZFSListFilesystemVersions(m.Local); err != nil { + log("cannot get filesystem versions, stopping...: %v\n", m.Local.ToString(), m, err) + return false + } + } + + var theirVersions []zfs.FilesystemVersion + theirVersions, err = remote.FilesystemVersionsRequest(rpc.FilesystemVersionsRequest{ + Filesystem: m.Remote, + }) + if err != nil { + log("cannot fetch remote filesystem versions, stopping: %s", err) + return false + } + + diff := zfs.MakeFilesystemDiff(versions, theirVersions) + log("diff: %#v\n", diff) + + if diff.IncrementalPath == nil { + log("performing initial sync, following policy: %#v", pull.InitialReplPolicy) + + if pull.InitialReplPolicy != rpc.InitialReplPolicyMostRecent { + panic(fmt.Sprintf("policy %#v not implemented", pull.InitialReplPolicy)) + } + + snapsOnly := make([]zfs.FilesystemVersion, 0, len(diff.MRCAPathRight)) + for s := range diff.MRCAPathRight { + if diff.MRCAPathRight[s].Type == zfs.Snapshot { + snapsOnly = append(snapsOnly, diff.MRCAPathRight[s]) + } + } + + if len(snapsOnly) < 1 { + log("cannot perform initial sync: no remote snapshots. stopping...") + return false + } + + r := rpc.InitialTransferRequest{ + Filesystem: m.Remote, + FilesystemVersion: snapsOnly[len(snapsOnly)-1], + } + + log("requesting initial transfer") + + var stream io.Reader + if stream, err = remote.InitialTransferRequest(r); err != nil { + log("error initial transfer request, stopping...: %s", err) + return false + } + + log("received initial transfer request response. zfs recv...") + + if err = zfs.ZFSRecv(m.Local, stream, "-u"); err != nil { + log("error receiving stream, stopping...: %s", err) + return false + } + + log("configuring properties of received filesystem") + + if err = zfs.ZFSSet(m.Local, "readonly", "on"); err != nil { + + } + + log("finished initial transfer") + + } else if len(diff.IncrementalPath) < 2 { + log("remote and local are in sync") + } else { + + log("incremental transfers using path: %#v", diff.IncrementalPath) + + for i := 0; i < len(diff.IncrementalPath)-1; i++ { + + from, to := diff.IncrementalPath[i], diff.IncrementalPath[i+1] + + log := func(format string, args ...interface{}) { + log("[%s => %s]: %s", from.Name, to.Name, fmt.Sprintf(format, args...)) + } + + r := rpc.IncrementalTransferRequest{ + Filesystem: m.Remote, + From: from, + To: to, + } + log("requesting incremental transfer: %#v", r) + + var stream io.Reader + if stream, err = remote.IncrementalTransferRequest(r); err != nil { + log("error requesting incremental transfer, stopping...: %s", err.Error()) + return false + } + + log("receving incremental transfer") + + if err = zfs.ZFSRecv(m.Local, stream); err != nil { + log("error receiving stream, stopping...: %s", err) + return false + } + + log("finished incremental transfer") + + } + + log("finished incremental transfer path") + + } + + return true + + }) + + return + +} diff --git a/rpc/structs.go b/rpc/structs.go index 96d73e3..0715776 100644 --- a/rpc/structs.go +++ b/rpc/structs.go @@ -59,6 +59,17 @@ type ByteStreamRPCProtocolVersionRequest struct { ClientVersion uint8 } +const LOCAL_TRANSPORT_IDENTITY string = "local" + +const DEFAULT_INITIAL_REPL_POLICY = InitialReplPolicyMostRecent + +type InitialReplPolicy string + +const ( + InitialReplPolicyMostRecent InitialReplPolicy = "most_recent" + InitialReplPolicyAll InitialReplPolicy = "all" +) + type CloseRequest struct { Goodbye string }