Skip to content
File

Blob: kv/aof/kv.go

go167 lines
1package aof
2 
3import (
4 "fmt"
5 "io/fs"
6 "path/filepath"
7 "sync"
8 "time"
9 
10 "go.miragespace.co/specter/kv/aof/proto"
11 "go.miragespace.co/specter/kv/memory"
12 "go.miragespace.co/specter/spec/chord"
13 
14 "github.com/tidwall/wal"
15 "go.uber.org/atomic"
16 "go.uber.org/zap"
17)
18 
19const (
20 LogDir = "wal"
21)
22 
23type DiskKV struct {
24 writeBarrier sync.RWMutex
25 logger *zap.Logger
26 memKv *memory.MemoryKV
27 queue chan *mutationReq
28 log *wal.Log
29 closeCh chan struct{}
30 closeWg sync.WaitGroup
31 closed *atomic.Bool
32 cfg Config
33 counter uint64
34 flushInterval time.Duration
35}
36 
37type Config struct {
38 Logger *zap.Logger
39 HasnFn chord.HashFn
40 DataDir string
41 FlushInterval time.Duration
42}
43 
44func (c Config) validate() error {
45 if c.Logger == nil {
46 return fmt.Errorf("nil Logger is invalid")
47 }
48 if c.HasnFn == nil {
49 return fmt.Errorf("nil HashFn is invalid")
50 }
51 if c.DataDir == "" {
52 return fmt.Errorf("empty DataDir is invalid")
53 }
54 if c.FlushInterval <= 0 {
55 return fmt.Errorf("non-positive FlushInterval is invalid")
56 }
57 return nil
58}
59 
60type mutationReq struct {
61 mut *proto.Mutation
62 err chan error
63}
64 
65func logPath(dir string) string {
66 return filepath.Join(dir, LogDir)
67}
68 
69func New(cfg Config) (*DiskKV, error) {
70 if err := cfg.validate(); err != nil {
71 return nil, err
72 }
73 // store log to wal/ subdirectory to support future snapshot
74 l, err := wal.Open(logPath(cfg.DataDir), &wal.Options{
75 SegmentSize: 2 * 1024 * 1024, // 2MB
76 SegmentCacheSize: 4, // 8MB
77 LogFormat: wal.Binary,
78 NoSync: true,
79 NoCopy: true,
80 })
81 if err != nil {
82 return nil, fmt.Errorf("error opening log: %w", err)
83 }
84 d := &DiskKV{
85 logger: cfg.Logger,
86 memKv: memory.WithHashFn(cfg.HasnFn),
87 queue: make(chan *mutationReq),
88 log: l,
89 closeCh: make(chan struct{}),
90 closed: atomic.NewBool(false),
91 flushInterval: cfg.FlushInterval,
92 cfg: cfg,
93 }
94 d.logger.Info("Using append only log for kv storage", zap.String("dir", cfg.DataDir))
95 
96 if err := d.replayLogs(); err != nil {
97 return nil, err
98 }
99 
100 d.closeWg.Add(1)
101 
102 return d, nil
103}
104 
105func (d *DiskKV) Start() {
106 ticker := time.NewTicker(d.flushInterval)
107 defer ticker.Stop()
108 
109 defer d.closeWg.Done()
110 
111 d.logger.Info("Periodically flushing logs to disk", zap.Duration("interval", d.flushInterval))
112 
113 dirty := false
114 for {
115 select {
116 case <-d.closeCh:
117 return
118 case <-ticker.C:
119 if !dirty {
120 continue
121 }
122 dirty = false
123 if err := d.log.Sync(); err != nil {
124 d.logger.Error("Error flushing logs periodically", zap.Error(err))
125 }
126 case m := <-d.queue:
127 var mutError error
128 if logError := d.appendLog(m.mut); logError == nil {
129 mutError = d.handleMutation(m.mut)
130 if mutError != nil {
131 d.rollbackOne(m.mut, mutError)
132 }
133 } else {
134 d.logger.Error("Error appending mutation log",
135 zap.String("mutation", m.mut.GetType().String()),
136 zap.Error(logError))
137 mutError = fs.ErrInvalid
138 }
139 dirty = true
140 m.err <- mutError
141 }
142 }
143}
144 
145func (d *DiskKV) Stop() {
146 d.writeBarrier.Lock()
147 defer d.writeBarrier.Unlock()
148 
149 if !d.closed.CompareAndSwap(false, true) {
150 return
151 }
152 
153 close(d.closeCh)
154 d.closeWg.Wait()
155 
156 d.logger.Info("Flushing logs to disk")
157 
158 if err := d.log.Sync(); err != nil {
159 d.logger.Error("Error flushing logs to disk", zap.Error(err))
160 }
161 if err := d.log.Close(); err != nil {
162 d.logger.Error("Error closing log file", zap.Error(err))
163 }
164}
165 
166var _ chord.KVProvider = (*DiskKV)(nil)