remote snapshot destruction & replication status zfs property
This commit is contained in:
+87
-31
@@ -7,7 +7,6 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"math/bits"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -41,6 +40,7 @@ type Sender interface {
|
||||
// If a non-nil io.ReadCloser is returned, it is guaranteed to be closed before
|
||||
// any next call to the parent github.com/zrepl/zrepl/replication.Endpoint.
|
||||
Send(ctx context.Context, r *pdu.SendReq) (*pdu.SendRes, io.ReadCloser, error)
|
||||
SnapshotReplicationStatus(ctx context.Context, r *pdu.SnapshotReplicationStatusReq) (*pdu.SnapshotReplicationStatusRes, error)
|
||||
}
|
||||
|
||||
// A Sender is usually part of a github.com/zrepl/zrepl/replication.Endpoint.
|
||||
@@ -76,17 +76,13 @@ const (
|
||||
)
|
||||
|
||||
func (s State) fsrsf() state {
|
||||
idx := bits.TrailingZeros(uint(s))
|
||||
if idx == bits.UintSize {
|
||||
panic(s)
|
||||
m := map[State]state{
|
||||
Ready: stateReady,
|
||||
RetryWait: stateRetryWait,
|
||||
PermanentError: nil,
|
||||
Completed: nil,
|
||||
}
|
||||
m := []state{
|
||||
stateReady,
|
||||
stateRetryWait,
|
||||
nil,
|
||||
nil,
|
||||
}
|
||||
return m[idx]
|
||||
return m[s]
|
||||
}
|
||||
|
||||
type Replication struct {
|
||||
@@ -115,7 +111,7 @@ func BuildReplication(fs string) *ReplicationBuilder {
|
||||
|
||||
func (b *ReplicationBuilder) AddStep(from, to FilesystemVersion) *ReplicationBuilder {
|
||||
step := &ReplicationStep{
|
||||
state: StepReady,
|
||||
state: StepReplicationReady,
|
||||
parent: b.r,
|
||||
from: from,
|
||||
to: to,
|
||||
@@ -147,15 +143,18 @@ func NewReplicationWithPermanentError(fs string, err error) *Replication {
|
||||
type StepState uint
|
||||
|
||||
const (
|
||||
StepReady StepState = 1 << iota
|
||||
StepRetry
|
||||
StepReplicationReady StepState = 1 << iota
|
||||
StepReplicationRetry
|
||||
StepMarkReplicatedReady
|
||||
StepMarkReplicatedRetry
|
||||
StepPermanentError
|
||||
StepCompleted
|
||||
)
|
||||
|
||||
type FilesystemVersion interface {
|
||||
SnapshotTime() time.Time
|
||||
RelName() string
|
||||
GetName() string // name without @ or #
|
||||
RelName() string // name with @ or #
|
||||
}
|
||||
|
||||
type ReplicationStep struct {
|
||||
@@ -223,7 +222,7 @@ func stateReady(ctx context.Context, sender Sender, receiver Receiver, u updater
|
||||
return s.fsrsf()
|
||||
}
|
||||
|
||||
stepState := current.do(ctx, sender, receiver)
|
||||
stepState := current.doReplication(ctx, sender, receiver)
|
||||
|
||||
return u(func(f *Replication) {
|
||||
switch stepState {
|
||||
@@ -235,7 +234,9 @@ func stateReady(ctx context.Context, sender Sender, receiver Receiver, u updater
|
||||
} else {
|
||||
f.state = Completed
|
||||
}
|
||||
case StepRetry:
|
||||
case StepReplicationRetry:
|
||||
fallthrough
|
||||
case StepMarkReplicatedRetry:
|
||||
f.retryWaitUntil = time.Now().Add(RetrySleepDuration)
|
||||
f.state = RetryWait
|
||||
case StepPermanentError:
|
||||
@@ -292,7 +293,22 @@ func (fsr *Replication) Report() *Report {
|
||||
return &rep
|
||||
}
|
||||
|
||||
func (s *ReplicationStep) do(ctx context.Context, sender Sender, receiver Receiver) StepState {
|
||||
func shouldRetry(err error) bool {
|
||||
switch err {
|
||||
case io.EOF:
|
||||
fallthrough
|
||||
case io.ErrUnexpectedEOF:
|
||||
fallthrough
|
||||
case io.ErrClosedPipe:
|
||||
return true
|
||||
}
|
||||
if _, ok := err.(net.Error); ok {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func (s *ReplicationStep) doReplication(ctx context.Context, sender Sender, receiver Receiver) StepState {
|
||||
|
||||
fs := s.parent.fs
|
||||
|
||||
@@ -305,17 +321,8 @@ func (s *ReplicationStep) do(ctx context.Context, sender Sender, receiver Receiv
|
||||
defer s.lock.Unlock()
|
||||
|
||||
s.err = err
|
||||
switch err {
|
||||
case io.EOF:
|
||||
fallthrough
|
||||
case io.ErrUnexpectedEOF:
|
||||
fallthrough
|
||||
case io.ErrClosedPipe:
|
||||
s.state = StepRetry
|
||||
return s.state
|
||||
}
|
||||
if _, ok := err.(net.Error); ok {
|
||||
s.state = StepRetry
|
||||
if shouldRetry(s.err) {
|
||||
s.state = StepReplicationRetry
|
||||
return s.state
|
||||
}
|
||||
s.state = StepPermanentError
|
||||
@@ -326,7 +333,7 @@ func (s *ReplicationStep) do(ctx context.Context, sender Sender, receiver Receiv
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
s.err = nil
|
||||
s.state = StepCompleted
|
||||
s.state = StepMarkReplicatedReady
|
||||
return s.state
|
||||
}
|
||||
|
||||
@@ -371,8 +378,57 @@ func (s *ReplicationStep) do(ctx context.Context, sender Sender, receiver Receiv
|
||||
return updateStateError(err)
|
||||
}
|
||||
log.Info("receive finished")
|
||||
return updateStateCompleted()
|
||||
|
||||
updateStateCompleted()
|
||||
|
||||
return s.doMarkReplicated(ctx, sender)
|
||||
|
||||
}
|
||||
|
||||
func (s *ReplicationStep) doMarkReplicated(ctx context.Context, sender Sender) StepState {
|
||||
|
||||
log := getLogger(ctx).
|
||||
WithField("filesystem", s.parent.fs).
|
||||
WithField("step", s.String())
|
||||
|
||||
updateStateError := func(err error) StepState {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
|
||||
s.err = err
|
||||
if shouldRetry(s.err) {
|
||||
s.state = StepMarkReplicatedRetry
|
||||
return s.state
|
||||
}
|
||||
s.state = StepPermanentError
|
||||
return s.state
|
||||
}
|
||||
|
||||
updateStateCompleted := func() StepState {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
s.state = StepCompleted
|
||||
return s.state
|
||||
}
|
||||
|
||||
log.Info("mark snapshot as replicated")
|
||||
req := pdu.SnapshotReplicationStatusReq{
|
||||
Filesystem: s.parent.fs,
|
||||
Snapshot: s.to.GetName(),
|
||||
Op: pdu.SnapshotReplicationStatusReq_SetReplicated,
|
||||
}
|
||||
res, err := sender.SnapshotReplicationStatus(ctx, &req)
|
||||
if err != nil {
|
||||
log.WithError(err).Error("error marking snapshot as replicated")
|
||||
return updateStateError(err)
|
||||
}
|
||||
if res.Replicated != true {
|
||||
err := fmt.Errorf("sender did not report snapshot as replicated")
|
||||
log.Error(err.Error())
|
||||
return updateStateError(err)
|
||||
}
|
||||
|
||||
return updateStateCompleted()
|
||||
}
|
||||
|
||||
func (s *ReplicationStep) String() string {
|
||||
|
||||
@@ -5,13 +5,15 @@ package fsrep
|
||||
import "strconv"
|
||||
|
||||
const (
|
||||
_StepState_name_0 = "StepReadyStepRetry"
|
||||
_StepState_name_1 = "StepPermanentError"
|
||||
_StepState_name_2 = "StepCompleted"
|
||||
_StepState_name_0 = "StepReplicationReadyStepReplicationRetry"
|
||||
_StepState_name_1 = "StepMarkReplicatedReady"
|
||||
_StepState_name_2 = "StepMarkReplicatedRetry"
|
||||
_StepState_name_3 = "StepPermanentError"
|
||||
_StepState_name_4 = "StepCompleted"
|
||||
)
|
||||
|
||||
var (
|
||||
_StepState_index_0 = [...]uint8{0, 9, 18}
|
||||
_StepState_index_0 = [...]uint8{0, 20, 40}
|
||||
)
|
||||
|
||||
func (i StepState) String() string {
|
||||
@@ -23,6 +25,10 @@ func (i StepState) String() string {
|
||||
return _StepState_name_1
|
||||
case i == 8:
|
||||
return _StepState_name_2
|
||||
case i == 16:
|
||||
return _StepState_name_3
|
||||
case i == 32:
|
||||
return _StepState_name_4
|
||||
default:
|
||||
return "StepState(" + strconv.FormatInt(int64(i), 10) + ")"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user