Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
52 commits
Select commit Hold shift + click to select a range
ac75ab1
migration: add range version RPC handlers
bootjp Jul 19, 2026
0a1b9c8
Preserve staged key scan routing
bootjp Jul 19, 2026
e7815af
kv: cover staged tombstones in reverse auxiliary scans
bootjp Aug 22, 2026
3fa73ed
migration: exempt resolver-owned keys from the floor precheck
bootjp Aug 22, 2026
86c48e6
migration: mark route group on exact legacy list-delta scans
bootjp Aug 22, 2026
da3b967
migration: halt raft apply when an import fails on one voter
bootjp Aug 22, 2026
8668bdc
migration: tombstone staged rows on prefix deletes
bootjp Aug 23, 2026
b7aa882
Revert "migration: tombstone staged rows on prefix deletes"
bootjp Aug 23, 2026
ddfbdb0
migration: export filesystem chunk payloads
bootjp Aug 24, 2026
a32559a
migration: stop double-exporting chunks; normalize ownership lookups
bootjp Aug 24, 2026
98e9a1c
migration: export routed filesystem usage counters
bootjp Aug 24, 2026
b1896b7
migration: reject foreign explicit-group reads after promotion
bootjp Aug 27, 2026
c85adaa
migration: cap the promotion batch bounds server-side
bootjp Aug 27, 2026
b087bed
migration: bound the staged-visibility candidate range scan
bootjp Aug 27, 2026
a7c5967
migration: bound candidate probes to the exact key
bootjp Aug 27, 2026
0d52124
migration: reject user writes to migration control prefixes
bootjp Aug 28, 2026
5f07cc0
migration: cap export scan work server-side
bootjp Aug 28, 2026
9e0d8b6
migration: apply expiration to staged-only values
bootjp Aug 28, 2026
48a03d8
migration: preserve snapshot key limits when staging imports
bootjp Aug 28, 2026
473110b
Merge hotspot split store export base
bootjp Aug 28, 2026
5029acc
Fix FSM migration apply fences
bootjp Aug 28, 2026
f654d1b
Route promoted S3 bucket metadata
bootjp Aug 28, 2026
c044e7d
Batch staged prefix tombstones
bootjp Aug 28, 2026
11307fa
Preserve migration retry state
bootjp Aug 28, 2026
7b6d12d
Preserve legacy S3 auxiliary scans
bootjp Aug 28, 2026
3adfc31
Expire staged visibility winners
bootjp Aug 28, 2026
f99f909
Reject staged control raw writes
bootjp Aug 28, 2026
6059584
Route S3 auxiliary owner probes through leaders
bootjp Aug 28, 2026
70d852f
Preserve max-size keys during staged imports
bootjp Aug 28, 2026
adfcf26
Support group-aware owner probes
bootjp Aug 28, 2026
0539933
Merge branch 'design/hotspot-split-m2-store-export' of github.com:boo…
bootjp Aug 28, 2026
cd83c05
migration: resolve txn-wrapped keys through the partition resolver
bootjp Aug 28, 2026
3ad0687
migration: refuse export pages the transport cannot carry
bootjp Aug 28, 2026
f06fdf8
migration: split export pages before they overshoot the budget
bootjp Aug 28, 2026
8432cb7
Merge design/hotspot-split-m2-store-export into design/hotspot-split-…
bootjp Aug 28, 2026
d0dea91
kv: reject reserved control keys on the transactional paths
bootjp Aug 29, 2026
76e6b21
migration: order the staged and live probes against promotion
bootjp Aug 29, 2026
58104e0
kv: narrow the transactional reserved-key test to migration-only keys
bootjp Aug 29, 2026
ff29188
migration: read staged before live everywhere, not just the point read
bootjp Aug 29, 2026
3ca8aaf
Merge design/hotspot-split-m2-store-export into design/hotspot-split-…
bootjp Aug 29, 2026
d989293
kv: expect the staged owner probe first
bootjp Aug 29, 2026
2fb90ad
migration: read staged before live in OCC read validation too
bootjp Aug 29, 2026
1bbe007
store: emit the oldest snapshot layout that fits the state
bootjp Aug 29, 2026
0f499f9
Merge remote-tracking branch 'origin/design/hotspot-split-m2-store-ex…
bootjp Aug 31, 2026
712431a
migration: tighten staged transaction retries
bootjp Aug 31, 2026
7b9e199
migration: honor S3 auxiliary write fences
bootjp Aug 31, 2026
b16dd62
migration: narrow auxiliary write fences
bootjp Aug 31, 2026
ace20ae
migration: pin apply-time floor checks
bootjp Aug 31, 2026
986823c
migration: export filesystem auxiliary rows
bootjp Aug 31, 2026
da19644
migration: continue write fence validation
bootjp Aug 31, 2026
30be14d
migration: route grouped reverse scans
bootjp Aug 31, 2026
8fc01b3
migration: route bucket auxiliary ownership
bootjp Aug 31, 2026
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
77 changes: 76 additions & 1 deletion adapter/distribution_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,7 @@ var (
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")
errDistributionCatalogMutationInvalid = errors.New("catalog store mutation is invalid")
)

Expand Down Expand Up @@ -167,7 +168,7 @@ func (s *DistributionServer) GetRoute(ctx context.Context, req *pb.GetRouteReque
if err := s.requireReadReady(); err != nil {
return nil, err
}
r, ok := s.engine.GetRoute(kv.RouteKey(req.Key))
r, ok := s.engine.GetRoute(kv.RouteOwnershipKey(req.Key))
if !ok {
return &pb.GetRouteResponse{}, nil
}
Expand Down Expand Up @@ -200,6 +201,67 @@ 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())
Comment thread
bootjp marked this conversation as resolved.
if err != nil {
return nil, err
}
// Normalized exactly like GetRoute above. An internal storage key -- a
// filesystem chunk, a Redis collection row -- routes by its logical key,
// so looking the raw bytes up in the snapshot answers with the owner of
// the raw family prefix instead of the group that actually owned the key
// at that catalog version.
route, ok := snapshot.RouteOf(kv.RouteOwnershipKey(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
}

// GetCatalogCapabilities negotiates the durable delta-watch protocol.
func (s *DistributionServer) GetCatalogCapabilities(ctx context.Context, _ *pb.CatalogCapabilitiesRequest) (*pb.CatalogCapabilitiesResponse, error) {
if s.catalog == nil {
Expand Down Expand Up @@ -899,6 +961,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
243 changes: 242 additions & 1 deletion adapter/distribution_server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (

"github.com/bootjp/elastickv/distribution"
"github.com/bootjp/elastickv/internal/fskeys"
"github.com/bootjp/elastickv/internal/s3keys"
"github.com/bootjp/elastickv/kv"
pb "github.com/bootjp/elastickv/proto"
"github.com/bootjp/elastickv/store"
Expand Down Expand Up @@ -59,11 +60,37 @@ func TestDistributionServerGetRoute_NormalizesFilesystemChunkKeys(t *testing.T)
require.Equal(t, uint64(2), resp.RaftGroupId)
}

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

bucket := "bucket-a"
routeKey := s3keys.RoutePrefixForBucketAnyGeneration(bucket)
engine := distribution.NewEngine()
engine.UpdateRoute([]byte(""), routeKey, 1)
engine.UpdateRoute(routeKey, nil, 2)

s := NewDistributionServer(engine, nil)
for _, key := range [][]byte{
s3keys.BucketMetaKey(bucket),
s3keys.BucketGenerationKey(bucket),
} {
resp, err := s.GetRoute(context.Background(), &pb.GetRouteRequest{Key: key})
require.NoError(t, err)
require.Equal(t, routeKey, resp.Start)
require.Equal(t, uint64(2), resp.RaftGroupId)
}
}

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 +108,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 +235,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 Expand Up @@ -1250,3 +1427,67 @@ type recordingDistributionFilesystemObserver struct {
func (o *recordingDistributionFilesystemObserver) ObserveFilePinnedHotspot(reason string) {
o.reasons = append(o.reasons, reason)
}

// GetRouteOwnership answers the historical owner of a key, so it has to
// normalize the same way GetRoute does. An internal storage key -- here a
// filesystem chunk -- routes by its logical !fs|route|chk| key, and the raw
// !fs|chk| bytes sort into a different route entirely. Without normalization
// the RPC reports the raw-prefix owner rather than the group that owned the
// key at that catalog version.
func TestDistributionServerGetRouteOwnership_NormalizesFilesystemChunkKeys(t *testing.T) {
t.Parallel()

home := uint64(11)
inode := uint64(22)
routeKey := fskeys.ChunkRouteKey(home, inode)
engine := distribution.NewEngine()
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 3,
Routes: []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte(""), End: routeKey, GroupID: 1, State: distribution.RouteStateActive},
{RouteID: 2, Start: routeKey, End: nil, GroupID: 2, State: distribution.RouteStateActive},
},
}))

s := NewDistributionServer(engine, nil)
resp, err := s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: fskeys.ChunkKey(home, inode, 99),
CatalogVersion: 3,
})
require.NoError(t, err)
require.True(t, resp.Found)
require.Equal(t, uint64(2), resp.Route.RaftGroupId,
"the chunk must resolve to its logical route owner, not the raw-prefix owner")
require.Equal(t, uint64(2), resp.Route.RouteId)
}

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

bucket := "bucket-a"
routeKey := s3keys.RoutePrefixForBucketAnyGeneration(bucket)
engine := distribution.NewEngine()
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 3,
Routes: []distribution.RouteDescriptor{
{RouteID: 1, Start: []byte(""), End: routeKey, GroupID: 1, State: distribution.RouteStateActive},
{RouteID: 2, Start: routeKey, End: nil, GroupID: 2, State: distribution.RouteStateActive},
},
}))

s := NewDistributionServer(engine, nil)
for _, key := range [][]byte{
s3keys.BucketMetaKey(bucket),
s3keys.BucketGenerationKey(bucket),
} {
resp, err := s.GetRouteOwnership(context.Background(), &pb.GetRouteOwnershipRequest{
Key: key,
CatalogVersion: 3,
})
require.NoError(t, err)
require.True(t, resp.Found)
require.Equal(t, uint64(2), resp.Route.RaftGroupId,
"bucket auxiliary keys must resolve to the bucket route owner, not the raw-prefix owner")
require.Equal(t, uint64(2), resp.Route.RouteId)
}
}
Loading
Loading