Skip to content
Open
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
68 changes: 33 additions & 35 deletions adapter/distribution_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -455,7 +455,7 @@ func (s *DistributionServer) saveSplitResultViaCoordinator(
left.SplitAtHLC = commitTS
right.SplitAtHLC = commitTS

ops, err := buildCatalogSplitOps(parentID, left, right, nextVersion, nextRouteID)
ops, err := buildCatalogSplitOps(parentID, left, right, nextVersion, nextRouteID, s.catalog.AllowsRouteDescriptorV2Writes())
if err != nil {
return distribution.CatalogSnapshot{}, grpcStatusErrorf(codes.Internal, "build split mutations: %v", err)
}
Expand Down Expand Up @@ -522,6 +522,7 @@ func buildCatalogSplitOps(
right distribution.RouteDescriptor,
nextVersion uint64,
nextRouteID uint64,
allowRouteDescriptorV2Writes bool,
) ([]*kv.Elem[kv.OP], error) {
// SplitRange mutates the catalog surgically: delete one parent route, add two
// children, bump the version, and advance the next-route-id counter.
Expand All @@ -531,14 +532,10 @@ func buildCatalogSplitOps(
Key: distribution.CatalogRouteKey(parentID),
})
for _, route := range []distribution.RouteDescriptor{left, right} {
encoded, err := distribution.EncodeRouteDescriptor(route)
encoded, patchOffset, err := distribution.EncodeRouteDescriptorForCatalogWriteWithSplitAtHLCOffset(route, allowRouteDescriptorV2Writes)
if err != nil {
return nil, errors.WithStack(err)
}
patchOffset, err := splitAtHLCPatchOffset(encoded)
if err != nil {
return nil, err
}
ops = append(ops, &kv.Elem[kv.OP]{
Op: kv.Put,
Key: distribution.CatalogRouteKey(route.RouteID),
Expand All @@ -559,14 +556,6 @@ func buildCatalogSplitOps(
return ops, nil
}

func splitAtHLCPatchOffset(encoded []byte) (uint64, error) {
const splitAtHLCTailBytes = 8
if len(encoded) < splitAtHLCTailBytes {
return 0, errors.WithStack(distribution.ErrCatalogInvalidRouteRecord)
}
return uint64(len(encoded) - splitAtHLCTailBytes), nil //nolint:gosec // len was checked to be at least splitAtHLCTailBytes.
}

func splitChildrenFromSnapshot(snapshot distribution.CatalogSnapshot, leftID uint64, rightID uint64) (distribution.RouteDescriptor, distribution.RouteDescriptor, error) {
left, found := findRouteByID(snapshot.Routes, leftID)
if !found {
Expand Down Expand Up @@ -692,22 +681,28 @@ func splitCatalogRoutes(
) (distribution.RouteDescriptor, distribution.RouteDescriptor) {
// parent and splitKey are already cloned before this point and are immutable here.
left := distribution.RouteDescriptor{
RouteID: leftID,
Start: parent.Start,
End: splitKey,
GroupID: parent.GroupID,
State: parent.State,
ParentRouteID: parent.RouteID,
SplitAtHLC: splitAtHLC,
RouteID: leftID,
Start: parent.Start,
End: splitKey,
GroupID: parent.GroupID,
State: parent.State,
ParentRouteID: parent.RouteID,
SplitAtHLC: splitAtHLC,
StagedVisibilityActive: parent.StagedVisibilityActive,
MigrationJobID: parent.MigrationJobID,
MinWriteTSExclusive: parent.MinWriteTSExclusive,
}
right := distribution.RouteDescriptor{
RouteID: rightID,
Start: splitKey,
End: parent.End,
GroupID: parent.GroupID,
State: parent.State,
ParentRouteID: parent.RouteID,
SplitAtHLC: splitAtHLC,
RouteID: rightID,
Start: splitKey,
End: parent.End,
GroupID: parent.GroupID,
State: parent.State,
ParentRouteID: parent.RouteID,
SplitAtHLC: splitAtHLC,
StagedVisibilityActive: parent.StagedVisibilityActive,
MigrationJobID: parent.MigrationJobID,
MinWriteTSExclusive: parent.MinWriteTSExclusive,
}
return left, right
}
Expand Down Expand Up @@ -757,13 +752,16 @@ func toProtoRouteDescriptors(routes []distribution.RouteDescriptor) []*pb.RouteD

func toProtoRouteDescriptor(route distribution.RouteDescriptor) *pb.RouteDescriptor {
return &pb.RouteDescriptor{
RouteId: route.RouteID,
Start: distribution.CloneBytes(route.Start),
End: distribution.CloneBytes(route.End),
RaftGroupId: route.GroupID,
State: toProtoRouteState(route.State),
ParentRouteId: route.ParentRouteID,
SplitAtHlc: route.SplitAtHLC,
RouteId: route.RouteID,
Start: distribution.CloneBytes(route.Start),
End: distribution.CloneBytes(route.End),
RaftGroupId: route.GroupID,
State: toProtoRouteState(route.State),
ParentRouteId: route.ParentRouteID,
SplitAtHlc: route.SplitAtHLC,
StagedVisibilityActive: route.StagedVisibilityActive,
MigrationJobId: route.MigrationJobID,
MinWriteTsExclusive: route.MinWriteTSExclusive,
Comment thread
bootjp marked this conversation as resolved.
}
}

Expand Down
47 changes: 30 additions & 17 deletions adapter/distribution_server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,16 +97,19 @@ func TestDistributionServerListRoutes_ReadsDurableCatalog(t *testing.T) {
t.Parallel()

ctx := context.Background()
catalog := distribution.NewCatalogStore(store.NewMVCCStore())
catalog := distribution.NewCatalogStore(store.NewMVCCStore(), distribution.WithCatalogRouteDescriptorV2Writes(true))
saved, err := catalog.Save(ctx, 0, []distribution.RouteDescriptor{
{
RouteID: 2,
Start: []byte("m"),
End: nil,
GroupID: 2,
State: distribution.RouteStateWriteFenced,
ParentRouteID: 1,
SplitAtHLC: 99,
RouteID: 2,
Start: []byte("m"),
End: nil,
GroupID: 2,
State: distribution.RouteStateWriteFenced,
ParentRouteID: 1,
SplitAtHLC: 77,
StagedVisibilityActive: true,
MigrationJobID: 42,
MinWriteTSExclusive: 99,
},
{
RouteID: 1,
Expand All @@ -133,7 +136,10 @@ func TestDistributionServerListRoutes_ReadsDurableCatalog(t *testing.T) {
require.Equal(t, uint64(2), resp.Routes[1].RouteId)
require.Nil(t, resp.Routes[1].End)
require.Equal(t, pb.RouteState_ROUTE_STATE_WRITE_FENCED, resp.Routes[1].State)
require.Equal(t, uint64(99), resp.Routes[1].SplitAtHlc)
require.Equal(t, uint64(77), resp.Routes[1].SplitAtHlc)
require.True(t, resp.Routes[1].StagedVisibilityActive)
require.Equal(t, uint64(42), resp.Routes[1].MigrationJobId)
require.Equal(t, uint64(99), resp.Routes[1].MinWriteTsExclusive)
}

func TestDistributionServerListRoutes_RequiresCatalog(t *testing.T) {
Expand All @@ -151,15 +157,16 @@ func TestDistributionServerSplitRange_Success(t *testing.T) {

ctx := context.Background()
baseStore := store.NewMVCCStore()
catalog := distribution.NewCatalogStore(baseStore)
catalog := distribution.NewCatalogStore(baseStore, distribution.WithCatalogRouteDescriptorV2Writes(true))
saved, err := catalog.Save(ctx, 0, []distribution.RouteDescriptor{
{
RouteID: 1,
Start: []byte(""),
End: []byte("m"),
GroupID: 1,
State: distribution.RouteStateActive,
ParentRouteID: 0,
RouteID: 1,
Start: []byte(""),
End: []byte("m"),
GroupID: 1,
State: distribution.RouteStateActive,
ParentRouteID: 0,
MinWriteTSExclusive: 99,
},
{
RouteID: 2,
Expand Down Expand Up @@ -192,21 +199,25 @@ func TestDistributionServerSplitRange_Success(t *testing.T) {
require.Equal(t, []byte("g"), resp.Left.End)
require.Equal(t, uint64(1), resp.Left.RaftGroupId)
require.Equal(t, uint64(1), resp.Left.ParentRouteId)
require.Equal(t, uint64(99), resp.Left.MinWriteTsExclusive)
require.Equal(t, uint64(4), resp.Right.RouteId)
require.Equal(t, []byte("g"), resp.Right.Start)
require.Equal(t, []byte("m"), resp.Right.End)
require.Equal(t, uint64(1), resp.Right.RaftGroupId)
require.Equal(t, uint64(1), resp.Right.ParentRouteId)
require.NotZero(t, resp.Left.SplitAtHlc)
require.Equal(t, resp.Left.SplitAtHlc, resp.Right.SplitAtHlc)
require.Equal(t, uint64(99), resp.Right.MinWriteTsExclusive)

snapshot, err := catalog.Snapshot(ctx)
require.NoError(t, err)
require.Equal(t, uint64(2), snapshot.Version)
require.Len(t, snapshot.Routes, 3)
// Catalog snapshots are sorted by range start key.
require.Equal(t, uint64(3), snapshot.Routes[0].RouteID)
require.Equal(t, uint64(99), snapshot.Routes[0].MinWriteTSExclusive)
require.Equal(t, uint64(4), snapshot.Routes[1].RouteID)
require.Equal(t, uint64(99), snapshot.Routes[1].MinWriteTSExclusive)
require.Equal(t, uint64(2), snapshot.Routes[2].RouteID)
require.NotZero(t, snapshot.Routes[0].SplitAtHLC)
require.Equal(t, snapshot.Routes[0].SplitAtHLC, snapshot.Routes[1].SplitAtHLC)
Expand All @@ -216,9 +227,11 @@ func TestDistributionServerSplitRange_Success(t *testing.T) {
leftRoute, ok := engine.GetRoute([]byte("b"))
require.True(t, ok)
require.Equal(t, uint64(3), leftRoute.RouteID)
require.Equal(t, uint64(99), leftRoute.MinWriteTSExclusive)
rightRoute, ok := engine.GetRoute([]byte("h"))
require.True(t, ok)
require.Equal(t, uint64(4), rightRoute.RouteID)
require.Equal(t, uint64(99), rightRoute.MinWriteTSExclusive)
}

func TestDistributionServerSplitRange_SnapsFilesystemChunkKeyToFileBoundary(t *testing.T) {
Expand Down Expand Up @@ -754,7 +767,7 @@ func TestBuildCatalogSplitOps_UsesSurgicalSplitMutations(t *testing.T) {
ParentRouteID: 1,
}

ops, err := buildCatalogSplitOps(1, left, right, 2, 5)
ops, err := buildCatalogSplitOps(1, left, right, 2, 5, false)
require.NoError(t, err)
require.Len(t, ops, 5)
require.Equal(t, kv.Del, ops[0].Op)
Expand Down
Loading
Loading