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
88 changes: 79 additions & 9 deletions adapter/distribution_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,15 +92,16 @@ var (
defaultCatalogReloadRetryAttempts = 20
defaultCatalogReloadRetryInterval = 10 * time.Millisecond

errDistributionCatalogNotConfigured = errors.New("route catalog is not configured")
errDistributionUnknownRoute = errors.New("unknown route")
errDistributionInvalidSplitKey = errors.New("invalid split key")
errDistributionSplitKeyAtBoundary = errors.New("split key at route boundary")
errDistributionCatalogConflict = errors.New("catalog version conflict")
errDistributionRouteIDOverflow = errors.New("route id overflow")
errDistributionNotLeader = errors.New("not leader for distribution catalog")
errDistributionCoordinatorRequired = errors.New("distribution coordinator is not configured")
errDistributionEngineNotConfigured = errors.New("distribution engine is not configured")
errDistributionCatalogNotConfigured = errors.New("route catalog is not configured")
errDistributionUnknownRoute = errors.New("unknown route")
errDistributionInvalidSplitKey = errors.New("invalid split key")
errDistributionSplitKeyAtBoundary = errors.New("split key at route boundary")
errDistributionCatalogConflict = errors.New("catalog version conflict")
errDistributionRouteIDOverflow = errors.New("route id overflow")
errDistributionNotLeader = errors.New("not leader for distribution catalog")
errDistributionCoordinatorRequired = errors.New("distribution coordinator is not configured")
errDistributionEngineNotConfigured = errors.New("distribution engine is not configured")
errDistributionCatalogVersionNotFound = errors.New("route catalog version not found")
)

// NewDistributionServer creates a new server.
Expand Down Expand Up @@ -176,6 +177,62 @@ func (s *DistributionServer) ListRoutes(ctx context.Context, req *pb.ListRoutesR
}, nil
}

func (s *DistributionServer) GetRouteOwnership(ctx context.Context, req *pb.GetRouteOwnershipRequest) (*pb.GetRouteOwnershipResponse, error) {
if err := s.requireReadReady(); err != nil {
return nil, err
}
snapshot, err := s.routeSnapshotAt(req.GetCatalogVersion())

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Gate route-history reads during startup

When readBlocked is set during startup/rotation, GetRoute and ListRoutes return Unavailable, but this new route-history path skips requireReadReady() and reads the in-memory history directly. A client or migrator calling GetRouteOwnership (and the analogous GetIntersectingRoutes below) before the catalog watcher has applied the durable snapshot can receive stale or NotFound ownership and route/fence migration work against the wrong catalog state; add the same read-gate check before routeSnapshotAt.

Useful? React with 👍 / 👎.

if err != nil {
return nil, err
}
route, ok := snapshot.RouteOf(req.GetKey())
if !ok {
return &pb.GetRouteOwnershipResponse{
CatalogVersion: snapshot.Version(),
Found: false,
}, nil
}
return &pb.GetRouteOwnershipResponse{
Route: toProtoRoute(route),
CatalogVersion: snapshot.Version(),
Found: true,
}, nil
}

func (s *DistributionServer) GetIntersectingRoutes(ctx context.Context, req *pb.GetIntersectingRoutesRequest) (*pb.GetIntersectingRoutesResponse, error) {
if err := s.requireReadReady(); err != nil {
return nil, err
}
snapshot, err := s.routeSnapshotAt(req.GetCatalogVersion())
if err != nil {
return nil, err
}
end := req.GetEnd()
if len(end) == 0 {
end = nil
}
routes := snapshot.IntersectingRoutes(req.GetStart(), end)
out := make([]*pb.RouteDescriptor, 0, len(routes))
for _, route := range routes {
out = append(out, toProtoRoute(route))
}
return &pb.GetIntersectingRoutesResponse{
Routes: out,
CatalogVersion: snapshot.Version(),
}, nil
}

func (s *DistributionServer) routeSnapshotAt(version uint64) (distribution.RouteHistorySnapshot, error) {
if s.engine == nil {
return distribution.RouteHistorySnapshot{}, grpcStatusError(codes.FailedPrecondition, errDistributionEngineNotConfigured.Error())
}
snapshot, ok := s.engine.SnapshotAt(version)
if !ok {
return distribution.RouteHistorySnapshot{}, grpcStatusErrorf(codes.NotFound, "%s: %d", errDistributionCatalogVersionNotFound, version)
}
return snapshot, nil
}

// SplitRange splits a route into two child routes in the same raft group.
func (s *DistributionServer) SplitRange(ctx context.Context, req *pb.SplitRangeRequest) (*pb.SplitRangeResponse, error) {
// SplitRange performs a read-modify-write cycle across catalog and engine.
Expand Down Expand Up @@ -659,6 +716,19 @@ func toProtoRouteDescriptor(route distribution.RouteDescriptor) *pb.RouteDescrip
}
}

func toProtoRoute(route distribution.Route) *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),
StagedVisibilityActive: route.StagedVisibilityActive,
MigrationJobId: route.MigrationJobID,
MinWriteTsExclusive: route.MinWriteTSExclusive,
}
}

func toProtoRouteState(state distribution.RouteState) pb.RouteState {
switch state {
case distribution.RouteStateActive:
Expand Down
157 changes: 156 additions & 1 deletion adapter/distribution_server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,12 @@ func TestDistributionServerRouteReadsHonorStartupGate(t *testing.T) {
t.Parallel()

engine := distribution.NewEngine()
engine.UpdateRoute([]byte("a"), nil, 1)
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 1,
Routes: []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte("a"), End: nil, GroupID: 1, State: distribution.RouteStateActive},
},
}))
catalog := distribution.NewCatalogStore(store.NewMVCCStore())
_, err := catalog.Save(context.Background(), 0, []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte("a"), End: nil, GroupID: 1, State: distribution.RouteStateActive},
Expand All @@ -81,11 +86,37 @@ func TestDistributionServerRouteReadsHonorStartupGate(t *testing.T) {
require.Error(t, err)
require.Equal(t, codes.Unavailable, status.Code(err))

_, err = s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: []byte("a"),
CatalogVersion: engine.Version(),
})
require.Error(t, err)
require.Equal(t, codes.Unavailable, status.Code(err))

_, err = s.GetIntersectingRoutes(context.Background(), &pb.GetIntersectingRoutesRequest{
Start: []byte("a"),
End: []byte("z"),
CatalogVersion: engine.Version(),
})
require.Error(t, err)
require.Equal(t, codes.Unavailable, status.Code(err))

blocked = false
_, err = s.GetRoute(context.Background(), &pb.GetRouteRequest{Key: []byte("a")})
require.NoError(t, err)
_, err = s.ListRoutes(context.Background(), &pb.ListRoutesRequest{})
require.NoError(t, err)
_, err = s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: []byte("a"),
CatalogVersion: engine.Version(),
})
require.NoError(t, err)
_, err = s.GetIntersectingRoutes(context.Background(), &pb.GetIntersectingRoutesRequest{
Start: []byte("a"),
End: []byte("z"),
CatalogVersion: engine.Version(),
})
require.NoError(t, err)
}

func TestDistributionServerGetTimestamp_IsMonotonic(t *testing.T) {
Expand Down Expand Up @@ -182,6 +213,130 @@ func TestDistributionServerListRoutes_RequiresCatalog(t *testing.T) {
require.ErrorContains(t, err, errDistributionCatalogNotConfigured.Error())
}

func TestDistributionServerGetRouteOwnership_UsesExactVersionSnapshot(t *testing.T) {
t.Parallel()

engine := distribution.NewEngine()
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 7,
Routes: []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte("a"), End: []byte("m"), GroupID: 1, State: distribution.RouteStateActive},
{
RouteID: 2,
Start: []byte("m"),
End: nil,
GroupID: 2,
State: distribution.RouteStateMigratingTarget,
StagedVisibilityActive: true,
MigrationJobID: 44,
MinWriteTSExclusive: 55,
},
},
}))

s := NewDistributionServer(engine, nil)
resp, err := s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: []byte("t"),
CatalogVersion: 7,
})
require.NoError(t, err)
require.True(t, resp.Found)
require.Equal(t, uint64(7), resp.CatalogVersion)
require.Equal(t, uint64(2), resp.Route.RouteId)
require.Equal(t, uint64(2), resp.Route.RaftGroupId)
require.Equal(t, pb.RouteState_ROUTE_STATE_MIGRATING_TARGET, resp.Route.State)
require.True(t, resp.Route.StagedVisibilityActive)
require.Equal(t, uint64(44), resp.Route.MigrationJobId)
require.Equal(t, uint64(55), resp.Route.MinWriteTsExclusive)

miss, err := s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: []byte("0"),
CatalogVersion: 7,
})
require.NoError(t, err)
require.False(t, miss.Found)
require.Equal(t, uint64(7), miss.CatalogVersion)
require.Nil(t, miss.Route)
}

func TestDistributionServerGetIntersectingRoutes_UsesExactVersionSnapshot(t *testing.T) {
t.Parallel()

engine := distribution.NewEngine()
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 9,
Routes: []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte(""), End: []byte("g"), GroupID: 1, State: distribution.RouteStateActive},
{RouteID: 2, Start: []byte("g"), End: []byte("m"), GroupID: 2, State: distribution.RouteStateWriteFenced},
{RouteID: 3, Start: []byte("m"), End: nil, GroupID: 3, State: distribution.RouteStateActive},
},
}))

s := NewDistributionServer(engine, nil)
resp, err := s.GetIntersectingRoutes(context.Background(), &pb.GetIntersectingRoutesRequest{
Start: []byte("f"),
End: []byte("z"),
CatalogVersion: 9,
})
require.NoError(t, err)
require.Equal(t, uint64(9), resp.CatalogVersion)
require.Len(t, resp.Routes, 3)
require.Equal(t, []uint64{1, 2, 3}, []uint64{resp.Routes[0].RouteId, resp.Routes[1].RouteId, resp.Routes[2].RouteId})

rightOpen, err := s.GetIntersectingRoutes(context.Background(), &pb.GetIntersectingRoutesRequest{
Start: []byte("m"),
End: nil,
CatalogVersion: 9,
})
require.NoError(t, err)
require.Len(t, rightOpen.Routes, 1)
require.Equal(t, uint64(3), rightOpen.Routes[0].RouteId)
}

func TestDistributionServerOwnershipRPCs_RejectUnknownCatalogVersion(t *testing.T) {
t.Parallel()

engine := distribution.NewEngine()
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 1,
Routes: []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte(""), End: nil, GroupID: 1, State: distribution.RouteStateActive},
},
}))

s := NewDistributionServer(engine, nil)
_, err := s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: []byte("a"),
CatalogVersion: 2,
})
require.Error(t, err)
require.Equal(t, codes.NotFound, status.Code(err))
require.ErrorContains(t, err, errDistributionCatalogVersionNotFound.Error())

_, err = s.GetIntersectingRoutes(context.Background(), &pb.GetIntersectingRoutesRequest{
Start: []byte(""),
CatalogVersion: 2,
})
require.Error(t, err)
require.Equal(t, codes.NotFound, status.Code(err))
require.ErrorContains(t, err, errDistributionCatalogVersionNotFound.Error())
}

func TestDistributionServerOwnershipRPCs_RequireEngine(t *testing.T) {
t.Parallel()

s := NewDistributionServer(nil, nil)
_, err := s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{CatalogVersion: 1})
require.Error(t, err)
require.Equal(t, codes.FailedPrecondition, status.Code(err))
require.ErrorContains(t, err, errDistributionEngineNotConfigured.Error())

_, err = s.GetIntersectingRoutes(context.Background(), &pb.GetIntersectingRoutesRequest{CatalogVersion: 1})
require.Error(t, err)
require.Equal(t, codes.FailedPrecondition, status.Code(err))
require.ErrorContains(t, err, errDistributionEngineNotConfigured.Error())
}

func TestDistributionServerSplitRange_Success(t *testing.T) {
t.Parallel()

Expand Down
Loading
Loading