Wireframe main executable.

This commit is contained in:
Christian Schwarz 2017-04-26 20:21:18 +02:00
parent 00231ecb73
commit 9750bf3123
3 changed files with 106 additions and 55 deletions

11
cmd/config_test.go Normal file
View File

@ -0,0 +1,11 @@
package main
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestSampleConfigFileIsParsedWithoutErrors(t *testing.T) {
_, err := ParseConfig("./sampleconf/zrepl.yml")
assert.Nil(t, err)
}

34
cmd/handler.go Normal file
View File

@ -0,0 +1,34 @@
package main
import (
"io"
"github.com/zrepl/zrepl/zfs"
"github.com/zrepl/zrepl/model"
"github.com/zrepl/zrepl/rpc"
)
type Handler struct {}
func (h Handler) HandleFilesystemRequest(r rpc.FilesystemRequest) (roots []model.Filesystem, err error) {
roots = make([]model.Filesystem, 0, 10)
for _, root := range r.Roots {
var zfsRoot model.Filesystem
if zfsRoot, err = zfs.FilesystemsAtRoot(root); err != nil {
return
}
roots = append(roots, zfsRoot)
}
return
}
func (h Handler) HandleInitialTransferRequest(r rpc.InitialTransferRequest) (io.Reader, error) {
// TODO ACL
return zfs.InitialSend(r.Snapshot)
}
func (h Handler) HandleIncrementalTransferRequest(r rpc.IncrementalTransferRequest) (io.Reader, error) {
// TODO ACL
return zfs.IncrementalSend(r.FromSnapshot, r.ToSnapshot)
}

View File

@ -1,5 +1,14 @@
package main
import (
"github.com/urfave/cli"
"errors"
"fmt"
"io"
"github.com/zrepl/zrepl/sshbytestream"
"github.com/zrepl/zrepl/rpc"
)
type Role uint
const (
@ -7,67 +16,64 @@ const (
ROLE_ACTION Role = iota
)
var conf Config
var handler Handler
func main() {
role = ROLE_IPC // TODO: if argv[1] == ipc then...
switch (role) {
case ROLE_IPC:
doIPC()
case ROLE_ACTION:
doAction()
app := cli.NewApp()
app.Name = "zrepl"
app.Usage = "replicate zfs datasets"
app.EnableBashCompletion = true
app.Flags = []cli.Flag{
cli.StringFlag{Name: "config"},
}
app.Before = func (c *cli.Context) (err error) {
if !c.GlobalIsSet("config") {
return errors.New("config flag not set")
}
if conf, err = ParseConfig(c.GlobalString("config")); err != nil {
return
}
func doIPC() {
sshByteStream = sshbytestream.Incoming()
handler = Handler{}
if err := ListenByteStreamRPC(sshByteStream, handler); err != nil {
// PANIC
return
}
app.Commands = []cli.Command{
{
Name: "sink",
Aliases: []string{"s"},
Usage: "start in sink mode",
Flags: []cli.Flag{
cli.StringFlag{Name: "identity"},
},
Action: doSink,
},
{
Name: "run",
Aliases: []string{"r"},
Usage: "do replication",
Action: doRun,
},
}
// exit(0)
app.RunAndExitOnError()
}
func doAction() {
sshByteStream = sshbytestream.Outgoing(model.SSHTransport{})
remote,_ := ConnectByteStreamRPC(sshByteStream)
request := NewFilesystemRequest(["zroot/var/db", "zroot/home"])
forest, _ := remote.FilesystemRequest(request)
for tree := forest {
fmt.Println(tree)
}
}
type Handler struct {}
func (h Handler) HandleFilesystemRequest(r FilesystemRequest) (roots []model.Filesystem, err error) {
roots = make([]model.Filesystem, 0, 10)
for _, root := range r.Roots {
if zfsRoot, err := zfs.FilesystemsAtRoot(root); err != nil {
return nil, err
}
roots = append(roots, zfsRoot)
}
func doSink(c *cli.Context) (err error) {
var sshByteStream io.ReadWriteCloser
if sshByteStream, err = sshbytestream.Incoming(); err != nil {
return
}
func (h Handler) HandleInitialTransferRequest(r InitialTransferRequest) (io.Read, error) {
// TODO ACL
return zfs.InitialSend(r.Snapshot)
return rpc.ListenByteStreamRPC(sshByteStream, handler)
}
func (h Handler) HandleIncrementalTransferRequestRequest(r IncrementalTransferRequest) (io.Read, error) {
// TODO ACL
return zfs.IncrementalSend(r.FromSnapshot, r.ToSnapshot)
func doRun(c *cli.Context) error {
fmt.Printf("%#v", conf)
return nil
}