Skip to content

Commit 5615f49

Browse files
authored
fix: replication of placeholder filesystems (#744)
fixes #742 Before this PR, when chaining replication from A => B => C, if B had placeholders and the `filesystems` included these placeholders, we'd incorrectly fail the planning phase with error `sender does not have any versions`. The non-placeholder child filesystems of these placeholders would then fail to replicate because of the initial-replication-dependency-tracking that we do, i.e., their parent failed to initially replication, hence they fail to replicate as well (`parent(s) failed during initial replication`). We can do better than that because we have the information whether a sender-side filesystem is a placeholder. This PR makes the planner act on that information. The outcome is that placeholders are replicated as placeholders (albeit the receiver remains in control of how these placeholders are created, i.e., `recv.placeholders`) The mechanism to do it is: 1. Don't plan any replication steps for filesystems that are placeholders on the sender. 2. Ensure that, if a receiving-side filesystem exists, it is indeed a placeholder. Check (2) may seem overly restrictive, but, the goal here is not just to mirror all non-placeholder filesystems, but also to mirror the hierarchy. Testing performed: - [x] confirm with issue reporter that this PR fixes their issue - [x] add a regression test that fails without the changes in this PR
1 parent 440b074 commit 5615f49

5 files changed

Lines changed: 158 additions & 23 deletions

File tree

‎endpoint/endpoint.go‎

Lines changed: 20 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -104,13 +104,28 @@ func (s *Sender) ListFilesystems(ctx context.Context, r *pdu.ListFilesystemReq)
104104
if err != nil {
105105
return nil, err
106106
}
107-
rfss := make([]*pdu.Filesystem, len(fss))
108-
for i := range fss {
109-
rfss[i] = &pdu.Filesystem{
110-
Path: fss[i].ToString(),
107+
rfss := make([]*pdu.Filesystem, 0, len(fss))
108+
for _, a := range fss {
109+
// TODO: dedup code with Receiver.ListFilesystems
110+
l := getLogger(ctx).WithField("fs", a)
111+
ph, err := zfs.ZFSGetFilesystemPlaceholderState(ctx, a)
112+
if err != nil {
113+
l.WithError(err).Error("error getting placeholder state")
114+
return nil, errors.Wrapf(err, "cannot get placeholder state for fs %q", a)
115+
}
116+
l.WithField("placeholder_state", fmt.Sprintf("%#v", ph)).Debug("placeholder state")
117+
if !ph.FSExists {
118+
l.Error("inconsistent placeholder state: filesystem must exists")
119+
err := errors.Errorf("inconsistent placeholder state: filesystem %q must exist in this context", a.ToString())
120+
return nil, err
121+
}
122+
123+
fs := &pdu.Filesystem{
124+
Path: a.ToString(),
111125
// ResumeToken does not make sense from Sender
112-
IsPlaceholder: false, // sender FSs are never placeholders
126+
IsPlaceholder: ph.IsPlaceholder,
113127
}
128+
rfss = append(rfss, fs)
114129
}
115130
res := &pdu.ListFilesystemRes{Filesystems: rfss}
116131
return res, nil

‎platformtest/tests/generated_cases.go‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎platformtest/tests/replication.go‎

Lines changed: 97 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1483,3 +1483,100 @@ func ReplicationInitialFail(ctx *platformtest.Context) {
14831483
require.NotNil(ctx, report.Attempts[0].Filesystems[0].PlanError)
14841484
require.Contains(ctx, report.Attempts[0].Filesystems[0].PlanError.Err, "automatic conflict resolution for initial replication is disabled in config")
14851485
}
1486+
1487+
// https://github.com/zrepl/zrepl/issues/742
1488+
func ReplicationOfPlaceholderFilesystemsInChainedReplicationScenario(ctx *platformtest.Context) {
1489+
1490+
//
1491+
// Setup datasets
1492+
//
1493+
platformtest.Run(ctx, platformtest.PanicErr, ctx.RootDataset, `
1494+
CREATEROOT
1495+
+ "host1"
1496+
+ "host1/a"
1497+
+ "host1/a/b"
1498+
+ "host1/a/b@1"
1499+
+ "host2"
1500+
+ "host2/sink"
1501+
+ "host3"
1502+
+ "host3/sink"
1503+
`)
1504+
1505+
// Replicate host1/a to host2/sink
1506+
host1_a := ctx.RootDataset + "/host1/a"
1507+
host1_b := ctx.RootDataset + "/host1/a/b"
1508+
1509+
host2_sink := ctx.RootDataset + "/host2/sink"
1510+
host2_a := host2_sink + "/" + host1_a
1511+
host2_b := host2_sink + "/" + host1_b
1512+
1513+
host3_sink := ctx.RootDataset + "/host3/sink"
1514+
host3_a := host3_sink + "/" + host2_a
1515+
host3_b := host3_sink + "/" + host2_b
1516+
1517+
type job struct {
1518+
sender, receiver endpoint.JobID
1519+
receiver_root string
1520+
sender_filesystems_filter map[string]bool
1521+
}
1522+
h1_to_h2 := job{
1523+
sender: endpoint.MustMakeJobID("h1-to-h2-sender"),
1524+
receiver: endpoint.MustMakeJobID("h1-to-h2-receiver"),
1525+
sender_filesystems_filter: map[string]bool{
1526+
// omit host1_a so that it becomes a placeholder on host2
1527+
host1_b: true,
1528+
},
1529+
receiver_root: host2_sink,
1530+
}
1531+
h2_to_h3 := job{
1532+
sender: endpoint.MustMakeJobID("h2-to-h3-sender"),
1533+
receiver: endpoint.MustMakeJobID("h2-to-h3-receiver"),
1534+
sender_filesystems_filter: map[string]bool{
1535+
host2_sink: false,
1536+
host2_sink + "<": true,
1537+
},
1538+
receiver_root: host3_sink,
1539+
}
1540+
1541+
do_repl := func(j job) {
1542+
1543+
sfilter := filters.NewDatasetMapFilter(len(j.sender_filesystems_filter), true)
1544+
1545+
for lhs, rhs := range j.sender_filesystems_filter {
1546+
var err error
1547+
if rhs {
1548+
err = sfilter.Add(lhs, filters.MapFilterResultOk)
1549+
} else {
1550+
err = sfilter.Add(lhs, filters.MapFilterResultOmit)
1551+
}
1552+
require.NoError(ctx, err)
1553+
}
1554+
1555+
rep := replicationInvocation{
1556+
sjid: j.sender,
1557+
rjid: j.receiver,
1558+
sfilter: sfilter,
1559+
rfsRoot: j.receiver_root,
1560+
guarantee: pdu.ReplicationConfigProtectionWithKind(pdu.ReplicationGuaranteeKind_GuaranteeResumability),
1561+
receiverConfigHook: func(rc *endpoint.ReceiverConfig) {
1562+
rc.PlaceholderEncryption = endpoint.PlaceholderCreationEncryptionPropertyOff
1563+
},
1564+
}
1565+
1566+
r := rep.Do(ctx)
1567+
ctx.Logf("\n%s", pretty.Sprint(r))
1568+
}
1569+
1570+
do_repl(h1_to_h2)
1571+
do_repl(h2_to_h3)
1572+
1573+
// assert that the replication worked
1574+
mustGetFilesystemVersion(ctx, host3_b+"@1")
1575+
1576+
// assert placeholder status
1577+
for _, fs := range []string{host2_a, host3_a} {
1578+
st, err := zfs.ZFSGetFilesystemPlaceholderState(ctx, mustDatasetPath(fs))
1579+
require.NoError(ctx, err)
1580+
require.True(ctx, st.IsPlaceholder)
1581+
}
1582+
}

‎replication/driver/replication_driver.go‎

Lines changed: 24 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -529,27 +529,31 @@ func (f *fs) do(ctx context.Context, pq *stepQueue, prev *fs) {
529529
// find the highest of the previously uncompleted steps for which we can also find a step
530530
// in our current plan
531531
prevUncompleted := prev.planned.steps[prev.planned.step:]
532-
var target struct{ prev, cur int }
533-
target.prev = -1
534-
target.cur = -1
535-
out:
536-
for p := len(prevUncompleted) - 1; p >= 0; p-- {
537-
for q := len(f.planned.steps) - 1; q >= 0; q-- {
538-
if prevUncompleted[p].step.TargetEquals(f.planned.steps[q].step) {
539-
target.prev = p
540-
target.cur = q
541-
break out
532+
if len(prevUncompleted) == 0 || len(f.planned.steps) == 0 {
533+
f.debug("no steps planned in previous attempt or this attempt, no correlation necessary len(prevUncompleted)=%d len(f.planned.steps)=%d", len(prevUncompleted), len(f.planned.steps))
534+
} else {
535+
var target struct{ prev, cur int }
536+
target.prev = -1
537+
target.cur = -1
538+
out:
539+
for p := len(prevUncompleted) - 1; p >= 0; p-- {
540+
for q := len(f.planned.steps) - 1; q >= 0; q-- {
541+
if prevUncompleted[p].step.TargetEquals(f.planned.steps[q].step) {
542+
target.prev = p
543+
target.cur = q
544+
break out
545+
}
542546
}
543547
}
544-
}
545-
if target.prev == -1 || target.cur == -1 {
546-
f.debug("no correlation possible between previous attempt and this attempt's plan")
547-
f.planning.err = newTimedError(fmt.Errorf("cannot correlate previously failed attempt to current plan"), time.Now())
548-
return
549-
}
548+
if target.prev == -1 || target.cur == -1 {
549+
f.debug("no correlation possible between previous attempt and this attempt's plan")
550+
f.planning.err = newTimedError(fmt.Errorf("cannot correlate previously failed attempt to current plan"), time.Now())
551+
return
552+
}
550553

551-
f.planned.steps = f.planned.steps[0:target.cur]
552-
f.debug("found correlation, new steps are len(fs.planned.steps) = %d", len(f.planned.steps))
554+
f.planned.steps = f.planned.steps[0:target.cur]
555+
f.debug("found correlation, new steps are len(fs.planned.steps) = %d", len(f.planned.steps))
556+
}
553557
} else {
554558
f.debug("previous attempt does not exist or did not finish planning, no correlation possible, taking this attempt's plan as is")
555559
}
@@ -600,6 +604,8 @@ func (f *fs) do(ctx context.Context, pq *stepQueue, prev *fs) {
600604
f.debug("parentHasNoSteps=%v parentFirstStepIsIncremental=%v parentHasTakenAtLeastOneSuccessfulStep=%v",
601605
parentHasNoSteps, parentFirstStepIsIncremental, parentHasTakenAtLeastOneSuccessfulStep)
602606

607+
// If the parent is a placeholder on the sender, `parentHasNoSteps` is true because we plan no steps for sender placeholders.
608+
// The receiver will create the necessary placeholders when they start receiving the first non-placeholder child filesystem.
603609
parentPresentOnReceiver := parentHasNoSteps || parentFirstStepIsIncremental || parentHasTakenAtLeastOneSuccessfulStep
604610

605611
allParentsPresentOnReceiver = allParentsPresentOnReceiver && parentPresentOnReceiver // no shadow

‎replication/logic/replication_logic.go‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -333,6 +333,22 @@ func (fs *Filesystem) doPlanning(ctx context.Context) ([]*Step, error) {
333333

334334
log(ctx).Debug("assessing filesystem")
335335

336+
if fs.senderFS.IsPlaceholder {
337+
log(ctx).Debug("sender filesystem is placeholder")
338+
if fs.receiverFS != nil {
339+
if fs.receiverFS.IsPlaceholder {
340+
// all good, fall through
341+
log(ctx).Debug("receiver filesystem is placeholder")
342+
} else {
343+
err := fmt.Errorf("sender filesystem is placeholder, but receiver filesystem is not")
344+
log(ctx).Error(err.Error())
345+
return nil, err
346+
}
347+
}
348+
log(ctx).Debug("no steps required for replicating placeholders, the endpoint.Receiver will create a placeholder when we receive the first non-placeholder child filesystem")
349+
return nil, nil
350+
}
351+
336352
sfsvsres, err := fs.sender.ListFilesystemVersions(ctx, &pdu.ListFilesystemVersionsReq{Filesystem: fs.Path})
337353
if err != nil {
338354
log(ctx).WithError(err).Error("cannot get remote filesystem versions")

0 commit comments

Comments
 (0)