File
Blob: kv/aof/log.go
| 1 | package aof |
| 2 | |
| 3 | import ( |
| 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 | |
| 13 | var crcTable = crc64.MakeTable(crc64.ECMA) |
| 14 | |
| 15 | func (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 | |
| 44 | func (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 | |
| 59 | func (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 | |
| 93 | func (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 | } |