Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
209 commits
Select commit Hold shift + click to select a range
43b4d73
store: add migration version import export
bootjp Jul 13, 2026
edf74ff
store: tighten migration export accounting
bootjp Jul 13, 2026
f95b557
distribution: add migration bracket planner
bootjp Jul 13, 2026
af847ef
Fix migration bracket edge cases
bootjp Jul 13, 2026
eab9622
Add migration fence drain guards
bootjp Jul 13, 2026
07a428e
Address migration guard review notes
bootjp Jul 13, 2026
7c17803
kv: fence broad mapped prefix deletes
bootjp Jul 13, 2026
2a85852
migration: add range version RPC handlers
bootjp Jul 13, 2026
fbd7f56
migration: raft-apply range imports
bootjp Jul 13, 2026
0fe341d
distribution: serve versioned ownership lookups
bootjp Jul 13, 2026
7d3b01c
migration: stage imported versions
bootjp Jul 13, 2026
e7f69ef
migration: merge staged visibility reads
bootjp Jul 13, 2026
b4f2477
migration: promote staged versions
bootjp Jul 13, 2026
2edefb8
migration: complete target promotion catalog state
bootjp Jul 13, 2026
2ea0d38
migration: add target readiness guard
bootjp Jul 13, 2026
b93f5b9
distribution: add split job management RPCs
bootjp Jul 13, 2026
c1afe27
store: persist promote applied index
bootjp Jul 13, 2026
ca1a050
migration: harden staged route writes
bootjp Jul 13, 2026
58c2128
distribution: page split job listing
bootjp Jul 13, 2026
e3079c0
store: fix migration export blockers
bootjp Jul 13, 2026
dd52db8
Gate migration promotion replay safety
bootjp Jul 13, 2026
5ba3e82
kv: enforce target readiness guards
bootjp Jul 13, 2026
b04d30f
store: keep promote replay ungated
bootjp Jul 13, 2026
a2ee041
Guard target readiness apply paths
bootjp Jul 13, 2026
643a5a6
store: fix migration export metadata handling
bootjp Jul 13, 2026
6ac85a1
Fix target readiness guard gaps
bootjp Jul 13, 2026
ce91b66
Guard manifest scan readiness routes
bootjp Jul 13, 2026
d33bf35
Harden staged migration visibility
bootjp Jul 13, 2026
f2bcd65
Fix target readiness route proofs
bootjp Jul 13, 2026
7cc0875
Merge cross-group migration base
bootjp Jul 13, 2026
f727f30
Enable split migration job creation
bootjp Jul 13, 2026
c72eace
Validate promote cursor before proposal
bootjp Jul 13, 2026
53d38e3
Harden split migration readiness checks
bootjp Jul 13, 2026
a624a21
Harden target readiness route proofs
bootjp Jul 13, 2026
fa769ac
Harden split migration start guards
bootjp Jul 13, 2026
4b8e345
Harden staged visibility OCC guards
bootjp Jul 13, 2026
4447dd3
Reject non-active migration sources
bootjp Jul 13, 2026
061275a
Bound migration export sparse scans
bootjp Jul 13, 2026
b5723b5
Fix target readiness staged delete paths
bootjp Jul 13, 2026
ca503cf
Fix staged migration read safety
bootjp Jul 13, 2026
f22c4ff
Stabilize urgent compactor pagination test
bootjp Jul 13, 2026
31ed507
Fix migration export metadata handling
bootjp Jul 13, 2026
2065650
Fix migration route bracket filtering
bootjp Jul 13, 2026
119dda2
Guard staged readiness probes
bootjp Jul 13, 2026
2b50129
Fix staged migration guard gaps
bootjp Jul 14, 2026
9f0b7f1
store: skip export writer registry rows
bootjp Jul 14, 2026
e6f88bd
Guard target readiness retry paths
bootjp Jul 14, 2026
83a3d99
Merge cross-group migration base
bootjp Jul 14, 2026
79b3409
Tighten target readiness proof
bootjp Jul 14, 2026
ce0c2e1
Wire split migration capability gate
bootjp Jul 14, 2026
f884ed2
Replicate target readiness guards
bootjp Jul 14, 2026
3cb3aad
store: advance promotion commit watermark
bootjp Jul 14, 2026
c5dae8f
Fix bounded store export edge cases
bootjp Jul 14, 2026
49adab3
Gate split migration on peer capability
bootjp Jul 14, 2026
c459255
Fix migration route edge cases
bootjp Jul 14, 2026
c64ff5f
migration: harden split capability gate
bootjp Jul 14, 2026
90fa654
store: honor raft promotion sync mode
bootjp Jul 14, 2026
a775fda
Trim Lua negative cache bound test
bootjp Jul 14, 2026
2b2cbd2
store: harden Pebble migration export
bootjp Jul 14, 2026
a9c9afe
Trim Lua negative cache bound test
bootjp Jul 14, 2026
67f3e9d
store: separate list delta key prefix
bootjp Jul 14, 2026
eebdbf2
Stabilize stream latency seed writes
bootjp Jul 14, 2026
c298060
Enable split migration capability gate
bootjp Jul 14, 2026
05b72cf
Bound migration export scan skips
bootjp Jul 14, 2026
5280996
Fix migration routing edge cases
bootjp Jul 14, 2026
a653226
Fence routed migration exports
bootjp Jul 14, 2026
e7467d5
Route legacy list deltas and stream scans
bootjp Jul 14, 2026
70fb630
Keep split migration capability fail closed
bootjp Jul 14, 2026
aa94570
Add migration bracket planner (#1086)
bootjp Jul 14, 2026
aff2a9f
migration: enable split migration job creation (#1093)
bootjp Jul 14, 2026
9f8dab7
Merge remote-tracking branch 'origin/design/hotspot-split-m2-wire' in…
bootjp Jul 14, 2026
6892dd9
Skip txn locks in migration export
bootjp Jul 14, 2026
fba7560
Filter legacy list deltas by value
bootjp Jul 14, 2026
9f80fa9
migration: tighten readiness guard publication
bootjp Jul 14, 2026
ccda608
migration: merge m2 wire and fix legacy deltas
bootjp Jul 14, 2026
585e5ac
Merge remote-tracking branch 'origin/design/hotspot-split-m2-store-ex…
bootjp Jul 14, 2026
727c1dd
Merge remote-tracking branch 'origin/design/hotspot-split-m2-fence-dr…
bootjp Jul 14, 2026
5acb33f
migration: fence S3 bucket metadata writes
bootjp Jul 14, 2026
db7a325
Merge remote-tracking branch 'origin/design/hotspot-split-m2-fence-dr…
bootjp Jul 14, 2026
14aa9e9
migration: guard read-only shard validation
bootjp Jul 14, 2026
5254a43
Merge cross-group migration base
bootjp Jul 14, 2026
295cdf6
migration: close cross-group floor gaps
bootjp Jul 14, 2026
b1381f3
migration: validate promotion resume state
bootjp Jul 14, 2026
dddda07
migration: validate staged readiness conflicts
bootjp Jul 14, 2026
84704f9
Close legacy list export gaps
bootjp Jul 14, 2026
650c450
migration: route staged s3 auxiliaries
bootjp Jul 14, 2026
b4b8664
Merge cross-group migration fixes
bootjp Jul 14, 2026
19a28e9
kv: keep floor checks out of raft replay
bootjp Jul 14, 2026
40c2f52
store: ignore stale promotion cursors
bootjp Jul 14, 2026
327cb1c
Merge remote-tracking branch 'origin/design/hotspot-split-m2-cross-gr…
bootjp Jul 14, 2026
f1aeeeb
Guard split fence retry paths
bootjp Jul 14, 2026
1db42a8
Merge remote-tracking branch 'origin/design/hotspot-split-m2-fence-dr…
bootjp Jul 14, 2026
4b1d950
Merge remote-tracking branch 'origin/design/hotspot-split-m2-cross-gr…
bootjp Jul 14, 2026
961a177
Merge remote-tracking branch 'origin/design/hotspot-split-m2-promote'…
bootjp Jul 14, 2026
9d9e051
Fix promotion completion merge compatibility
bootjp Jul 14, 2026
8e74d6e
Merge promotion completion into target readiness
bootjp Jul 14, 2026
2d57ae9
Fail closed read-only shard readiness checks
bootjp Jul 14, 2026
9994ec7
Merge remote-tracking branch 'origin/design/hotspot-split-m2-target-r…
bootjp Jul 14, 2026
8d11161
Guard snapshot spooling disk headroom
bootjp Jul 14, 2026
0215c38
Raise snapshot spool headroom reserve
bootjp Jul 14, 2026
774779b
Fix raft startup and snapshot spool guards
bootjp Jul 14, 2026
46722f6
Handle route fences in adapter retries
bootjp Jul 14, 2026
70cf6f8
Stabilize raft snapshot dispatch and maintenance gates
bootjp Jul 14, 2026
a241809
Honor partition resolver in fence prechecks
bootjp Jul 14, 2026
ee363c4
Make raft startup and fence checks deterministic
bootjp Jul 14, 2026
54ad999
Gate public startup after raft replay
bootjp Jul 14, 2026
e20475f
Stabilize raft dispatch and S3 cleanup fences
bootjp Jul 14, 2026
c3a8ecf
Stabilize raft startup and write fences
bootjp Jul 14, 2026
61feed4
Stabilize route-fence cleanup retries
bootjp Jul 14, 2026
f373cd3
Keep heartbeat responses coalescing under read-index load
bootjp Jul 14, 2026
8f821fe
Stabilize raft snapshot recovery checkpoints
bootjp Jul 14, 2026
3e988ef
Enforce current route fences at apply time
bootjp Jul 14, 2026
04e8053
Prioritize received raft snapshots
bootjp Jul 14, 2026
d818602
Stabilize snapshot catch-up and fence checks
bootjp Jul 14, 2026
db9843a
Stabilize startup reads and snapshot recovery
bootjp Jul 15, 2026
a6d7d76
Reset stale raft peer connections
bootjp Jul 15, 2026
c8f865e
Stabilize snapshot catch-up and fence routing
bootjp Jul 15, 2026
b553090
Keep SQS receive route fences visible
bootjp Jul 15, 2026
d2b22ac
Stabilize redis proxy and raft recovery
bootjp Jul 15, 2026
549a786
Add migration fence drain guards (#1087)
bootjp Jul 15, 2026
fefd809
Fix sparse Lua list pop fallback
bootjp Jul 16, 2026
e1b32b5
Align ElasticKV proxy socket timeouts
bootjp Jul 16, 2026
cf47858
Replay blocking zset pops deterministically
bootjp Jul 16, 2026
f697640
Cap ElasticKV secondary script concurrency
bootjp Jul 16, 2026
7e40979
Refresh ElasticKV leader after not-leader replies
bootjp Jul 16, 2026
3df9239
Avoid full stream rewrites for Lua XADD
bootjp Jul 16, 2026
d19b2a4
Handle Redis proxy replay edge cases
bootjp Jul 16, 2026
a850ee1
distribution: add split job management RPCs (#1092)
bootjp Jul 16, 2026
b8bb6ad
Honor Lua XADD cached state
bootjp Jul 16, 2026
adf5cdb
Cover Lua XADD maxlen zero append
bootjp Jul 16, 2026
01b02ad
Merge remote-tracking branch 'origin/main' into work/pr1088-conflict
bootjp Jul 16, 2026
f02d1c0
Merge remote-tracking branch 'origin/design/hotspot-split-m2-store-ex…
bootjp Jul 16, 2026
7407423
Enforce migration write floors during apply
bootjp Jul 16, 2026
2c920bc
Fence migration exports with applied reads
bootjp Jul 16, 2026
bdf05a6
Delete staged rows for migrated prefixes
bootjp Jul 16, 2026
c9d7ddc
Preserve staged scan visibility
bootjp Jul 16, 2026
bbc1876
migration: harden target readiness guards
bootjp Jul 16, 2026
a02b6c1
migration: add target readiness guard (#1091)
bootjp Jul 16, 2026
acc9305
migration: fix staged s3 bucket routing
bootjp Jul 16, 2026
f902f34
migration: persist promotion timestamp floor
bootjp Jul 16, 2026
e3e37ce
migration: keep invalid promote cursors non-halting
bootjp Jul 16, 2026
f6457de
migration: promote staged versions (#1089)
bootjp Jul 16, 2026
a3b7276
Merge remote-tracking branch 'origin/design/hotspot-split-m2-cross-gr…
bootjp Jul 16, 2026
f51ca32
Merge remote-tracking branch 'origin/design/hotspot-split-m2-promotio…
bootjp Jul 16, 2026
5dec8cb
migration: fix promotion complete review findings
bootjp Jul 16, 2026
00da46c
migration: make promotion completion retry idempotent
bootjp Jul 16, 2026
19961d3
migration: disambiguate legacy list exports
bootjp Jul 16, 2026
bd55701
store: preserve empty-key migration exports
bootjp Jul 16, 2026
8a687ca
migration: tighten routed cleanup safety
bootjp Jul 16, 2026
fd179f9
migration: preserve ambiguous list tombstones
bootjp Jul 16, 2026
c275022
adapter: preserve scan route groups for list cleanup
bootjp Jul 16, 2026
5cf01fd
adapter: backtrack list compactor accepted tails
bootjp Jul 16, 2026
7fe4112
kv: cover pinned primary transaction commits
bootjp Jul 16, 2026
83cd5bf
ci: generate redis proxy image metadata locally
bootjp Jul 16, 2026
0e76181
proxy: allow redis-only without secondary seeds
bootjp Jul 16, 2026
ab82447
ci: avoid failing lint on reviewdog API errors
bootjp Jul 16, 2026
5189172
ci: keep lint required when reviewdog is unavailable
bootjp Jul 16, 2026
319b442
migration: persist applied index on imports
bootjp Jul 16, 2026
2c29191
kv: preserve broad scans with staged routes
bootjp Jul 16, 2026
aaa96da
migration: gate import proposals during rollouts
bootjp Jul 17, 2026
d34d103
kv: reject raw prefix writes below floors
bootjp Jul 17, 2026
78a0400
migration: route partitioned exports by group
bootjp Jul 17, 2026
4851c6f
kv: scan staged routes for physical limits
bootjp Jul 17, 2026
052c64c
Merge cross-group migration fixes
bootjp Jul 17, 2026
2b7c3af
migration: keep split capability gated
bootjp Jul 17, 2026
eea5dec
adapter: pin sparse list fallback route bounds
bootjp Jul 17, 2026
1b70e39
redis: fix Lua XADD and blocking move replay
bootjp Jul 17, 2026
8df8acf
proxy: replay blocking multipop writes
bootjp Jul 17, 2026
9a0286d
proxy: avoid retrying user not-leader errors
bootjp Jul 17, 2026
7b12482
proxy: classify not-leader redis replies safely
bootjp Jul 17, 2026
5d78d7a
migration: reject armed readiness without write floor
bootjp Jul 17, 2026
178fda6
docs: mark hotspot m2 migration partial
bootjp Jul 17, 2026
c065ad9
migration: run split target promotion
bootjp Jul 17, 2026
252bd64
migration: harden split runner promotion
bootjp Jul 17, 2026
938a284
distribution: keep split migration gates closed
bootjp Jul 17, 2026
3e1c85f
Fix sparse Lua list pop fallback (#1094)
bootjp Jul 17, 2026
1c7ea5b
distribution: finish split promotion cleanup
bootjp Jul 17, 2026
07c4eea
Merge hotspot split M2 wire updates
bootjp Jul 18, 2026
3a7c3c2
Fix pinned migration cleanup routing
bootjp Jul 18, 2026
903e589
adapter: preserve stream metadata on Lua expiry
bootjp Jul 18, 2026
b3a3e7e
Merge remote-tracking branch 'origin/design/hotspot-split-m2-wire' in…
bootjp Jul 18, 2026
2c3d7d8
Merge remote-tracking branch 'origin/design/hotspot-split-m2-store-ex…
bootjp Jul 18, 2026
79405da
Merge remote-tracking branch 'origin/design/hotspot-split-m2-store-ex…
bootjp Jul 18, 2026
d9da086
raft: report cold-start replay gaps safely
bootjp Jul 18, 2026
5affdac
Merge remote-tracking branch 'origin/design/hotspot-split-m2-store-ex…
bootjp Jul 18, 2026
34706b0
Merge remote-tracking branch 'origin/design/hotspot-split-m2-cross-gr…
bootjp Jul 18, 2026
df4370b
Merge remote-tracking branch 'origin/design/hotspot-split-m2-promotio…
bootjp Jul 18, 2026
7c1f0aa
distribution: timestamp completed split jobs
bootjp Jul 18, 2026
737b930
kv: preserve migration fence bypasses
bootjp Jul 18, 2026
7cc8008
Merge hotspot split promotion updates
bootjp Jul 18, 2026
0c74b60
Complete split migration runner phases
bootjp Jul 18, 2026
634a75e
Complete hotspot split M2 migration lifecycle
bootjp Jul 18, 2026
1cbee6a
Track post-arm two-phase commits during split
bootjp Jul 18, 2026
8264ea3
Close split migration startup and retry gaps
bootjp Jul 18, 2026
28449f9
Complete hotspot split M2 migration lifecycle (#1096)
bootjp Jul 19, 2026
db933c0
docs: mark hotspot split M2 implemented
bootjp Jul 19, 2026
f270d8d
docs: link implemented hotspot split design
bootjp Jul 19, 2026
f2aad59
migration: complete target promotion catalog state
bootjp Jul 19, 2026
54cc7b9
migration: harden split cutover compatibility
bootjp Jul 19, 2026
48a4da8
migration: satisfy current lint gates
bootjp Jul 19, 2026
19c2657
Resolving merge conflicts with base branch
Copilot Jul 20, 2026
6b54f22
docs: resolve merge conflicts with base branch
Copilot Jul 20, 2026
20cebac
Keep hotspot M2 status PR docs-only
bootjp Jul 22, 2026
4ed488c
Mark hotspot split M2 design implemented (#1124)
bootjp Jul 23, 2026
ebf6878
proto: reserve raft admin status fields
bootjp Jul 23, 2026
f79bc00
migration: harden promotion readiness fences
bootjp Jul 23, 2026
4163665
adapter: fail closed on split capability regression
bootjp Jul 23, 2026
d57699d
migration: harden split cleanup guards
bootjp Jul 23, 2026
dcb7744
migration: preserve cleanup metadata barriers
bootjp Jul 24, 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
1,186 changes: 1,171 additions & 15 deletions adapter/distribution_server.go

Large diffs are not rendered by default.

1,662 changes: 1,621 additions & 41 deletions adapter/distribution_server_test.go

Large diffs are not rendered by default.

400 changes: 398 additions & 2 deletions adapter/internal.go

Large diffs are not rendered by default.

154 changes: 154 additions & 0 deletions adapter/internal_migration_probe_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
package adapter

import (
"context"
"testing"
"time"

"github.com/bootjp/elastickv/distribution"
"github.com/bootjp/elastickv/kv"
pb "github.com/bootjp/elastickv/proto"
"github.com/bootjp/elastickv/store"
"github.com/stretchr/testify/require"
)

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

ctx := context.Background()
st := store.NewMVCCStore()
engine := distribution.NewEngine()
require.NoError(t, engine.ApplySnapshot(distribution.CatalogSnapshot{
Version: 9,
Routes: []distribution.RouteDescriptor{{
RouteID: 2,
Start: []byte("m"),
GroupID: 2,
State: distribution.RouteStateActive,
MinWriteTSExclusive: 88,
}},
}))
tracker := kv.NewActiveTimestampTracker()
internal := NewInternalWithEngine(nil, nil, nil, nil,
WithInternalStore(st),
WithInternalRouteEngine(engine),
WithInternalActiveTimestampTracker(tracker),
)

state := store.TargetStagedReadinessState{
JobID: 7,
RouteStart: []byte("m"),
ExpectedCutoverVersion: 9,
MigrationJobID: 7,
MinWriteTSExclusive: 88,
Armed: true,
}
readinessWriter, ok := st.(store.MigrationTargetReadinessWriter)
require.True(t, ok)
require.NoError(t, readinessWriter.ApplyTargetStagedReadiness(ctx, state))
control, err := internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_CONTROL_APPLIED,
RouteStart: []byte("m"),
ExpectedCatalogVersion: 9,
MigrationJobId: 7,
MinWriteTsExclusive: 88,
})
require.NoError(t, err)
require.True(t, control.Ready)

cleared, err := internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_TARGET_DESCRIPTOR_CLEARED,
RouteStart: []byte("m"),
ExpectedCatalogVersion: 9,
ExpectedGroupId: 2,
MinWriteTsExclusive: 88,
})
require.NoError(t, err)
require.True(t, cleared.Ready)

pin := tracker.Pin(100)
drained, err := internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_SOURCE_READ_DRAINED,
ReadDrainNotBeforeMs: time.Now().Add(-time.Second).UnixMilli(),
})
require.NoError(t, err)
require.False(t, drained.Ready)
pin.Release()
drained, err = internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_SOURCE_READ_DRAINED,
ReadDrainNotBeforeMs: time.Now().Add(-time.Second).UnixMilli(),
})
require.NoError(t, err)
require.True(t, drained.Ready)

metadata, err := internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_METADATA_CLEARED,
})
require.NoError(t, err)
require.False(t, metadata.Ready)
cleaner, ok := st.(store.MigrationCleaner)
require.True(t, ok)
require.NoError(t, cleaner.ClearMigrationState(ctx, 7, 0))
metadata, err = internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_METADATA_CLEARED,
})
require.NoError(t, err)
require.True(t, metadata.Ready)
}

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

ctx := context.Background()
st := store.NewMVCCStore()
internal := NewInternalWithEngine(nil, nil, nil, nil, WithInternalStore(st))
_, err := st.ImportVersions(ctx, store.ImportVersionsOptions{
JobID: 7,
BracketID: 1,
BatchSeq: 1,
Cursor: []byte("cursor-1"),
Versions: []store.MVCCVersion{{
Key: []byte("m/key"),
Value: []byte("value"),
CommitTS: 42,
}},
})
require.NoError(t, err)

metadata, err := internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_METADATA_CLEARED,
})
require.NoError(t, err)
require.False(t, metadata.Ready)

cleaner, ok := st.(store.MigrationCleaner)
require.True(t, ok)
require.NoError(t, cleaner.ClearMigrationState(ctx, 7, 0))
metadata, err = internal.ProbeMigrationState(ctx, &pb.ProbeMigrationStateRequest{
JobId: 7,
Kind: pb.MigrationStateProbeKind_MIGRATION_STATE_PROBE_KIND_METADATA_CLEARED,
})
require.NoError(t, err)
require.True(t, metadata.Ready)
}

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

ctx := context.Background()
st := store.NewMVCCStore()
require.NoError(t, st.ApplyMutations(ctx, []*store.KVPairMutation{{Key: []byte("m"), Value: []byte("value")}}, nil, 50, 50))
internal := NewInternalWithEngine(nil, mockInternalLeader{}, nil, nil, WithInternalStore(st))

resp, err := internal.IssueMigrationTimestamp(ctx, &pb.IssueMigrationTimestampRequest{})
require.NoError(t, err)
require.Equal(t, uint64(50), resp.GetLastCommitTs())
require.Greater(t, resp.GetTimestamp(), resp.GetLastCommitTs())
}
134 changes: 134 additions & 0 deletions adapter/internal_migration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -659,3 +659,137 @@ func TestInternalPromoteStagedVersionsAppliesStoreBatch(t *testing.T) {
require.ErrorIs(t, err, store.ErrKeyNotFound)
require.GreaterOrEqual(t, clock.Current(), uint64(30))
}

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

ctx := context.Background()
st := store.NewMVCCStore()
proposer := &applyingMigrationProposer{
fsm: kv.NewKvFSMWithHLC(st, nil),
}
internal := NewInternalWithEngine(nil, mockInternalLeader{}, nil, nil,
WithInternalStore(st),
WithInternalMigrationProposer(proposer),
WithInternalMigrationPromoteGate(func(context.Context) error {
return status.Error(codes.FailedPrecondition, "migration opcode disabled for test")
}),
)

resp, err := internal.ApplyTargetStagedReadiness(ctx, &pb.TargetStagedReadinessRequest{
JobId: 9,
RouteStart: []byte("a"),
RouteEnd: []byte("z"),
ExpectedCutoverVersion: 3,
MigrationJobId: 7,
MinWriteTsExclusive: 100,
Armed: true,
})
require.NoError(t, err)
require.NotNil(t, resp)
require.Equal(t, uint64(1), proposer.calls)

reader, ok := st.(store.MigrationTargetReadinessReader)
require.True(t, ok)
states, err := reader.MigrationTargetReadinessStates(ctx)
require.NoError(t, err)
require.Equal(t, []store.TargetStagedReadinessState{{
JobID: 9,
RouteStart: []byte("a"),
RouteEnd: []byte("z"),
ExpectedCutoverVersion: 3,
MigrationJobID: 7,
MinWriteTSExclusive: 100,
Armed: true,
}}, states)
}

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

ctx := context.Background()
st := store.NewMVCCStore()
proposer := &applyingMigrationProposer{
fsm: kv.NewKvFSMWithHLC(st, nil),
}
internal := NewInternalWithEngine(nil, mockInternalLeader{}, nil, nil,
WithInternalStore(st),
WithInternalMigrationProposer(proposer),
WithInternalMigrationPromoteGate(func(context.Context) error { return nil }),
)

_, err := internal.ApplyTargetStagedReadiness(ctx, &pb.TargetStagedReadinessRequest{
JobId: 9,
RouteStart: []byte("a"),
RouteEnd: []byte("z"),
ExpectedCutoverVersion: 3,
MigrationJobId: 7,
MinWriteTsExclusive: 100,
Armed: true,
})
require.NoError(t, err)
require.Equal(t, uint64(1), proposer.calls)

reader, ok := st.(store.MigrationTargetReadinessReader)
require.True(t, ok)
states, err := reader.MigrationTargetReadinessStates(ctx)
require.NoError(t, err)
require.Equal(t, []store.TargetStagedReadinessState{{
JobID: 9,
RouteStart: []byte("a"),
RouteEnd: []byte("z"),
ExpectedCutoverVersion: 3,
MigrationJobID: 7,
MinWriteTSExclusive: 100,
Armed: true,
}}, states)
}

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

ctx := context.Background()
st := store.NewMVCCStore()
proposer := &applyingMigrationProposer{
fsm: kv.NewKvFSMWithHLC(st, nil),
}
internal := NewInternalWithEngine(nil, mockInternalLeader{}, nil, nil,
WithInternalStore(st),
WithInternalMigrationProposer(proposer),
WithInternalMigrationPromoteGate(func(context.Context) error { return nil }),
)

resp, err := internal.ApplyTargetStagedReadiness(ctx, &pb.TargetStagedReadinessRequest{
JobId: 9,
RouteStart: []byte("a"),
RouteEnd: []byte("z"),
ExpectedCutoverVersion: 3,
MigrationJobId: 7,
Armed: true,
})
require.Nil(t, resp)
require.Error(t, err)
require.Equal(t, codes.InvalidArgument, status.Code(err))
require.Equal(t, uint64(0), proposer.calls)
}

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

proposer := &applyingMigrationProposer{fsm: kv.NewKvFSMWithHLC(store.NewMVCCStore(), nil)}
internal := NewInternalWithEngine(nil, mockInternalLeader{}, nil, nil,
WithInternalMigrationProposer(proposer),
WithInternalMigrationCleanupGate(func(context.Context) error {
return status.Error(codes.FailedPrecondition, "migration cleanup disabled for test")
}),
)

resp, err := internal.CleanupMigration(context.Background(), &pb.CleanupMigrationRequest{
JobId: 9,
Mode: pb.MigrationCleanupMode_MIGRATION_CLEANUP_MODE_VERSIONS,
KeyFamily: distribution.MigrationFamilyUser,
})
require.Nil(t, resp)
require.Equal(t, codes.FailedPrecondition, status.Code(err))
require.Zero(t, proposer.calls)
}
24 changes: 16 additions & 8 deletions adapter/redis_compat_commands_stream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -279,19 +279,27 @@ func TestRedis_StreamXReadLatencyIsConstant(t *testing.T) {

rdb := redis.NewClient(&redis.Options{Addr: nodes[0].redisAddress})
defer func() { _ = rdb.Close() }()
ctx := context.Background()
ctx := t.Context()

const probes = 100
lastID := seedStreamEntriesForXReadLatency(t, nodes[0].redisServer, ctx, "stream-lat")

measure := func() time.Duration {
start := time.Now()
streams, err := rdb.XRead(ctx, &redis.XReadArgs{
Streams: []string{"stream-lat", lastID},
Count: 10,
Block: 10 * time.Millisecond,
}).Result()
elapsed := time.Since(start)
var (
streams []redis.XStream
elapsed time.Duration
)
err := retryNotLeader(ctx, func() error {
start := time.Now()
var xerr error
streams, xerr = rdb.XRead(ctx, &redis.XReadArgs{
Streams: []string{"stream-lat", lastID},
Count: 10,
Block: 10 * time.Millisecond,
}).Result()
elapsed = time.Since(start)
return xerr
})
require.True(t, errors.Is(err, redis.Nil) || err == nil)
require.Empty(t, streams)
return elapsed
Expand Down
7 changes: 5 additions & 2 deletions adapter/redis_delta_compactor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1684,8 +1684,11 @@ func TestDeltaCompactor_UrgentCompactionPagination(t *testing.T) {
require.NoError(t, st.PutAt(ctx, dKey, delta, i+1, 0))
}

// Queue and process the urgent compaction.
c.TriggerUrgentCompaction("hash", userKey)
// Exercise the urgent pagination loop directly. Running through c.Run would
// also start an initial background SyncOnce; the local test coordinator applies
// elems one-by-one instead of atomically, so a test read can observe the meta
// update before all delete elems have been applied under the race detector.
c.compactUrgentKey(ctx, urgentCompactionRequest{typeName: "hash", userKey: userKey})

runCtx, cancel := context.WithCancel(ctx)
defer cancel()
Expand Down
Loading
Loading