diff --git a/adapter/admin_grpc.go b/adapter/admin_grpc.go index 0eb0dd028..a8cb10a8a 100644 --- a/adapter/admin_grpc.go +++ b/adapter/admin_grpc.go @@ -699,6 +699,7 @@ func newKeyVizRowFrom(mr keyviz.MatrixRow, numCols int) *pb.KeyVizRow { } row := &pb.KeyVizRow{ BucketId: bucketIDFor(mr), + Label: string(mr.Label), Start: append([]byte(nil), mr.Start...), End: append([]byte(nil), mr.End...), Aggregate: mr.Aggregate, @@ -722,6 +723,9 @@ func bucketIDFor(mr keyviz.MatrixRow) string { return "virtual:" + strconv.FormatUint(mr.RouteID, 10) } id := "route:" + strconv.FormatUint(mr.RouteID, 10) + if mr.Label != keyviz.LabelLegacy { + id += ":" + string(mr.Label) + } // Sub-bucket suffix only for genuinely sub-divided routes, so K=1 / // aggregate / degenerate slots keep the exact legacy id. See §5.1. if mr.SubBucketCount > 1 { diff --git a/adapter/admin_grpc_subrange_test.go b/adapter/admin_grpc_subrange_test.go index 2c7c4edc1..8847ac1bf 100644 --- a/adapter/admin_grpc_subrange_test.go +++ b/adapter/admin_grpc_subrange_test.go @@ -37,6 +37,33 @@ func TestMatrixToProtoSubRowsDoNotCollide(t *testing.T) { require.Equal(t, []uint64{9}, byID["route:1#1"]) } +func TestMatrixToProtoLabelsDoNotCollide(t *testing.T) { + t.Parallel() + pick := func(r keyviz.MatrixRow) uint64 { return r.Writes } + cols := []keyviz.MatrixColumn{ + { + At: time.Unix(1_700_000_000, 0), + Rows: []keyviz.MatrixRow{ + {RouteID: 1, Label: keyviz.LabelDynamo, Start: []byte("a"), End: []byte("z"), SubBucketCount: 1, Writes: 5}, + {RouteID: 1, Label: keyviz.LabelRedis, Start: []byte("a"), End: []byte("z"), SubBucketCount: 1, Writes: 9}, + }, + }, + } + + resp := matrixToProto(cols, pick, 0) + require.Len(t, resp.Rows, 2) + byID := map[string]uint64{} + labels := map[string]string{} + for _, row := range resp.Rows { + byID[row.BucketId] = row.Values[0] + labels[row.BucketId] = row.Label + } + require.Equal(t, uint64(5), byID["route:1:dynamo"]) + require.Equal(t, "dynamo", labels["route:1:dynamo"]) + require.Equal(t, uint64(9), byID["route:1:redis"]) + require.Equal(t, "redis", labels["route:1:redis"]) +} + // TestNewKeyVizRowFromAggregateZeroTotalFallback pins the adapter/handler // harmonization (Claude round-2 follow-up): an aggregate row whose // MemberRoutesTotal is still 0 (a just-coalesced bucket serialized before diff --git a/adapter/dynamodb.go b/adapter/dynamodb.go index f53ed238b..742e89d23 100644 --- a/adapter/dynamodb.go +++ b/adapter/dynamodb.go @@ -12,6 +12,7 @@ import ( "sync" "time" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" "github.com/bootjp/elastickv/monitoring" "github.com/bootjp/elastickv/store" @@ -259,7 +260,7 @@ func NewDynamoDBServer(listen net.Listener, st store.MVCCStore, coordinate kv.Co d := &DynamoDBServer{ listen: listen, store: st, - coordinator: coordinate, + coordinator: kv.WithKeyVizLabel(coordinate, keyviz.LabelDynamo), onePhaseTxnDedup: os.Getenv("ELASTICKV_DYNAMODB_ONEPHASE_DEDUP") != "0", } d.targetHandlers = map[string]func(http.ResponseWriter, *http.Request){ diff --git a/adapter/grpc.go b/adapter/grpc.go index cd966afad..49aaa58cc 100644 --- a/adapter/grpc.go +++ b/adapter/grpc.go @@ -7,6 +7,7 @@ import ( "sync" "github.com/bootjp/elastickv/internal" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" pb "github.com/bootjp/elastickv/proto" "github.com/bootjp/elastickv/store" @@ -47,7 +48,7 @@ func NewGRPCServer(store store.MVCCStore, coordinate kv.Coordinator, opts ...GRP Level: slog.LevelWarn, })), grpcTranscoder: newGrpcGrpcTranscoder(), - coordinator: coordinate, + coordinator: kv.WithKeyVizLabel(coordinate, keyviz.LabelRawKV), store: store, } for _, opt := range opts { diff --git a/adapter/redis.go b/adapter/redis.go index 73eb0e5d1..bb243ca80 100644 --- a/adapter/redis.go +++ b/adapter/redis.go @@ -16,6 +16,7 @@ import ( "time" "github.com/bootjp/elastickv/internal/raftengine" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" "github.com/bootjp/elastickv/monitoring" pb "github.com/bootjp/elastickv/proto" @@ -476,7 +477,7 @@ func NewRedisServer(listen net.Listener, redisAddr string, store store.MVCCStore r := &RedisServer{ listen: listen, store: store, - coordinator: coordinate, + coordinator: kv.WithKeyVizLabel(coordinate, keyviz.LabelRedis), redisTranscoder: newRedisTranscoder(), redisAddr: redisAddr, relay: relay, diff --git a/adapter/redis_delta_compactor.go b/adapter/redis_delta_compactor.go index 50159edfb..0f64a197d 100644 --- a/adapter/redis_delta_compactor.go +++ b/adapter/redis_delta_compactor.go @@ -7,6 +7,7 @@ import ( "sync" "time" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" "github.com/bootjp/elastickv/store" "github.com/cockroachdb/errors" @@ -113,7 +114,7 @@ func WithDeltaCompactorLogger(l *slog.Logger) DeltaCompactorOption { func NewDeltaCompactor(st store.MVCCStore, coord kv.Coordinator, opts ...DeltaCompactorOption) *DeltaCompactor { c := &DeltaCompactor{ st: st, - coord: coord, + coord: kv.WithKeyVizLabel(coord, keyviz.LabelRedis), logger: slog.Default(), maxCount: defaultDeltaCompactorMaxDeltaCount, interval: defaultDeltaCompactorScanInterval, diff --git a/adapter/redis_delta_compactor_test.go b/adapter/redis_delta_compactor_test.go index bf357630b..326bfff71 100644 --- a/adapter/redis_delta_compactor_test.go +++ b/adapter/redis_delta_compactor_test.go @@ -5,6 +5,7 @@ import ( "testing" "time" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" "github.com/bootjp/elastickv/store" "github.com/stretchr/testify/require" @@ -20,6 +21,27 @@ func newDeltaCompactorTestFixture(t *testing.T) (store.MVCCStore, *DeltaCompacto return st, c } +type recordingCompactorCoordinator struct { + stubAdapterCoordinator + labels []keyviz.Label +} + +func (c *recordingCompactorCoordinator) Dispatch(_ context.Context, reqs *kv.OperationGroup[kv.OP]) (*kv.CoordinateResponse, error) { + if reqs != nil { + c.labels = append(c.labels, reqs.KeyVizLabel) + } + return &kv.CoordinateResponse{}, nil +} + +func TestDeltaCompactorStampsRedisKeyVizLabel(t *testing.T) { + t.Parallel() + rec := &recordingCompactorCoordinator{} + c := NewDeltaCompactor(store.NewMVCCStore(), rec) + _, err := c.coord.Dispatch(context.Background(), &kv.OperationGroup[kv.OP]{}) + require.NoError(t, err) + require.Equal(t, []keyviz.Label{keyviz.LabelRedis}, rec.labels) +} + func TestDeltaCompactor_ListDeltaFoldedIntoBaseMeta(t *testing.T) { t.Parallel() diff --git a/adapter/s3.go b/adapter/s3.go index db200417b..2fc5dec2d 100644 --- a/adapter/s3.go +++ b/adapter/s3.go @@ -24,6 +24,7 @@ import ( "time" "github.com/bootjp/elastickv/internal/s3keys" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" "github.com/bootjp/elastickv/store" "github.com/cockroachdb/errors" @@ -318,7 +319,7 @@ func NewS3Server(listen net.Listener, s3Addr string, st store.MVCCStore, coordin s3Addr: s3Addr, region: s3DefaultRegion, store: st, - coordinator: coordinate, + coordinator: kv.WithKeyVizLabel(coordinate, keyviz.LabelS3), leaderS3: cloneLeaderAddrMap(leaderS3), cleanupSem: make(chan struct{}, s3ManifestCleanupWorkers), } diff --git a/adapter/sqs.go b/adapter/sqs.go index b7da6ddfe..4a8bdd93c 100644 --- a/adapter/sqs.go +++ b/adapter/sqs.go @@ -10,6 +10,7 @@ import ( "sync" "time" + "github.com/bootjp/elastickv/keyviz" "github.com/bootjp/elastickv/kv" "github.com/bootjp/elastickv/store" "github.com/cockroachdb/errors" @@ -285,7 +286,7 @@ func NewSQSServer(listen net.Listener, st store.MVCCStore, coordinate kv.Coordin s := &SQSServer{ listen: listen, store: st, - coordinator: coordinate, + coordinator: kv.WithKeyVizLabel(coordinate, keyviz.LabelSQS), reaperCtx: reaperCtx, reaperCancel: reaperCancel, throttle: newBucketStoreDefault(), diff --git a/docs/admin_ui_key_visualizer_design.md b/docs/admin_ui_key_visualizer_design.md index 2f296a6a1..1c0794fa3 100644 --- a/docs/admin_ui_key_visualizer_design.md +++ b/docs/admin_ui_key_visualizer_design.md @@ -135,7 +135,7 @@ Response matrix format: `matrix[i][j]` is the value for bucket `i` at time colum A single call site is added at the dispatch entry of `kv.ShardedCoordinator` (see `kv/sharded_coordinator.go`), immediately after the request is resolved to a `RouteID`: ```go -sampler.Observe(routeID, op, keyLen, valueLen) +sampler.Observe(routeID, key, op, valueLen, label) ``` `sampler` is an interface; the default implementation is nil-safe (a nil sampler compiles to one branch and no allocation). The hook runs *before* Raft proposal so it measures offered load, not applied load. diff --git a/docs/design/2026_04_28_proposed_keyviz_adapter_labels.md b/docs/design/2026_04_28_implemented_keyviz_adapter_labels.md similarity index 91% rename from docs/design/2026_04_28_proposed_keyviz_adapter_labels.md rename to docs/design/2026_04_28_implemented_keyviz_adapter_labels.md index 88ad5e4ef..831f8cd18 100644 --- a/docs/design/2026_04_28_proposed_keyviz_adapter_labels.md +++ b/docs/design/2026_04_28_implemented_keyviz_adapter_labels.md @@ -1,13 +1,21 @@ --- -status: proposed +status: implemented phase: 2-C+ parent_design: docs/admin_ui_key_visualizer_design.md author: bootjp -date: 2026-04-28 +date: 2026-07-08 --- # KeyViz adapter / namespace labels +Implementation note: the shipped surface uses typed labels in +`keyviz/labels.go`, a `--keyvizLabelsEnabled` startup flag, adapter +coordinator wrappers for RawKV / Redis / DynamoDB / S3 / SQS, composite +`route::