Compare commits

..

5 Commits

Author SHA1 Message Date
Christian Schwarz 0eb7032735 daemon/control: envconst timeout for control socket server-side timeouts
refs #262
2020-02-17 22:36:32 +01:00
Christian Schwarz 3ff1966cab docs/installation: use && for early exit if build-in-docker step fails 2020-02-17 18:02:04 +01:00
Christian Schwarz a3842155c5 zrepl test filesystems: support snap job type 2020-02-17 18:02:04 +01:00
Christian Schwarz 02b3b4f80c fix some typos 2020-02-17 18:02:04 +01:00
Christian Schwarz 0882290595 README update
- donation links
- package overview
2020-02-17 18:02:04 +01:00
24 changed files with 67 additions and 318 deletions
+21 -9
View File
@@ -1,8 +1,9 @@
[![GitHub license](https://img.shields.io/github/license/zrepl/zrepl.svg)](https://github.com/zrepl/zrepl/blob/master/LICENSE) [![GitHub license](https://img.shields.io/github/license/zrepl/zrepl.svg)](https://github.com/zrepl/zrepl/blob/master/LICENSE)
[![Language: Go](https://img.shields.io/badge/language-Go-6ad7e5.svg)](https://golang.org/) [![Language: Go](https://img.shields.io/badge/language-Go-6ad7e5.svg)](https://golang.org/)
[![User Docs](https://img.shields.io/badge/docs-web-blue.svg)](https://zrepl.github.io) [![User Docs](https://img.shields.io/badge/docs-web-blue.svg)](https://zrepl.github.io)
[![Donate via PayPal](https://img.shields.io/badge/donate-paypal-yellow.svg)](https://www.paypal.com/cgi-bin/webscr?cmd=_s-xclick&hosted_button_id=R5QSXJVYHGX96) [![Donate via Patreon](https://img.shields.io/endpoint.svg?url=https%3A%2F%2Fshieldsio-patreon.herokuapp.com%2Fzrepl%2Fpledges&style=flat&color=yellow)](https://www.patreon.com/zrepl)
[![Donate via Liberapay](https://img.shields.io/liberapay/receives/zrepl.svg?logo=liberapay)](https://liberapay.com/zrepl/donate) [![Donate via Liberapay](https://img.shields.io/liberapay/receives/zrepl.svg?logo=liberapay)](https://liberapay.com/zrepl/donate)
[![Donate via PayPal](https://img.shields.io/badge/donate-paypal-yellow.svg)](https://www.paypal.com/cgi-bin/webscr?cmd=_s-xclick&hosted_button_id=R5QSXJVYHGX96)
[![Twitter](https://img.shields.io/twitter/url/https/github.com/zrepl/zrepl.svg?style=social)](https://twitter.com/intent/tweet?text=Wow:&url=https%3A%2F%2Fgithub.com%2Fzrepl%2Fzrepl) [![Twitter](https://img.shields.io/twitter/url/https/github.com/zrepl/zrepl.svg?style=social)](https://twitter.com/intent/tweet?text=Wow:&url=https%3A%2F%2Fgithub.com%2Fzrepl%2Fzrepl)
# zrepl # zrepl
@@ -23,6 +24,7 @@ zrepl is a one-stop ZFS backup & replication solution.
If so, think of an expressive configuration example. If so, think of an expressive configuration example.
2. Think of at least one use case that generalizes from your concrete application. 2. Think of at least one use case that generalizes from your concrete application.
3. Open an issue on GitHub with example conf & use case attached. 3. Open an issue on GitHub with example conf & use case attached.
4. **Optional**: [Post a bounty](https://www.bountysource.com/teams/zrepl) on the issue, or [contact Christian Schwarz](https://cschwarz.com) for contract work.
The above does not apply if you already implemented everything. The above does not apply if you already implemented everything.
Check out the *Coding Workflow* section below for details. Check out the *Coding Workflow* section below for details.
@@ -46,11 +48,8 @@ zrepl is written in [Go](https://golang.org) and uses [Go modules](https://githu
The documentation is written in [ReStructured Text](http://docutils.sourceforge.net/rst.html) using the [Sphinx](https://www.sphinx-doc.org) framework. The documentation is written in [ReStructured Text](http://docutils.sourceforge.net/rst.html) using the [Sphinx](https://www.sphinx-doc.org) framework.
To get started, run `./lazy.sh devsetup` to easily install build dependencies and read `docs/installation.rst -> Compiling from Source`. To get started, run `./lazy.sh devsetup` to easily install build dependencies and read `docs/installation.rst -> Compiling from Source`.
`lazy.sh` uses `python3-pip` to fetch the build dependencies for the docs - you might want to use a [venv](https://docs.python.org/3/library/venv.html).
### Overall Architecture If you just want to install the Go dependencies, run `./lazy.sh godep`.
The application architecture is documented as part of the user docs in the *Implementation* section (`docs/content/impl`).
Make sure to develop an understanding how zrepl is typically used by studying the user docs first.
### Project Structure ### Project Structure
@@ -62,14 +61,15 @@ Make sure to develop an understanding how zrepl is typically used by studying th
│   └── samples │   └── samples
├── daemon # the implementation of `zrepl daemon` subcommand ├── daemon # the implementation of `zrepl daemon` subcommand
│   ├── filters │   ├── filters
│   ├── hooks # snapshot hooks
│   ├── job # job implementations │   ├── job # job implementations
│   ├── logging # logging outlets + formatters │   ├── logging # logging outlets + formatters
│   ├── nethelpers │   ├── nethelpers
│   ├── prometheus │   ├── prometheus
│   ├── pruner # pruner implementation │   ├── pruner # pruner implementation
│   ├── snapper # snapshotter implementation │   ├── snapper # snapshotter implementation
├── docs # sphinx-based documentation
├── dist # supplemental material for users & package maintainers ├── dist # supplemental material for users & package maintainers
├── docs # sphinx-based documentation
│   ├── **/*.rst # documentation in reStructuredText │   ├── **/*.rst # documentation in reStructuredText
│   ├── sphinxconf │   ├── sphinxconf
│   │   └── conf.py # sphinx config (see commit 445a280 why its not in docs/) │   │   └── conf.py # sphinx config (see commit 445a280 why its not in docs/)
@@ -78,6 +78,7 @@ Make sure to develop an understanding how zrepl is typically used by studying th
│   └── public_git # checkout of zrepl.github.io managed by above shell script │   └── public_git # checkout of zrepl.github.io managed by above shell script
├── endpoint # implementation of replication endpoints (=> package replication) ├── endpoint # implementation of replication endpoints (=> package replication)
├── logger # our own logger package ├── logger # our own logger package
├── platformtest # test suite for our zfs abstractions (error classification, etc)
├── pruning # pruning rules (the logic, not the actual execution) ├── pruning # pruning rules (the logic, not the actual execution)
│   └── retentiongrid │   └── retentiongrid
├── replication ├── replication
@@ -92,14 +93,13 @@ Make sure to develop an understanding how zrepl is typically used by studying th
│ ├── transportmux # TCP connecter and listener used to split control & data traffic │ ├── transportmux # TCP connecter and listener used to split control & data traffic
│ └── versionhandshake # replication protocol version handshake perfomed on newly established connections │ └── versionhandshake # replication protocol version handshake perfomed on newly established connections
├── tlsconf # abstraction for Go TLS server + client config ├── tlsconf # abstraction for Go TLS server + client config
├── transport # transports implementation ├── transport # transport implementations
│ ├── fromconfig │ ├── fromconfig
│ ├── local │ ├── local
│ ├── ssh │ ├── ssh
│ ├── tcp │ ├── tcp
│ └── tls │ └── tls
├── util ├── util
├── vendor # managed by dep
├── version # abstraction for versions (filled during build by Makefile) ├── version # abstraction for versions (filled during build by Makefile)
└── zfs # zfs(8) wrappers └── zfs # zfs(8) wrappers
``` ```
@@ -134,3 +134,15 @@ There will not be a big refactoring (an attempt was made, but it's destroying to
However, new contributions & patches should fix naming without further notice in the commit message. However, new contributions & patches should fix naming without further notice in the commit message.
### RPC debugging
Optionally, there are various RPC-related environment varibles, that if set to something != `""` will produce additional debug output on stderr:
https://github.com/zrepl/zrepl/blob/master/rpc/rpc_debug.go#L11
https://github.com/zrepl/zrepl/blob/master/rpc/dataconn/dataconn_debug.go#L11
https://github.com/zrepl/zrepl/blob/master/rpc/dataconn/stream/stream_debug.go#L11
https://github.com/zrepl/zrepl/blob/master/rpc/dataconn/heartbeatconn/heartbeatconn_debug.go#L11
+2
View File
@@ -60,6 +60,8 @@ func runTestFilterCmd(subcommand *cli.Subcommand, args []string) error {
confFilter = j.Filesystems confFilter = j.Filesystems
case *config.PushJob: case *config.PushJob:
confFilter = j.Filesystems confFilter = j.Filesystems
case *config.SnapJob:
confFilter = j.Filesystems
default: default:
return fmt.Errorf("job type %T does not have filesystems filter", j) return fmt.Errorf("job type %T does not have filesystems filter", j)
} }
+2 -2
View File
@@ -147,8 +147,8 @@ func (j *controlJob) Run(ctx context.Context) {
server := http.Server{ server := http.Server{
Handler: mux, Handler: mux,
// control socket is local, 1s timeout should be more than sufficient, even on a loaded system // control socket is local, 1s timeout should be more than sufficient, even on a loaded system
WriteTimeout: 1 * time.Second, WriteTimeout: envconst.Duration("ZREPL_DAEMON_CONTROL_WRITE_TIMEOUT", 1*time.Second),
ReadTimeout: 1 * time.Second, ReadTimeout: envconst.Duration("ZREPL_DAEMON_CONTROL_READ_TIMEOUT", 1*time.Second),
} }
outer: outer:
+2 -31
View File
@@ -4,7 +4,6 @@ import (
"context" "context"
"fmt" "fmt"
"github.com/golang/protobuf/proto"
"github.com/pkg/errors" "github.com/pkg/errors"
"github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus"
@@ -181,39 +180,11 @@ func (j *PassiveSide) Run(ctx context.Context) {
} }
ctxInterceptor := func(handlerCtx context.Context) context.Context { ctxInterceptor := func(handlerCtx context.Context) context.Context {
reqLog := log.WithField(logging.ReqIDField, rpc.MustGetRequestID(handlerCtx)) return logging.WithSubsystemLoggers(handlerCtx, log)
return logging.WithSubsystemLoggers(handlerCtx, reqLog)
}
// ctxInterceptor already injected request loggers (with request id)
pre := func(ctx context.Context, method string, req interface{}) {
l := endpoint.GetLogger(ctx).WithField("method", method).WithField("reqT", fmt.Sprintf("%T", req))
if p, ok := req.(proto.Message); ok {
l = l.WithField("req", proto.CompactTextString(p))
}
ci, ok := ctx.Value(endpoint.ClientIdentityKey).(string)
if !ok {
// like endpoint, we must assume that rpc.Server sets it
panic(method)
}
l = l.WithField("clientIdentity", ci)
l.Debug("incoming")
}
post := func(ctx context.Context, res interface{}, err error) {
l := endpoint.GetLogger(ctx).
WithField("responseT", fmt.Sprintf("%T", res)).
WithField("errT", fmt.Sprintf("%T", err))
if s, ok := res.(fmt.Stringer); ok {
l = l.WithField("response", s.String())
}
if err != nil {
l = l.WithError(err)
}
l.Debug("handler done")
} }
rpcLoggers := rpc.GetLoggersOrPanic(ctx) // WithSubsystemLoggers above rpcLoggers := rpc.GetLoggersOrPanic(ctx) // WithSubsystemLoggers above
server := rpc.NewServer(handler, rpcLoggers, ctxInterceptor, pre, post) server := rpc.NewServer(handler, rpcLoggers, ctxInterceptor)
listener, err := j.listen() listener, err := j.listen()
if err != nil { if err != nil {
-5
View File
@@ -96,11 +96,6 @@ func WithSubsystemLoggers(ctx context.Context, log logger.Logger) context.Contex
Control: log.WithField(SubsysField, SubsysRPCControl), Control: log.WithField(SubsysField, SubsysRPCControl),
Data: log.WithField(SubsysField, SubsysRPCData), Data: log.WithField(SubsysField, SubsysRPCData),
}, },
// rpc.Loggers{
// General: logger.NewNullLogger(),
// Control: logger.NewNullLogger(),
// Data: logger.NewNullLogger(),
// },
) )
return ctx return ctx
} }
+1 -2
View File
@@ -22,7 +22,6 @@ const (
const ( const (
JobField string = "job" JobField string = "job"
SubsysField string = "subsystem" SubsysField string = "subsystem"
ReqIDField string = "reqid"
) )
type MetadataFlags int64 type MetadataFlags int64
@@ -86,7 +85,7 @@ func (f *HumanFormatter) Format(e *logger.Entry) (out []byte, err error) {
fmt.Fprintf(&line, "[%s]", col.Sprint(e.Level.Short())) fmt.Fprintf(&line, "[%s]", col.Sprint(e.Level.Short()))
} }
prefixFields := []string{JobField, SubsysField, ReqIDField} prefixFields := []string{JobField, SubsysField}
prefixed := make(map[string]bool, len(prefixFields)+2) prefixed := make(map[string]bool, len(prefixFields)+2)
for _, field := range prefixFields { for _, field := range prefixFields {
val, ok := e.Fields[field] val, ok := e.Fields[field]
+3 -3
View File
@@ -101,9 +101,9 @@ and serves as a reference for build dependencies and procedure:
:: ::
git clone https://github.com/zrepl/zrepl.git git clone https://github.com/zrepl/zrepl.git && \
cd zrepl cd zrepl && \
sudo docker build -t zrepl_build -f build.Dockerfile . sudo docker build -t zrepl_build -f build.Dockerfile . && \
sudo docker run -it --rm \ sudo docker run -it --rm \
-v "${PWD}:/src" \ -v "${PWD}:/src" \
--user "$(id -u):$(id -g)" \ --user "$(id -u):$(id -g)" \
-4
View File
@@ -19,10 +19,6 @@ func WithLogger(ctx context.Context, log Logger) context.Context {
return context.WithValue(ctx, contextKeyLogger, log) return context.WithValue(ctx, contextKeyLogger, log)
} }
func GetLogger(ctx context.Context) Logger {
return getLogger(ctx)
}
func getLogger(ctx context.Context) Logger { func getLogger(ctx context.Context) Logger {
if l, ok := ctx.Value(contextKeyLogger).(Logger); ok { if l, ok := ctx.Value(contextKeyLogger).(Logger); ok {
return l return l
+2 -4
View File
@@ -159,7 +159,6 @@ func sendArgsFromPDUAndValidateExists(ctx context.Context, fs string, fsv *pdu.F
} }
func (s *Sender) Send(ctx context.Context, r *pdu.SendReq) (*pdu.SendRes, zfs.StreamCopier, error) { func (s *Sender) Send(ctx context.Context, r *pdu.SendReq) (*pdu.SendRes, zfs.StreamCopier, error) {
getLogger(ctx).WithField("req", r.String()).Debug("incoming Send request")
_, err := s.filterCheckFS(r.Filesystem) _, err := s.filterCheckFS(r.Filesystem)
if err != nil { if err != nil {
@@ -260,7 +259,6 @@ func (s *Sender) Send(ctx context.Context, r *pdu.SendReq) (*pdu.SendRes, zfs.St
} }
func (p *Sender) SendCompleted(ctx context.Context, r *pdu.SendCompletedReq) (*pdu.SendCompletedRes, error) { func (p *Sender) SendCompleted(ctx context.Context, r *pdu.SendCompletedReq) (*pdu.SendCompletedRes, error) {
getLogger(ctx).WithField("req", r.String()).Debug("incoming SendCompleted request")
orig := r.GetOriginalReq() // may be nil, always use proto getters orig := r.GetOriginalReq() // may be nil, always use proto getters
fs := orig.GetFilesystem() fs := orig.GetFilesystem()
@@ -512,7 +510,7 @@ func (s *Receiver) ListFilesystems(ctx context.Context, req *pdu.ListFilesystemR
// present filesystem without the root_fs prefix // present filesystem without the root_fs prefix
fss := make([]*pdu.Filesystem, 0, len(filtered)) fss := make([]*pdu.Filesystem, 0, len(filtered))
for _, a := range filtered { for _, a := range filtered {
l := getLogger(ctx).WithField("fs", a.ToString()) l := getLogger(ctx).WithField("fs", a)
ph, err := zfs.ZFSGetFilesystemPlaceholderState(a) ph, err := zfs.ZFSGetFilesystemPlaceholderState(a)
if err != nil { if err != nil {
l.WithError(err).Error("error getting placeholder state") l.WithError(err).Error("error getting placeholder state")
@@ -599,7 +597,7 @@ func (s *Receiver) Send(ctx context.Context, req *pdu.SendReq) (*pdu.SendRes, zf
var maxConcurrentZFSRecvSemaphore = semaphore.New(envconst.Int64("ZREPL_ENDPOINT_MAX_CONCURRENT_RECV", 10)) var maxConcurrentZFSRecvSemaphore = semaphore.New(envconst.Int64("ZREPL_ENDPOINT_MAX_CONCURRENT_RECV", 10))
func (s *Receiver) Receive(ctx context.Context, req *pdu.ReceiveReq, receive zfs.StreamCopier) (*pdu.ReceiveRes, error) { func (s *Receiver) Receive(ctx context.Context, req *pdu.ReceiveReq, receive zfs.StreamCopier) (*pdu.ReceiveRes, error) {
getLogger(ctx).WithField("rr", req.String()).Debug("incoming Receive") getLogger(ctx).Debug("incoming Receive")
defer receive.Close() defer receive.Close()
root := s.clientRootFromCtx(ctx) root := s.clientRootFromCtx(ctx)
-1
View File
@@ -7,7 +7,6 @@ require (
github.com/gdamore/tcell v1.2.0 github.com/gdamore/tcell v1.2.0
github.com/go-logfmt/logfmt v0.4.0 github.com/go-logfmt/logfmt v0.4.0
github.com/go-sql-driver/mysql v1.4.1-0.20190907122137-b2c03bcae3d4 github.com/go-sql-driver/mysql v1.4.1-0.20190907122137-b2c03bcae3d4
github.com/gogo/protobuf v1.1.1
github.com/golang/protobuf v1.3.2 github.com/golang/protobuf v1.3.2
github.com/google/uuid v1.1.1 github.com/google/uuid v1.1.1
github.com/jinzhu/copier v0.0.0-20170922082739-db4671f3a9b8 github.com/jinzhu/copier v0.0.0-20170922082739-db4671f3a9b8
-1
View File
@@ -68,7 +68,6 @@ github.com/go-toolsmith/strparse v1.0.0/go.mod h1:YI2nUKP9YGZnL/L1/DLFBfixrcjslW
github.com/go-toolsmith/typep v0.0.0-20181030061450-d63dc7650676/go.mod h1:JSQCQMUPdRlMZFswiq3TGpNp1GMktqkR2Ns5AIQkATU= github.com/go-toolsmith/typep v0.0.0-20181030061450-d63dc7650676/go.mod h1:JSQCQMUPdRlMZFswiq3TGpNp1GMktqkR2Ns5AIQkATU=
github.com/go-toolsmith/typep v1.0.0/go.mod h1:JSQCQMUPdRlMZFswiq3TGpNp1GMktqkR2Ns5AIQkATU= github.com/go-toolsmith/typep v1.0.0/go.mod h1:JSQCQMUPdRlMZFswiq3TGpNp1GMktqkR2Ns5AIQkATU=
github.com/gobwas/glob v0.2.3/go.mod h1:d3Ez4x06l9bZtSvzIay5+Yzi0fmZzPgnTbPcKjJAkT8= github.com/gobwas/glob v0.2.3/go.mod h1:d3Ez4x06l9bZtSvzIay5+Yzi0fmZzPgnTbPcKjJAkT8=
github.com/gogo/protobuf v1.1.1 h1:72R+M5VuhED/KujmZVcIquuo8mBgX4oVda//DQb3PXo=
github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ= github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ=
github.com/gogo/protobuf v1.2.1/go.mod h1:hp+jE20tsWTFYpLwKvXlhS1hjn+gTNwPg2I6zVXpSg4= github.com/gogo/protobuf v1.2.1/go.mod h1:hp+jE20tsWTFYpLwKvXlhS1hjn+gTNwPg2I6zVXpSg4=
github.com/golang/glog v0.0.0-20160126235308-23def4e6c14b h1:VKtxabqXZkF25pY9ekfRL6a582T4P37/31XEstQ5p58= github.com/golang/glog v0.0.0-20160126235308-23def4e6c14b h1:VKtxabqXZkF25pY9ekfRL6a582T4P37/31XEstQ5p58=
+2 -7
View File
@@ -5,7 +5,6 @@ import (
"errors" "errors"
"fmt" "fmt"
"net" "net"
"os"
"sort" "sort"
"strings" "strings"
"sync" "sync"
@@ -15,7 +14,6 @@ import (
"google.golang.org/grpc/status" "google.golang.org/grpc/status"
"github.com/zrepl/zrepl/replication/report" "github.com/zrepl/zrepl/replication/report"
"github.com/zrepl/zrepl/tracing"
"github.com/zrepl/zrepl/util/chainlock" "github.com/zrepl/zrepl/util/chainlock"
"github.com/zrepl/zrepl/util/envconst" "github.com/zrepl/zrepl/util/envconst"
) )
@@ -208,7 +206,6 @@ func Do(ctx context.Context, planner Planner) (ReportFunc, WaitFunc) {
} }
run.attempts = append(run.attempts, cur) run.attempts = append(run.attempts, cur)
run.l.DropWhile(func() { run.l.DropWhile(func() {
ctx := tracing.Child(ctx, fmt.Sprintf("attempt#%d", ano)) // shadow
cur.do(ctx, prev) cur.do(ctx, prev)
}) })
prev = cur prev = cur
@@ -284,7 +281,7 @@ func Do(ctx context.Context, planner Planner) (ReportFunc, WaitFunc) {
} }
func (a *attempt) do(ctx context.Context, prev *attempt) { func (a *attempt) do(ctx context.Context, prev *attempt) {
pfss, err := a.planner.Plan(tracing.Child(ctx, "plan")) pfss, err := a.planner.Plan(ctx)
errTime := time.Now() errTime := time.Now()
defer a.l.Lock().Unlock() defer a.l.Lock().Unlock()
if err != nil { if err != nil {
@@ -364,7 +361,7 @@ func (a *attempt) do(ctx context.Context, prev *attempt) {
fssesDone.Add(1) fssesDone.Add(1)
go func(f *fs) { go func(f *fs) {
defer fssesDone.Done() defer fssesDone.Done()
f.do(tracing.Child(ctx, f.fs.ReportInfo().Name), stepQueue, prevs[f]) f.do(ctx, stepQueue, prevs[f])
}(f) }(f)
} }
a.l.DropWhile(func() { a.l.DropWhile(func() {
@@ -375,8 +372,6 @@ func (a *attempt) do(ctx context.Context, prev *attempt) {
func (fs *fs) do(ctx context.Context, pq *stepQueue, prev *fs) { func (fs *fs) do(ctx context.Context, pq *stepQueue, prev *fs) {
fmt.Fprintf(os.Stderr, "CHILD STACK: %v\n", tracing.GetStack(ctx))
defer fs.l.Lock().Unlock() defer fs.l.Lock().Unlock()
// get planned steps from replication logic // get planned steps from replication logic
+5 -7
View File
@@ -262,7 +262,6 @@ func (p *Planner) doPlanning(ctx context.Context) ([]*Filesystem, error) {
return nil, err return nil, err
} }
sfss := slfssres.GetFilesystems() sfss := slfssres.GetFilesystems()
// no progress here since we could run in a live-lock on connectivity issues
rlfssres, err := p.receiver.ListFilesystems(ctx, &pdu.ListFilesystemReq{}) rlfssres, err := p.receiver.ListFilesystems(ctx, &pdu.ListFilesystemReq{})
if err != nil { if err != nil {
@@ -606,7 +605,7 @@ func (s *Step) doReplication(ctx context.Context) error {
log := getLogger(ctx) log := getLogger(ctx)
sr := s.buildSendRequest(false) sr := s.buildSendRequest(false)
log.WithField("sr", sr.String()).Debug("initiate send request") log.Debug("initiate send request")
sres, sstreamCopier, err := s.sender.Send(ctx, sr) sres, sstreamCopier, err := s.sender.Send(ctx, sr)
if err != nil { if err != nil {
log.WithError(err).Error("send request failed") log.WithError(err).Error("send request failed")
@@ -633,7 +632,7 @@ func (s *Step) doReplication(ctx context.Context) error {
To: sr.GetTo(), To: sr.GetTo(),
ClearResumeToken: !sres.UsedResumeToken, ClearResumeToken: !sres.UsedResumeToken,
} }
log.WithField("rr", rr.String()).Debug("initiate receive request") log.Debug("initiate receive request")
_, err = s.receiver.Receive(ctx, rr, byteCountingStream) _, err = s.receiver.Receive(ctx, rr, byteCountingStream)
if err != nil { if err != nil {
log. log.
@@ -649,11 +648,10 @@ func (s *Step) doReplication(ctx context.Context) error {
} }
log.Debug("receive finished") log.Debug("receive finished")
scr := &pdu.SendCompletedReq{ log.Debug("tell sender replication completed")
_, err = s.sender.SendCompleted(ctx, &pdu.SendCompletedReq{
OriginalReq: sr, OriginalReq: sr,
} })
log.WithField("scr", scr.String()).Debug("tell sender replication completed")
_, err = s.sender.SendCompleted(ctx, scr)
if err != nil { if err != nil {
log.WithError(err).Error("error telling sender that replication completed successfully") log.WithError(err).Error("error telling sender that replication completed successfully")
return err return err
+1 -1
View File
@@ -75,7 +75,7 @@ func (c *Client) recv(ctx context.Context, conn *stream.Conn, res proto.Message)
if err != nil { if err != nil {
return err return err
} }
header := string(headerBuf) // FIXME header := string(headerBuf)
if strings.HasPrefix(header, responseHeaderHandlerErrorPrefix) { if strings.HasPrefix(header, responseHeaderHandlerErrorPrefix) {
// FIXME distinguishable error type // FIXME distinguishable error type
return &RemoteHandlerError{strings.TrimPrefix(header, responseHeaderHandlerErrorPrefix)} return &RemoteHandlerError{strings.TrimPrefix(header, responseHeaderHandlerErrorPrefix)}
+10 -23
View File
@@ -14,8 +14,8 @@ import (
"github.com/zrepl/zrepl/zfs" "github.com/zrepl/zrepl/zfs"
) )
// ReqInterceptor has a chance to exchange the context and connection on each request. // WireInterceptor has a chance to exchange the context and connection on each client connection.
type ReqInterceptor func(ctx context.Context, rawConn *transport.AuthConn) (context.Context, *transport.AuthConn) type WireInterceptor func(ctx context.Context, rawConn *transport.AuthConn) (context.Context, *transport.AuthConn)
// Handler implements the functionality that is exposed by Server to the Client. // Handler implements the functionality that is exposed by Server to the Client.
type Handler interface { type Handler interface {
@@ -30,27 +30,19 @@ type Handler interface {
PingDataconn(ctx context.Context, r *pdu.PingReq) (*pdu.PingRes, error) PingDataconn(ctx context.Context, r *pdu.PingReq) (*pdu.PingRes, error)
} }
type PreHandlerInspector = func(ctx context.Context, endpoint string, req interface{})
type PostHandlerInspector = func(ctx context.Context, response interface{}, err error)
type Logger = logger.Logger type Logger = logger.Logger
type Server struct { type Server struct {
h Handler h Handler
ri ReqInterceptor wi WireInterceptor
log Logger log Logger
pre PreHandlerInspector
post PostHandlerInspector
} }
func NewServer(ri ReqInterceptor, logger Logger, pre PreHandlerInspector, handler Handler, post PostHandlerInspector) *Server { func NewServer(wi WireInterceptor, logger Logger, handler Handler) *Server {
return &Server{ return &Server{
h: handler, h: handler,
ri: ri, wi: wi,
log: logger, log: logger,
pre: pre,
post: post,
} }
} }
@@ -91,8 +83,8 @@ func (s *Server) serveConn(nc *transport.AuthConn) {
defer s.log.Debug("serveConn done") defer s.log.Debug("serveConn done")
ctx := context.Background() ctx := context.Background()
if s.ri != nil { if s.wi != nil {
ctx, nc = s.ri(ctx, nc) ctx, nc = s.wi(ctx, nc)
} }
c := stream.Wrap(nc, HeartbeatInterval, HeartbeatPeerTimeout) c := stream.Wrap(nc, HeartbeatInterval, HeartbeatPeerTimeout)
@@ -108,7 +100,7 @@ func (s *Server) serveConn(nc *transport.AuthConn) {
s.log.WithError(err).Error("error reading structured part") s.log.WithError(err).Error("error reading structured part")
return return
} }
endpoint := string(header) // FIXME endpoint := string(header)
reqStructured, err := c.ReadStreamedMessage(ctx, RequestStructuredMaxSize, ReqStructured) reqStructured, err := c.ReadStreamedMessage(ctx, RequestStructuredMaxSize, ReqStructured)
if err != nil { if err != nil {
@@ -128,7 +120,6 @@ func (s *Server) serveConn(nc *transport.AuthConn) {
s.log.WithError(err).Error("cannot unmarshal send request") s.log.WithError(err).Error("cannot unmarshal send request")
return return
} }
s.pre(ctx, endpoint, &req)
res, sendStream, handlerErr = s.h.Send(ctx, &req) // SHADOWING res, sendStream, handlerErr = s.h.Send(ctx, &req) // SHADOWING
case EndpointRecv: case EndpointRecv:
var req pdu.ReceiveReq var req pdu.ReceiveReq
@@ -136,7 +127,6 @@ func (s *Server) serveConn(nc *transport.AuthConn) {
s.log.WithError(err).Error("cannot unmarshal receive request") s.log.WithError(err).Error("cannot unmarshal receive request")
return return
} }
s.pre(ctx, endpoint, &req)
res, handlerErr = s.h.Receive(ctx, &req, &streamCopier{streamConn: c, closeStreamOnClose: false}) // SHADOWING res, handlerErr = s.h.Receive(ctx, &req, &streamCopier{streamConn: c, closeStreamOnClose: false}) // SHADOWING
case EndpointPing: case EndpointPing:
var req pdu.PingReq var req pdu.PingReq
@@ -144,7 +134,6 @@ func (s *Server) serveConn(nc *transport.AuthConn) {
s.log.WithError(err).Error("cannot unmarshal ping request") s.log.WithError(err).Error("cannot unmarshal ping request")
return return
} }
s.pre(ctx, endpoint, &req)
res, handlerErr = s.h.PingDataconn(ctx, &req) // SHADOWING res, handlerErr = s.h.PingDataconn(ctx, &req) // SHADOWING
default: default:
s.log.WithField("endpoint", endpoint).Error("unknown endpoint") s.log.WithField("endpoint", endpoint).Error("unknown endpoint")
@@ -153,8 +142,6 @@ func (s *Server) serveConn(nc *transport.AuthConn) {
s.log.WithField("endpoint", endpoint).WithField("errType", fmt.Sprintf("%T", handlerErr)).Debug("handler returned") s.log.WithField("endpoint", endpoint).WithField("errType", fmt.Sprintf("%T", handlerErr)).Debug("handler returned")
s.post(ctx, res, handlerErr)
// prepare protobuf now to return the protobuf error in the header // prepare protobuf now to return the protobuf error in the header
// if marshaling fails. We consider failed marshaling a handler error // if marshaling fails. We consider failed marshaling a handler error
var protobuf *bytes.Buffer var protobuf *bytes.Buffer
+1 -1
View File
@@ -49,7 +49,7 @@ type Wire interface {
// No data that could otherwise be Read is lost as a consequence of this call. // No data that could otherwise be Read is lost as a consequence of this call.
// The use case for this API is abortive connection shutdown. // The use case for this API is abortive connection shutdown.
// To provide any value over draining Wire using io.Read, an implementation // To provide any value over draining Wire using io.Read, an implementation
// will likely use out-of-bounds messaging mechanisms. // will likely use out-of-band messaging mechanisms.
// TODO WaitForPeerClose() (supported bool, err error) // TODO WaitForPeerClose() (supported bool, err error)
} }
@@ -101,11 +101,7 @@ func (*transportCredentials) OverrideServerName(string) error {
type ContextInterceptor = func(ctx context.Context) context.Context type ContextInterceptor = func(ctx context.Context) context.Context
type PreHandlerInspector func(ctx context.Context, endpoint string, req interface{}) func NewInterceptors(logger Logger, clientIdentityKey interface{}, ctxInterceptor ContextInterceptor) (unary grpc.UnaryServerInterceptor, stream grpc.StreamServerInterceptor) {
type PostHandlerInspector func(ctx context.Context, response interface{}, err error)
func NewInterceptors(logger Logger, clientIdentityKey interface{}, ctxInterceptor ContextInterceptor, pre PreHandlerInspector, post PostHandlerInspector) (unary grpc.UnaryServerInterceptor, stream grpc.StreamServerInterceptor) {
unary = func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (resp interface{}, err error) { unary = func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (resp interface{}, err error) {
logger.WithField("fullMethod", info.FullMethod).Debug("request") logger.WithField("fullMethod", info.FullMethod).Debug("request")
p, ok := peer.FromContext(ctx) p, ok := peer.FromContext(ctx)
@@ -122,10 +118,7 @@ func NewInterceptors(logger Logger, clientIdentityKey interface{}, ctxIntercepto
if ctxInterceptor != nil { if ctxInterceptor != nil {
ctx = ctxInterceptor(ctx) ctx = ctxInterceptor(ctx)
} }
pre(ctx, info.FullMethod, req) return handler(ctx, req)
res, err := handler(ctx, req)
post(ctx, res, err)
return res, err
} }
stream = func(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error { stream = func(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
panic("unimplemented") panic("unimplemented")
@@ -4,18 +4,14 @@ package grpchelper
import ( import (
"context" "context"
"fmt"
"strings"
"time" "time"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/keepalive" "google.golang.org/grpc/keepalive"
"google.golang.org/grpc/metadata"
"github.com/zrepl/zrepl/logger" "github.com/zrepl/zrepl/logger"
"github.com/zrepl/zrepl/rpc/grpcclientidentity" "github.com/zrepl/zrepl/rpc/grpcclientidentity"
"github.com/zrepl/zrepl/rpc/netadaptor" "github.com/zrepl/zrepl/rpc/netadaptor"
"github.com/zrepl/zrepl/tracing"
"github.com/zrepl/zrepl/transport" "github.com/zrepl/zrepl/transport"
) )
@@ -32,40 +28,15 @@ type Logger = logger.Logger
// ClientConn is an easy-to-use wrapper around the Dialer and TransportCredentials interface // ClientConn is an easy-to-use wrapper around the Dialer and TransportCredentials interface
// to produce a grpc.ClientConn // to produce a grpc.ClientConn
func ClientConn(cn transport.Connecter, log Logger, genReqID func() string) *grpc.ClientConn { func ClientConn(cn transport.Connecter, log Logger) *grpc.ClientConn {
ka := grpc.WithKeepaliveParams(keepalive.ClientParameters{ ka := grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: StartKeepalivesAfterInactivityDuration, Time: StartKeepalivesAfterInactivityDuration,
Timeout: KeepalivePeerTimeout, Timeout: KeepalivePeerTimeout,
PermitWithoutStream: true, PermitWithoutStream: true,
}) })
unaryIntcpt := grpc.WithUnaryInterceptor(func(ctx context.Context, method string, req, reply interface{}, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
traceStack := tracing.GetStack(ctx)
if len(traceStack) == 0 {
panic("implementation error: expecting trace stack")
}
reqId := genReqID() + " " + strings.Join(traceStack, " <> ")
l := log
l = l.WithField("reqID", reqId)
l = l.WithField("reqT", fmt.Sprintf("%T", req))
l = l.WithField("repT", fmt.Sprintf("%T", reply))
l = l.WithField("ctxT", fmt.Sprintf("%T", ctx))
l.Debug("unary intercepted")
ctx = metadata.AppendToOutgoingContext(ctx, "zrepl-req-id", reqId)
err := invoker(ctx, method, req, reply, cc, opts...)
l = log.WithField("repT", fmt.Sprintf("%T", reply))
l = log.WithField("errT", fmt.Sprintf("%T", err))
l.Debug("reply intercepted")
return err
})
dialerOption := grpc.WithDialer(grpcclientidentity.NewDialer(log, cn)) dialerOption := grpc.WithDialer(grpcclientidentity.NewDialer(log, cn))
cred := grpc.WithTransportCredentials(grpcclientidentity.NewTransportCredentials(log)) cred := grpc.WithTransportCredentials(grpcclientidentity.NewTransportCredentials(log))
cc, err := grpc.DialContext(context.Background(), "doesn't matter done by dialer", dialerOption, cred, ka, unaryIntcpt) cc, err := grpc.DialContext(context.Background(), "doesn't matter done by dialer", dialerOption, cred, ka)
if err != nil { if err != nil {
log.WithError(err).Error("cannot create gRPC client conn (non-blocking)") log.WithError(err).Error("cannot create gRPC client conn (non-blocking)")
// It's ok to panic here: the we call grpc.DialContext without the // It's ok to panic here: the we call grpc.DialContext without the
@@ -79,7 +50,7 @@ func ClientConn(cn transport.Connecter, log Logger, genReqID func() string) *grp
} }
// NewServer is a convenience interface around the TransportCredentials and Interceptors interface. // NewServer is a convenience interface around the TransportCredentials and Interceptors interface.
func NewServer(authListener transport.AuthenticatedListener, clientIdentityKey interface{}, logger grpcclientidentity.Logger, interceptor grpcclientidentity.ContextInterceptor, pre grpcclientidentity.PreHandlerInspector, post grpcclientidentity.PostHandlerInspector) (srv *grpc.Server, serve func() error) { func NewServer(authListener transport.AuthenticatedListener, clientIdentityKey interface{}, logger grpcclientidentity.Logger, ctxInterceptor grpcclientidentity.ContextInterceptor) (srv *grpc.Server, serve func() error) {
ka := grpc.KeepaliveParams(keepalive.ServerParameters{ ka := grpc.KeepaliveParams(keepalive.ServerParameters{
Time: StartKeepalivesAfterInactivityDuration, Time: StartKeepalivesAfterInactivityDuration,
Timeout: KeepalivePeerTimeout, Timeout: KeepalivePeerTimeout,
@@ -89,17 +60,8 @@ func NewServer(authListener transport.AuthenticatedListener, clientIdentityKey i
PermitWithoutStream: true, PermitWithoutStream: true,
}) })
tcs := grpcclientidentity.NewTransportCredentials(logger) tcs := grpcclientidentity.NewTransportCredentials(logger)
unary, stream := grpcclientidentity.NewInterceptors(logger, clientIdentityKey, interceptor, pre, post) unary, stream := grpcclientidentity.NewInterceptors(logger, clientIdentityKey, ctxInterceptor)
unary2 := func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (resp interface{}, err error) { srv = grpc.NewServer(grpc.Creds(tcs), grpc.UnaryInterceptor(unary), grpc.StreamInterceptor(stream), ka, ep)
md, ok := metadata.FromIncomingContext(ctx)
if !ok {
logger.Error("no request id in incoming request")
} else {
logger.Info(fmt.Sprintf("req-id = %v", md.Get("zrepl-req-id")))
}
return unary(ctx, req, info, handler)
}
srv = grpc.NewServer(grpc.Creds(tcs), grpc.UnaryInterceptor(unary2), grpc.StreamInterceptor(stream), ka, ep)
serve = func() error { serve = func() error {
if err := srv.Serve(netadaptor.New(authListener, logger)); err != nil { if err := srv.Serve(netadaptor.New(authListener, logger)); err != nil {
+3 -7
View File
@@ -26,7 +26,6 @@ import (
// Client implements the active side of a replication setup. // Client implements the active side of a replication setup.
// It satisfies the Endpoint, Sender and Receiver interface defined by package replication. // It satisfies the Endpoint, Sender and Receiver interface defined by package replication.
type Client struct { type Client struct {
reqIDGen *requestIDGenerator
dataClient *dataconn.Client dataClient *dataconn.Client
controlClient pdu.ReplicationClient // this the grpc client instance, see constructor controlClient pdu.ReplicationClient // this the grpc client instance, see constructor
controlConn *grpc.ClientConn controlConn *grpc.ClientConn
@@ -48,13 +47,10 @@ func NewClient(cn transport.Connecter, loggers Loggers) *Client {
muxedConnecter := mux(cn) muxedConnecter := mux(cn)
c := &Client{ c := &Client{
loggers: loggers, loggers: loggers,
closed: make(chan struct{}), closed: make(chan struct{}),
reqIDGen: newRequestIDGenerator(),
} }
grpcConn := grpchelper.ClientConn(muxedConnecter.control, loggers.Control, func() string { grpcConn := grpchelper.ClientConn(muxedConnecter.control, loggers.Control)
return c.reqIDGen.newID().String()
})
go func() { go func() {
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
-58
View File
@@ -1,58 +0,0 @@
package rpc
import (
"context"
"encoding/base64"
"strings"
"github.com/google/uuid"
)
type requestIDGenerator struct{}
func newRequestIDGenerator() *requestIDGenerator {
return &requestIDGenerator{}
}
type RequestID struct{ s string }
func (r RequestID) String() string { return r.s }
type requestIDContextKey int
const (
requestIDContextKeyRequestID requestIDContextKey = 1 + iota
)
func (g *requestIDGenerator) newID() RequestID {
id := uuid.New()
var buf strings.Builder
enc := base64.NewEncoder(base64.RawStdEncoding, &buf)
n, err := enc.Write(id[:])
if err != nil {
panic(err)
} else if n != len(id) {
panic(n)
}
if err := enc.Close(); err != nil {
panic(err)
}
return RequestID{buf.String()}
}
func (g *requestIDGenerator) inject(ctx context.Context) context.Context {
return context.WithValue(ctx, requestIDContextKeyRequestID, g.newID())
}
func GetRequestID(ctx context.Context) (id RequestID, ok bool) {
id, ok = ctx.Value(requestIDContextKeyRequestID).(RequestID)
return id, ok
}
func MustGetRequestID(ctx context.Context) (id RequestID) {
id, ok := GetRequestID(ctx)
if !ok {
panic("calling context expectes request id to bet set")
}
return id
}
+4 -25
View File
@@ -7,7 +7,6 @@ import (
"github.com/zrepl/zrepl/endpoint" "github.com/zrepl/zrepl/endpoint"
"github.com/zrepl/zrepl/replication/logic/pdu" "github.com/zrepl/zrepl/replication/logic/pdu"
"github.com/zrepl/zrepl/rpc/dataconn" "github.com/zrepl/zrepl/rpc/dataconn"
"github.com/zrepl/zrepl/rpc/grpcclientidentity"
"github.com/zrepl/zrepl/rpc/grpcclientidentity/grpchelper" "github.com/zrepl/zrepl/rpc/grpcclientidentity/grpchelper"
"github.com/zrepl/zrepl/rpc/versionhandshake" "github.com/zrepl/zrepl/rpc/versionhandshake"
"github.com/zrepl/zrepl/transport" "github.com/zrepl/zrepl/transport"
@@ -31,35 +30,15 @@ type Server struct {
dataServerServe serveFunc dataServerServe serveFunc
} }
type CtxInterceptor = grpcclientidentity.ContextInterceptor type HandlerContextInterceptor func(ctx context.Context) context.Context
func chainInterceptors(interceptors []CtxInterceptor) CtxInterceptor {
return func(ctx context.Context) (_ context.Context) {
for _, i := range interceptors {
ctx = i(ctx)
}
return ctx
}
}
type PreHandlerInspector = func(ctx context.Context, endpoint string, req interface{})
type PostHandlerInspector = func(ctx context.Context, response interface{}, err error)
// config must be valid (use its Validate function). // config must be valid (use its Validate function).
func NewServer(handler Handler, loggers Loggers, ctxInterceptor CtxInterceptor, pre PreHandlerInspector, post PostHandlerInspector) *Server { func NewServer(handler Handler, loggers Loggers, ctxInterceptor HandlerContextInterceptor) *Server {
reqIdGen := newRequestIDGenerator()
ctxInterceptor = chainInterceptors([]CtxInterceptor{
reqIdGen.inject,
ctxInterceptor,
})
// setup control server // setup control server
controlServerServe := func(ctx context.Context, controlListener transport.AuthenticatedListener, errOut chan<- error) { controlServerServe := func(ctx context.Context, controlListener transport.AuthenticatedListener, errOut chan<- error) {
controlServer, serve := grpchelper.NewServer(controlListener, endpoint.ClientIdentityKey, loggers.Control, ctxInterceptor, pre, post) controlServer, serve := grpchelper.NewServer(controlListener, endpoint.ClientIdentityKey, loggers.Control, ctxInterceptor)
pdu.RegisterReplicationServer(controlServer, handler) pdu.RegisterReplicationServer(controlServer, handler)
// give time for graceful stop until deadline expires, then hard stop // give time for graceful stop until deadline expires, then hard stop
@@ -84,7 +63,7 @@ func NewServer(handler Handler, loggers Loggers, ctxInterceptor CtxInterceptor,
} }
return ctx, wire return ctx, wire
} }
dataServer := dataconn.NewServer(dataServerClientIdentitySetter, loggers.Data, pre, handler, post) dataServer := dataconn.NewServer(dataServerClientIdentitySetter, loggers.Data, handler)
dataServerServe := func(ctx context.Context, dataListener transport.AuthenticatedListener, errOut chan<- error) { dataServerServe := func(ctx context.Context, dataListener transport.AuthenticatedListener, errOut chan<- error) {
dataServer.Serve(ctx, dataListener) dataServer.Serve(ctx, dataListener)
errOut <- nil // TODO bad design of dataServer? errOut <- nil // TODO bad design of dataServer?
-53
View File
@@ -1,53 +0,0 @@
package tracing
import "context"
type tracingContextKey int
const (
CallerContext tracingContextKey = 1 + iota
)
type jobSubtree struct {
jobid string
}
type ctx struct {
parent *ctx
job *jobSubtree
ident string
}
var root = &ctx{nil, nil, ""}
func getParentOrRoot(c context.Context) *ctx {
parent, ok := c.Value(CallerContext).(*ctx)
if !ok {
parent = root
}
return parent
}
func makeChild(c context.Context, child *ctx) context.Context {
if child.parent == nil {
panic(child)
}
return context.WithValue(c, CallerContext, child)
}
func Child(c context.Context, ident string) context.Context {
parent := getParentOrRoot(c)
return makeChild(c, &ctx{parent: parent, ident: ident})
}
func GetStack(c context.Context) (idents []string) {
ct, ok := c.Value(CallerContext).(*ctx)
if !ok {
return idents
}
for ct.parent != nil {
idents = append(idents, ct.ident)
ct = ct.parent
}
return idents
}
-21
View File
@@ -1,21 +0,0 @@
package tracing
import (
"context"
"strings"
"testing"
"github.com/stretchr/testify/assert"
)
func TestIt(t *testing.T) {
ctx := context.Background()
ctx = Child(ctx, "a")
ctx = Child(ctx, "b")
ctx = Child(ctx, "c")
ctx = Child(ctx, "d")
assert.Equal(t, "dcba", strings.Join(GetStack(ctx), ""))
}
+1 -1
View File
@@ -32,7 +32,7 @@ func TLSListenerFactoryFromConfig(c *config.Global, in *config.TLSServe) (transp
serverCert, err := tls.LoadX509KeyPair(in.Cert, in.Key) serverCert, err := tls.LoadX509KeyPair(in.Cert, in.Key)
if err != nil { if err != nil {
return nil, errors.Wrap(err, "cannot parse cer/key pair") return nil, errors.Wrap(err, "cannot parse cert/key pair")
} }
clientCNs := make(map[string]struct{}, len(in.ClientCNs)) clientCNs := make(map[string]struct{}, len(in.ClientCNs))