Skip to content

Commit e4d91ae

Browse files
kvstorage/wag: prepare for WAG truncation
This commit sets up the premitives for WAG truncations by: 1) Adding a `Delete` function for removing WAG nodes by index. 2) Changing the WAG `Iterator.Iter` method to return `iter.Seq2[uint64, wagpb.Node]` pair. This will be useful when we iterate over the WAG and truncate the nodes that have been applied and synced. Epic: none Release note: None Co-Authored-By: roachdev-claude <roachdev-claude-bot@cockroachlabs.com>
1 parent 6cea22f commit e4d91ae

2 files changed

Lines changed: 68 additions & 13 deletions

File tree

pkg/kv/kvserver/kvstorage/wag/store.go

Lines changed: 17 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -76,11 +76,16 @@ func Write(w storage.Writer, index uint64, node wagpb.Node) error {
7676
return w.PutUnversioned(keys.StoreWAGNodeKey(index), data)
7777
}
7878

79+
// Delete removes the WAG node at the given sequence number.
80+
func Delete(w storage.Writer, index uint64) error {
81+
return w.ClearUnversioned(keys.StoreWAGNodeKey(index), storage.ClearOptions{})
82+
}
83+
7984
// Iterator helps to scan the WAG sequence.
8085
//
8186
// var iter wag.Iterator
82-
// for node := range iter.Iter(ctx, reader) {
83-
// // process node
87+
// for index, node := range iter.Iter(ctx, reader) {
88+
// // process index, node
8489
// }
8590
// if err := iter.Error(); err != nil {
8691
// return err
@@ -92,8 +97,9 @@ type Iterator struct {
9297
err error
9398
}
9499

95-
// Iter returns an iterator that scans the WAG sequence.
96-
func (it *Iterator) Iter(ctx context.Context, r storage.Reader) iter.Seq[wagpb.Node] {
100+
// Iter returns an iterator that scans the WAG sequence. The iterator yields a
101+
// pair containing the WAG node index and the WAG node itself.
102+
func (it *Iterator) Iter(ctx context.Context, r storage.Reader) iter.Seq2[uint64, wagpb.Node] {
97103
prefix := keys.StoreWAGPrefix()
98104
mi, err := r.NewMVCCIterator(ctx, storage.MVCCKeyIterKind, storage.IterOptions{
99105
UpperBound: prefix.PrefixEnd(),
@@ -104,13 +110,18 @@ func (it *Iterator) Iter(ctx context.Context, r storage.Reader) iter.Seq[wagpb.N
104110
}
105111
mi.SeekGE(storage.MakeMVCCMetadataKey(prefix))
106112

107-
return func(yield func(wagpb.Node) bool) {
113+
return func(yield func(uint64, wagpb.Node) bool) {
108114
defer mi.Close()
109115
for ; ; mi.Next() {
110116
if ok, err := mi.Valid(); err != nil || !ok {
111117
it.err = err
112118
return
113119
}
120+
index, err := keys.DecodeWAGNodeKey(mi.UnsafeKey().Key)
121+
if err != nil {
122+
it.err = err
123+
return
124+
}
114125
v, err := mi.UnsafeValue()
115126
if err != nil {
116127
it.err = err
@@ -120,7 +131,7 @@ func (it *Iterator) Iter(ctx context.Context, r storage.Reader) iter.Seq[wagpb.N
120131
if it.err = node.Unmarshal(v); it.err != nil { // nolint:protounmarshal
121132
return
122133
}
123-
if !yield(node) {
134+
if !yield(index, node) {
124135
return
125136
}
126137
}

pkg/kv/kvserver/kvstorage/wag/store_test.go

Lines changed: 51 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -55,16 +55,60 @@ func TestWrite(t *testing.T) {
5555
out = strings.ReplaceAll(out, "\n\n", "\n")
5656
echotest.Require(t, out, filepath.Join("testdata", t.Name()+".txt"))
5757

58-
// Smoke check that the iterator works.
59-
var iter Iterator
60-
count := 0
61-
for range iter.Iter(context.Background(), s.eng) {
62-
count++
58+
// readIndices returns the WAG node indices by scanning the engine.
59+
readIndices := func() []uint64 {
60+
var it Iterator
61+
var res []uint64
62+
for index := range it.Iter(context.Background(), s.eng) {
63+
res = append(res, index)
64+
}
65+
require.NoError(t, it.Error())
66+
return res
6367
}
64-
require.NoError(t, iter.Error())
68+
69+
// Verify that the iterator returns nodes with the correct indices.
6570
// 3 WAG nodes: create, init, split. The split is a single node with two
6671
// events (Split + Init) rather than two separate nodes (dep + event).
67-
require.Equal(t, 3, count)
72+
require.Equal(t, []uint64{1, 2, 3}, readIndices())
73+
}
74+
75+
func TestDelete(t *testing.T) {
76+
defer leaktest.AfterTest(t)()
77+
defer log.Scope(t).Close(t)
78+
79+
eng := storage.NewDefaultInMemForTesting()
80+
defer eng.Close()
81+
82+
id := roachpb.FullReplicaID{RangeID: 1, ReplicaID: 1}
83+
node := wagpb.Node{
84+
Events: []wagpb.Event{
85+
{Addr: wagpb.MakeAddr(id, 0), Type: wagpb.EventApply},
86+
},
87+
}
88+
89+
// Write 5 WAG nodes with indices 1 through 5.
90+
for i := uint64(1); i <= 5; i++ {
91+
require.NoError(t, Write(eng, i, node))
92+
}
93+
94+
// Read back all indices.
95+
readIndices := func() []uint64 {
96+
var it Iterator
97+
var res []uint64
98+
for index := range it.Iter(context.Background(), eng) {
99+
res = append(res, index)
100+
}
101+
require.NoError(t, it.Error())
102+
return res
103+
}
104+
require.Equal(t, []uint64{1, 2, 3, 4, 5}, readIndices())
105+
106+
// Delete nodes 2 and 4.
107+
require.NoError(t, Delete(eng, 2))
108+
require.NoError(t, Delete(eng, 4))
109+
110+
// Verify that only nodes 1, 3, 5 remain.
111+
require.Equal(t, []uint64{1, 3, 5}, readIndices())
68112
}
69113

70114
type store struct {

0 commit comments

Comments
 (0)