Skip to content
File

Blob: kv/aof/log.go

go106 lines
1package aof
2 
3import (
4 "fmt"
5 "hash/crc64"
6 
7 "go.miragespace.co/specter/kv/aof/proto"
8 
9 bufPool "github.com/libp2p/go-buffer-pool"
10 "go.uber.org/zap"
11)
12 
13var crcTable = crc64.MakeTable(crc64.ECMA)
14 
15func (d *DiskKV) replayLogs() error {
16 index, err := d.log.LastIndex()
17 if err != nil {
18 return fmt.Errorf("error reading last log index: %w", err)
19 }
20 d.logger.Info("Replaying mutation logs", zap.Uint64("index", index))
21 mut := &proto.Mutation{}
22 entry := &proto.LogEntry{}
23 for i := uint64(1); i <= index; i++ {
24 buf, err := d.log.Read(i)
25 if err != nil {
26 return fmt.Errorf("error reading log at index %d: %w", i, err)
27 }
28 if err := entry.UnmarshalVT(buf); err != nil {
29 return fmt.Errorf("error deserializing log at index %d: %w", i, err)
30 }
31 if err := d.decodeEntry(entry, mut); err != nil {
32 return fmt.Errorf("error decoding entry to mutation at index %d: %w", i, err)
33 }
34 if err := d.handleMutation(mut); err != nil {
35 return fmt.Errorf("error apply mutation to memory state at index %d: %w", i, err)
36 }
37 entry.Reset()
38 mut.Reset()
39 }
40 d.counter = index + 1
41 return nil
42}
43 
44func (d *DiskKV) decodeEntry(entry *proto.LogEntry, mut *proto.Mutation) (err error) {
45 if entry.GetChecksum() != crc64.Checksum(entry.GetData(), crcTable) {
46 err = fmt.Errorf("log entry checksum does not match, possibly corrupted log")
47 return
48 }
49 switch entry.GetVersion() {
50 case proto.LogVersion_V1:
51 // uncompressed
52 err = mut.UnmarshalVT(entry.Data)
53 default:
54 err = fmt.Errorf("unknown log version: %s", entry.GetVersion())
55 }
56 return
57}
58 
59func (d *DiskKV) appendLog(mut *proto.Mutation) error {
60 mutBuf := bufPool.Get(mut.SizeVT())
61 defer bufPool.Put(mutBuf)
62 
63 _, err := mut.MarshalToSizedBufferVT(mutBuf)
64 if err != nil {
65 d.logger.Error("Error serializing mutation", zap.String("mutation", mut.GetType().String()), zap.Error(err))
66 return err
67 }
68 
69 entry := proto.LogEntryFromVTPool()
70 defer entry.ReturnToVTPool()
71 
72 entry.Version = proto.LogVersion_V1
73 entry.Data = mutBuf
74 entry.Checksum = crc64.Checksum(mutBuf, crcTable)
75 
76 logBuf := bufPool.Get(entry.SizeVT())
77 defer bufPool.Put(logBuf)
78 
79 _, err = entry.MarshalToSizedBufferVT(logBuf)
80 if err != nil {
81 d.logger.Error("Error serializing log entry", zap.String("version", entry.GetVersion().String()), zap.String("mutation", mut.GetType().String()), zap.Error(err))
82 return err
83 }
84 
85 if err := d.log.Write(d.counter, logBuf); err != nil {
86 d.logger.Error("Error appending to log", zap.Uint64("counter", d.counter), zap.String("mutation", mut.GetType().String()), zap.Error(err))
87 return err
88 }
89 d.counter += 1
90 return nil
91}
92 
93func (d *DiskKV) rollbackOne(mut *proto.Mutation, err error) {
94 d.logger.Warn("Rolling back last mutation because of an error",
95 zap.String("mutation", mut.GetType().String()),
96 zap.Uint64("truncate", d.counter-2),
97 zap.Uint64("index", d.counter-1),
98 zap.Error(err),
99 )
100 d.counter -= 1
101 if err := d.log.TruncateBack(d.counter - 1); err != nil {
102 d.logger.Error("Error applying rollback to the last mutation",
103 zap.Error(err))
104 }
105}