Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
134 changes: 134 additions & 0 deletions oxiad/coordinator/metadata/metadata.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,13 @@ type Metadata interface {
GetShardStatus(namespace string, shard int64) (commonobject.Borrowed[*commonproto.ShardMetadata], bool)
UpdateShardStatus(namespace string, shard int64, shardMetadata *commonproto.ShardMetadata) error
DeleteShardStatus(namespace string, shard int64) error
// InitShardSplit records the start of a shard split: the parent is marked
// as splitting and both children are created in the same update, so a
// reader never sees one without the other. It fails with
// ErrFailedPrecondition when the parent cannot be split or the split point
// and child ids are not usable.
InitShardSplit(namespace string, parent, left, right int64, splitPoint uint32,
leftEnsemble, rightEnsemble []*commonproto.DataServerIdentity) error

GetConfig() commonobject.Borrowed[*commonproto.ClusterConfiguration]
SubscribeConfig() *commonwatch.Receiver[provider.Versioned[*commonproto.ClusterConfiguration]]
Expand Down Expand Up @@ -412,6 +419,133 @@ func (m *coordinatorMetadata) DeleteShardStatus(namespace string, shard int64) e
})
}

// InitShardSplit atomically records the start of a shard split. The parent is
// marked as splitting and both child shards are created within a single update,
// so no reader ever observes a parent that claims children that do not exist,
// or children whose parent is not splitting.
//
// Only the values that cannot be derived from the parent are taken from the
// caller: the reserved child ids, the split point and the children's ensembles.
// The children's hash ranges are cut from the parent's current range, and the
// rest of their metadata is built here, so a child can never disagree with the
// parent it was split from.
func (m *coordinatorMetadata) InitShardSplit(namespace string, parent, left, right int64, splitPoint uint32,
leftEnsemble, rightEnsemble []*commonproto.DataServerIdentity) error {
var precondition error
err := backoff.RetryNotify(func() error {
precondition = nil
return m.computeStatus(func(clusterStatus *commonproto.ClusterStatus, _ metadatacommon.Version) (*commonproto.ClusterStatus, bool) {
namespaceStatus, exist := clusterStatus.Namespaces[namespace]
if !exist {
precondition = fmt.Errorf("%w: namespace %q", metadatacommon.ErrNotFound, namespace)
return clusterStatus, false
}
parentMetadata, exist := namespaceStatus.Shards[parent]
if !exist {
precondition = fmt.Errorf("%w: shard %d of namespace %q", metadatacommon.ErrNotFound, parent, namespace)
return clusterStatus, false
}
if precondition = validateShardSplit(namespaceStatus, parentMetadata, parent, left, right,
splitPoint, leftEnsemble, rightEnsemble); precondition != nil {
return clusterStatus, false
}

hashRange := parentMetadata.GetInt32HashRange()
parentMetadata.Split = &commonproto.SplitMetadata{
Phase: commonproto.SplitPhaseBootstrap,
ChildShardIds: []int64{left, right},
SplitPoint: splitPoint,
}
namespaceStatus.Shards[left] = newChildShardMetadata(parent, splitPoint,
leftEnsemble, hashRange.GetMin(), splitPoint)
namespaceStatus.Shards[right] = newChildShardMetadata(parent, splitPoint,
rightEnsemble, splitPoint+1, hashRange.GetMax())
return clusterStatus, true
})
}, oxiatime.NewBackOff(m.ctx), func(err error, duration time.Duration) {
m.logger.Warn(
"failed to initiate the shard split",
slog.String("namespace", namespace),
slog.Int64("shard", parent),
slog.Any("error", err),
slog.Duration("retry-after", duration),
)
})
if err != nil {
return err
}
return precondition
}

// validateShardSplit checks, against the freshly loaded status, that the parent
// is in a state where a split can start and that the caller's child ids, split
// point and ensembles can be used to build the children.
func validateShardSplit(namespaceStatus *commonproto.NamespaceStatus, parentMetadata *commonproto.ShardMetadata,
parent, left, right int64, splitPoint uint32,
leftEnsemble, rightEnsemble []*commonproto.DataServerIdentity) error {
if status := parentMetadata.GetStatusOrDefault(); status != commonproto.ShardStatusSteadyState {
return fmt.Errorf("%w: shard %d is not in steady state (status=%s)",
metadatacommon.ErrFailedPrecondition, parent, status)
}
if parentMetadata.Split != nil {
return fmt.Errorf("%w: shard %d already has an active split",
metadatacommon.ErrFailedPrecondition, parent)
}
if len(parentMetadata.PendingDeleteShardNodes) > 0 {
return fmt.Errorf("%w: shard %d has pending ensemble changes",
metadatacommon.ErrFailedPrecondition, parent)
}
if left == right {
return fmt.Errorf("%w: the two children of shard %d must have distinct ids, got %d twice",
metadatacommon.ErrFailedPrecondition, parent, left)
}
for _, child := range []int64{left, right} {
if _, exist := namespaceStatus.Shards[child]; exist {
return fmt.Errorf("%w: child shard %d is already in use",
metadatacommon.ErrFailedPrecondition, child)
}
}
if len(leftEnsemble) == 0 || len(rightEnsemble) == 0 {
return fmt.Errorf("%w: both children of shard %d need an ensemble",
metadatacommon.ErrFailedPrecondition, parent)
}
hashRange := parentMetadata.GetInt32HashRange()
if hashRange.GetMax()-hashRange.GetMin() < 1 {
return fmt.Errorf("%w: shard %d hash range is too small to split",
metadatacommon.ErrFailedPrecondition, parent)
}
if splitPoint < hashRange.GetMin() || splitPoint >= hashRange.GetMax() {
return fmt.Errorf("%w: split point %d is outside the hash range [%d, %d] of shard %d",
metadatacommon.ErrFailedPrecondition, splitPoint, hashRange.GetMin(), hashRange.GetMax(), parent)
}
return nil
}

// newChildShardMetadata builds a child shard as it looks at birth: it serves
// its slice of the parent's hash range, and carries the split metadata that
// marks it as a child until the split completes.
func newChildShardMetadata(parent int64, splitPoint uint32,
ensemble []*commonproto.DataServerIdentity, minHash, maxHash uint32) *commonproto.ShardMetadata {
clonedEnsemble := make([]*commonproto.DataServerIdentity, len(ensemble))
for idx, dataServer := range ensemble {
clonedEnsemble[idx] = gproto.CloneOf(dataServer)
}
return &commonproto.ShardMetadata{
Status: commonproto.ShardStatusSteadyState,
Term: 0,
Ensemble: clonedEnsemble,
Int32HashRange: &commonproto.HashRange{
Min: minHash,
Max: maxHash,
},
Split: &commonproto.SplitMetadata{
Phase: commonproto.SplitPhaseBootstrap,
ParentShardId: parent,
SplitPoint: splitPoint,
},
}
}

func (m *coordinatorMetadata) GetConfig() commonobject.Borrowed[*commonproto.ClusterConfiguration] {
return commonobject.Borrow(m.configProvider.Watch().Load().Value)
}
Expand Down
152 changes: 152 additions & 0 deletions oxiad/coordinator/metadata/metadata_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -205,3 +205,155 @@
internal: s1:8191
`), 0600))
}

func newSplitTestMetadata(t *testing.T, parent *commonproto.ShardMetadata) Metadata {
t.Helper()

statusProvider := memory.NewProvider(metadatacodec.ClusterStatusCodec, metadataconstant.WatchDisabled, "")
configProvider := memory.NewProvider(metadatacodec.ClusterConfigCodec, metadataconstant.WatchEnabled, "")
metadata := newMetadata(t.Context(), statusProvider, configProvider, "")
require.True(t, metadata.CreateNamespaceStatus("default", &commonproto.NamespaceStatus{
Shards: map[int64]*commonproto.ShardMetadata{0: parent},
}))
t.Cleanup(func() { require.NoError(t, metadata.Close()) })
return metadata
}

func splittableShard() *commonproto.ShardMetadata {
return &commonproto.ShardMetadata{
Status: commonproto.ShardStatusSteadyState,
Term: 3,
Ensemble: []*commonproto.DataServerIdentity{{Internal: "s1:8191"}},
Int32HashRange: &commonproto.HashRange{Min: 0, Max: 100},
}
}

func splitEnsembles() ([]*commonproto.DataServerIdentity, []*commonproto.DataServerIdentity) {

Check failure on line 231 in oxiad/coordinator/metadata/metadata_test.go

View workflow job for this annotation

GitHub Actions / Build & Test

confusing-results: unnamed results of the same type may be confusing, consider using named results (revive)
return []*commonproto.DataServerIdentity{{Internal: "s1:8191"}},
[]*commonproto.DataServerIdentity{{Internal: "s2:8191"}}
}

func TestMetadataInitShardSplitCreatesParentAndChildrenTogether(t *testing.T) {
metadata := newSplitTestMetadata(t, splittableShard())
left, right := splitEnsembles()

require.NoError(t, metadata.InitShardSplit("default", 0, 1, 2, 40, left, right))

parent, exists := metadata.GetShardStatus("default", 0)
require.True(t, exists)
split := parent.UnsafeBorrow().GetSplit()
require.Equal(t, commonproto.SplitPhaseBootstrap, split.GetPhaseOrDefault())
require.Equal(t, []int64{1, 2}, split.GetChildShardIds())
require.EqualValues(t, 40, split.GetSplitPoint())
// The parent keeps serving its whole range until the split completes
require.EqualValues(t, 3, parent.UnsafeBorrow().GetTerm())
require.EqualValues(t, 0, parent.UnsafeBorrow().GetInt32HashRange().GetMin())
require.EqualValues(t, 100, parent.UnsafeBorrow().GetInt32HashRange().GetMax())

// The children partition the parent's range at the split point
leftChild, exists := metadata.GetShardStatus("default", 1)
require.True(t, exists)
require.EqualValues(t, 0, leftChild.UnsafeBorrow().GetInt32HashRange().GetMin())
require.EqualValues(t, 40, leftChild.UnsafeBorrow().GetInt32HashRange().GetMax())

rightChild, exists := metadata.GetShardStatus("default", 2)
require.True(t, exists)
require.EqualValues(t, 41, rightChild.UnsafeBorrow().GetInt32HashRange().GetMin())
require.EqualValues(t, 100, rightChild.UnsafeBorrow().GetInt32HashRange().GetMax())

// Each child starts at term 0, on its own ensemble, and points back at the parent
for shard, internal := range map[int64]string{1: "s1:8191", 2: "s2:8191"} {
child, exists := metadata.GetShardStatus("default", shard)
require.True(t, exists)
metadata := child.UnsafeBorrow()
require.EqualValues(t, 0, metadata.GetTerm())
require.Equal(t, commonproto.ShardStatusSteadyState, metadata.GetStatusOrDefault())
require.Len(t, metadata.GetEnsemble(), 1)
require.Equal(t, internal, metadata.GetEnsemble()[0].GetInternal())
require.EqualValues(t, 0, metadata.GetSplit().GetParentShardId())
require.EqualValues(t, 40, metadata.GetSplit().GetSplitPoint())
require.Empty(t, metadata.GetSplit().GetChildShardIds())
}
}

func TestMetadataInitShardSplitRejectsUnsplittableParent(t *testing.T) {
left, right := splitEnsembles()

for name, testCase := range map[string]struct {
parent *commonproto.ShardMetadata
splitPoint uint32
leftChild int64
rightChild int64
leftEnsemb []*commonproto.DataServerIdentity
rightEnsemb []*commonproto.DataServerIdentity
}{
"not in steady state": {parent: func() *commonproto.ShardMetadata {
shard := splittableShard()
shard.Status = commonproto.ShardStatusDeleting
return shard
}(), splitPoint: 40, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"split already active": {parent: func() *commonproto.ShardMetadata {
shard := splittableShard()
shard.Split = &commonproto.SplitMetadata{ChildShardIds: []int64{7, 8}}
return shard
}(), splitPoint: 40, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"pending ensemble changes": {parent: func() *commonproto.ShardMetadata {
shard := splittableShard()
shard.PendingDeleteShardNodes = []*commonproto.DataServerIdentity{{Internal: "s9:8191"}}
return shard
}(), splitPoint: 40, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"hash range too small": {parent: func() *commonproto.ShardMetadata {
shard := splittableShard()
shard.Int32HashRange = &commonproto.HashRange{Min: 7, Max: 7}
return shard
}(), splitPoint: 7, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"split point below the range": {parent: func() *commonproto.ShardMetadata {
shard := splittableShard()
shard.Int32HashRange = &commonproto.HashRange{Min: 10, Max: 100}
return shard
}(), splitPoint: 5, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"split point at the range end": {parent: splittableShard(), splitPoint: 100, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"child id already in use": {parent: splittableShard(), splitPoint: 40, leftChild: 0, rightChild: 2, leftEnsemb: left, rightEnsemb: right},
"children share an id": {parent: splittableShard(), splitPoint: 40, leftChild: 1, rightChild: 1, leftEnsemb: left, rightEnsemb: right},
"child without an ensemble": {parent: splittableShard(), splitPoint: 40, leftChild: 1, rightChild: 2, leftEnsemb: left, rightEnsemb: nil},
} {
t.Run(name, func(t *testing.T) {
metadata := newSplitTestMetadata(t, testCase.parent)

err := metadata.InitShardSplit("default", 0, testCase.leftChild, testCase.rightChild,
testCase.splitPoint, testCase.leftEnsemb, testCase.rightEnsemb)
require.ErrorIs(t, err, metadataconstant.ErrFailedPrecondition)

// Nothing was persisted: no child shards, and the parent is untouched
for _, shard := range []int64{1, 2} {
_, exists := metadata.GetShardStatus("default", shard)
require.False(t, exists)
}
parent, exists := metadata.GetShardStatus("default", 0)
require.True(t, exists)
require.Equal(t, testCase.parent.GetSplit(), parent.UnsafeBorrow().GetSplit())
})
}
}

func TestMetadataInitShardSplitFailsWhenTargetIsGone(t *testing.T) {
metadata := newSplitTestMetadata(t, splittableShard())
left, right := splitEnsembles()

require.ErrorIs(t, metadata.InitShardSplit("other", 0, 1, 2, 40, left, right), metadataconstant.ErrNotFound)
require.ErrorIs(t, metadata.InitShardSplit("default", 7, 1, 2, 40, left, right), metadataconstant.ErrNotFound)
}

func TestMetadataInitShardSplitKeepsCallerEnsemblesIsolated(t *testing.T) {
metadata := newSplitTestMetadata(t, splittableShard())
left, right := splitEnsembles()

require.NoError(t, metadata.InitShardSplit("default", 0, 1, 2, 40, left, right))

// A caller that reuses its ensemble slices must not reach into the status
left[0].Internal = "mutated:8191"

child, exists := metadata.GetShardStatus("default", 1)
require.True(t, exists)
require.Equal(t, "s1:8191", child.UnsafeBorrow().GetEnsemble()[0].GetInternal())
}
5 changes: 5 additions & 0 deletions oxiad/coordinator/reconciler/namespace_reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -183,6 +183,11 @@ func (m *mockNamespaceMetadata) UpdateShardStatus(namespace string, shard int64,
return nil
}

func (*mockNamespaceMetadata) InitShardSplit(string, int64, int64, int64, uint32,
[]*proto.DataServerIdentity, []*proto.DataServerIdentity) error {
return nil
}

func (*mockNamespaceMetadata) DeleteShardStatus(string, int64) error { return nil }

func (*mockNamespaceMetadata) CreateNamespace(*proto.Namespace) error {
Expand Down
1 change: 1 addition & 0 deletions oxiad/coordinator/runtime/action/action.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ type Type string
const (
SwapNode Type = "swap-node"
Election Type = "election"
Split Type = "split"
)

type Action interface {
Expand Down
Loading
Loading