Skip to content
File

Blob: kv/aof/mutation.go

go112 lines
1package aof
2 
3import (
4 "context"
5 "io/fs"
6 "sync"
7 
8 "go.miragespace.co/specter/kv/aof/proto"
9 "go.miragespace.co/specter/spec/protocol"
10)
11 
12var reqPool = sync.Pool{
13 New: func() any {
14 return &mutationReq{
15 err: make(chan error),
16 mut: proto.MutationFromVTPool(),
17 }
18 },
19}
20 
21func (d *DiskKV) handleMutation(mut *proto.Mutation) error {
22 var err error
23 
24 switch mut.GetType() {
25 case proto.MutationType_SIMPLE_PUT:
26 err = d.memKv.Put(context.Background(), mut.GetKey(), mut.GetValue())
27 
28 case proto.MutationType_SIMPLE_DELETE:
29 err = d.memKv.Delete(context.Background(), mut.GetKey())
30 
31 case proto.MutationType_PREFIX_APPEND:
32 err = d.memKv.PrefixAppend(context.Background(), mut.GetKey(), mut.GetValue())
33 
34 case proto.MutationType_PREFIX_REMOVE:
35 err = d.memKv.PrefixRemove(context.Background(), mut.GetKey(), mut.GetValue())
36 
37 case proto.MutationType_IMPORT:
38 err = d.memKv.Import(context.Background(), mut.GetKeys(), mut.GetValues())
39 
40 case proto.MutationType_REMOVE_KEYS:
41 err = d.memKv.RemoveKeys(context.Background(), mut.GetKeys())
42 
43 }
44 return err
45}
46 
47func (d *DiskKV) mutationHandler(fn func(mut *proto.Mutation)) error {
48 d.writeBarrier.RLock()
49 defer d.writeBarrier.RUnlock()
50 if d.closed.Load() {
51 return fs.ErrClosed
52 }
53 
54 req := reqPool.Get().(*mutationReq)
55 defer func() {
56 req.mut.ResetVT()
57 reqPool.Put(req)
58 }()
59 
60 fn(req.mut)
61 
62 d.queue <- req
63 return <-req.err
64}
65 
66func (d *DiskKV) Put(ctx context.Context, key []byte, value []byte) error {
67 return d.mutationHandler(func(mut *proto.Mutation) {
68 mut.Type = proto.MutationType_SIMPLE_PUT
69 mut.Key = key
70 mut.Value = value
71 })
72}
73 
74func (d *DiskKV) Delete(ctx context.Context, key []byte) error {
75 return d.mutationHandler(func(mut *proto.Mutation) {
76 mut.Type = proto.MutationType_SIMPLE_DELETE
77 mut.Key = key
78 })
79}
80 
81func (d *DiskKV) PrefixAppend(ctx context.Context, prefix []byte, child []byte) error {
82 return d.mutationHandler(func(mut *proto.Mutation) {
83 mut.Type = proto.MutationType_PREFIX_APPEND
84 mut.Key = prefix
85 mut.Value = child
86 })
87}
88 
89func (d *DiskKV) PrefixRemove(ctx context.Context, prefix []byte, child []byte) error {
90 return d.mutationHandler(func(mut *proto.Mutation) {
91 mut.Type = proto.MutationType_PREFIX_REMOVE
92 mut.Key = prefix
93 mut.Value = child
94 })
95}
96 
97func (d *DiskKV) Import(ctx context.Context, keys [][]byte, values []*protocol.KVTransfer) error {
98 return d.mutationHandler(func(mut *proto.Mutation) {
99 mut.Type = proto.MutationType_IMPORT
100 mut.Keys = keys
101 mut.Values = values
102 })
103}
104 
105func (d *DiskKV) RemoveKeys(ctx context.Context, keys [][]byte) error {
106 d.mutationHandler(func(mut *proto.Mutation) {
107 mut.Type = proto.MutationType_REMOVE_KEYS
108 mut.Keys = keys
109 })
110 return nil
111}