From ae97446835f67c6132548f4f362dfde7b9c6f9ea Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Wed, 15 Apr 2026 09:51:36 +0530 Subject: [PATCH 01/14] introduce configurable fsync strategies --- cmd/server/main.go | 3 ++- internal/wal/wal.go | 39 +++++++++++++++++++++++++++++++++++---- 2 files changed, 37 insertions(+), 5 deletions(-) diff --git a/cmd/server/main.go b/cmd/server/main.go index 6ef4c1c..492fed5 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -13,6 +13,7 @@ import ( const port = ":9000" // TODO configs const walFilePath = "bin/wal.log" // TODO configs +const fsyncStrategy = wal.EVERY_SEC // TODO configs func main() { ctx, cancel := context.WithCancel(context.Background()) @@ -21,7 +22,7 @@ func main() { go handleInterrupts(sigChan, cancel) store := store.NewStore() - wal, err := wal.NewWal(walFilePath) + wal, err := wal.NewWal(walFilePath, fsyncStrategy, ctx) if err != nil { log.Fatal("unable to open wal file: ", err) } diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 00773d8..7c42a71 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -5,25 +5,40 @@ import ( "os" "strings" "sync" + "context" + "time" ) +type FsyncStrategy int + +const ( + ALWAYS FsyncStrategy = iota + EVERY_SEC + NEVER +) type Wal struct { mu sync.Mutex file *os.File closed bool + fsyncStrategy FsyncStrategy } -func NewWal(filepath string) (*Wal, error) { +func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context) (*Wal, error) { file, err := os.OpenFile(filepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { return nil, err } - return &Wal{ + w := &Wal{ mu: sync.Mutex{}, file: file, closed: false, - }, nil + } + if fsyncStrategy == EVERY_SEC { + go scheduleFsyncEverySec(w, ctx) + } + + return w, nil } func (w *Wal) Append(cmd ... string) (int, error) { @@ -32,7 +47,9 @@ func (w *Wal) Append(cmd ... string) (int, error) { if !w.closed { bytes, err := fmt.Fprintf(w.file, "%s\n", strings.Join(cmd, " ")) if err == nil { - err = w.file.Sync() + if w.fsyncStrategy == ALWAYS { + err = w.file.Sync() + } } return bytes, err } @@ -45,4 +62,18 @@ func (w *Wal) Close() error { w.closed = true return w.file.Close() +} + +func scheduleFsyncEverySec(w *Wal, ctx context.Context) { + ticker := time.NewTicker(1 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + w.file.Sync() + case <-ctx.Done(): + return + } + } } \ No newline at end of file From d8ff8fd0b814b24fefd0bc5ce60fd944f51e7845 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Wed, 15 Apr 2026 10:00:37 +0530 Subject: [PATCH 02/14] stop server if fsync is failing --- cmd/server/main.go | 2 +- internal/wal/wal.go | 15 ++++++++++----- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/cmd/server/main.go b/cmd/server/main.go index 492fed5..38eb46d 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -22,7 +22,7 @@ func main() { go handleInterrupts(sigChan, cancel) store := store.NewStore() - wal, err := wal.NewWal(walFilePath, fsyncStrategy, ctx) + wal, err := wal.NewWal(walFilePath, fsyncStrategy, ctx, cancel) if err != nil { log.Fatal("unable to open wal file: ", err) } diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 7c42a71..c7d23f7 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -1,11 +1,12 @@ package wal import ( + "context" "fmt" + "log" "os" "strings" "sync" - "context" "time" ) @@ -24,7 +25,7 @@ type Wal struct { fsyncStrategy FsyncStrategy } -func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context) (*Wal, error) { +func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, cancel context.CancelFunc) (*Wal, error) { file, err := os.OpenFile(filepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { return nil, err @@ -35,7 +36,7 @@ func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context) ( closed: false, } if fsyncStrategy == EVERY_SEC { - go scheduleFsyncEverySec(w, ctx) + go scheduleFsyncEverySec(w, ctx, cancel) } return w, nil @@ -64,14 +65,18 @@ func (w *Wal) Close() error { return w.file.Close() } -func scheduleFsyncEverySec(w *Wal, ctx context.Context) { +func scheduleFsyncEverySec(w *Wal, ctx context.Context, cancel context.CancelFunc) { ticker := time.NewTicker(1 * time.Second) defer ticker.Stop() for { select { case <-ticker.C: - w.file.Sync() + err := w.file.Sync() + if err != nil { + log.Printf("Fsync failure: %s, stopping the server", err) + cancel() + } case <-ctx.Done(): return } From 20e683dc37db3dd0f93bf8231d7c5b2f67e1f136 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Thu, 16 Apr 2026 09:18:13 +0530 Subject: [PATCH 03/14] create a wal file iterator --- internal/wal/wal.go | 22 +++++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index c7d23f7..601fc75 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -1,8 +1,11 @@ package wal import ( + "bufio" "context" "fmt" + "io" + "iter" "log" "os" "strings" @@ -81,4 +84,21 @@ func scheduleFsyncEverySec(w *Wal, ctx context.Context, cancel context.CancelFun return } } -} \ No newline at end of file +} + +func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { + return func(yield func(string, error) bool) { + f, err := os.Open(w.file.Name()) // create a separate fd + if err != nil { yield("", err); return } + defer f.Close() + + _, err = f.Seek(fromOffset, io.SeekStart) + if err != nil { yield("", err); return } + + scanner := bufio.NewScanner(f) + for scanner.Scan() { + if !yield(scanner.Text(), nil) { return } + } + if err := scanner.Err(); err != nil { yield("", err) } + } +} From 562ed53d954f84d2e6ce90907a8c670120d1ca2f Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Thu, 16 Apr 2026 09:18:50 +0530 Subject: [PATCH 04/14] replay wal on startup --- internal/handler/handler.go | 14 ++++++++++++++ internal/server/tcp_server.go | 12 ++++++++++++ 2 files changed, 26 insertions(+) diff --git a/internal/handler/handler.go b/internal/handler/handler.go index 3ba55b5..06ab6eb 100644 --- a/internal/handler/handler.go +++ b/internal/handler/handler.go @@ -2,6 +2,8 @@ package handler import ( "fmt" + "log" + "com.github.SantanuKar43/simple-kv/internal/protocol" "com.github.SantanuKar43/simple-kv/internal/store" "com.github.SantanuKar43/simple-kv/internal/wal" @@ -46,4 +48,16 @@ func Handle(input string, store *store.Store, wal *wal.Wal) string { default: return fmt.Sprintf("ERR unknown command %s", cmd.Name) } +} + +func Replay(input string, store *store.Store) { + cmd := protocol.Parse(input) + switch cmd.Name { + case "SET": + store.Set(cmd.Args[0], cmd.Args[1]) + case "DEL": + store.Delete(cmd.Args[0]) + default: + log.Printf("invalid command parsed %s\n", cmd.Name) + } } \ No newline at end of file diff --git a/internal/server/tcp_server.go b/internal/server/tcp_server.go index e2464fc..29c512d 100644 --- a/internal/server/tcp_server.go +++ b/internal/server/tcp_server.go @@ -26,6 +26,7 @@ func NewTCPServer(addr string, store *store.Store, wal *wal.Wal) *TCPServer { } func (s *TCPServer) Start(ctx context.Context) error { + err := replayWal(s.store, s.wal, 0) ln, err := net.Listen("tcp", s.addr) if err != nil { return err @@ -71,4 +72,15 @@ func (s *TCPServer) handleConnection(conn net.Conn) { response := handler.Handle(line, s.store, s.wal) conn.Write([]byte(response + "\n")) } +} + +func replayWal(store *store.Store, wal *wal.Wal, offset int64) error { + for line, err := range wal.WALIterator(offset) { + if err != nil { + log.Printf("Unable to replay WAL %s\n", err) + return err + } + handler.Replay(line, store) + } + return nil } \ No newline at end of file From 2b392df279ac56fba02dbb95673c85b3ee16ff23 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 18 Apr 2026 08:21:37 +0530 Subject: [PATCH 05/14] add checksum to wal entry, fix wal entry format, fix wal file on startup if corrupted --- internal/server/tcp_server.go | 6 +- internal/wal/wal.go | 111 +++++++++++++++++++++++++++++++--- 2 files changed, 108 insertions(+), 9 deletions(-) diff --git a/internal/server/tcp_server.go b/internal/server/tcp_server.go index 29c512d..7576de1 100644 --- a/internal/server/tcp_server.go +++ b/internal/server/tcp_server.go @@ -26,7 +26,11 @@ func NewTCPServer(addr string, store *store.Store, wal *wal.Wal) *TCPServer { } func (s *TCPServer) Start(ctx context.Context) error { - err := replayWal(s.store, s.wal, 0) + err := s.wal.FixIfBroken() + if err != nil { + return err + } + err = replayWal(s.store, s.wal, 0) ln, err := net.Listen("tcp", s.addr) if err != nil { return err diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 601fc75..300aa22 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -3,7 +3,9 @@ package wal import ( "bufio" "context" + "encoding/binary" "fmt" + "hash/crc32" "io" "iter" "log" @@ -49,7 +51,19 @@ func (w *Wal) Append(cmd ... string) (int, error) { w.mu.Lock() defer w.mu.Unlock() if !w.closed { - bytes, err := fmt.Fprintf(w.file, "%s\n", strings.Join(cmd, " ")) + // wal entry format - [length: 4 bytes][line: bytes][checksum: 4 bytes] + line := strings.Join(cmd, " ") + lineBytes := []byte(line) + length := len(lineBytes) + entry := make([]byte, length + 8) + + binary.LittleEndian.PutUint32(entry[:4], uint32(length)) + copy(entry[4:], lineBytes) + + checksum := crc32.ChecksumIEEE(entry[:(length + 4)]) + binary.LittleEndian.PutUint32(entry[(length + 4):], checksum) + + bytes, err := w.file.Write(entry) if err == nil { if w.fsyncStrategy == ALWAYS { err = w.file.Sync() @@ -86,19 +100,100 @@ func scheduleFsyncEverySec(w *Wal, ctx context.Context, cancel context.CancelFun } } +func (w *Wal) FixIfBroken() error { + f, err := os.OpenFile(w.file.Name(), os.O_RDWR, 0644) // create a separate fd + if err != nil { + return fmt.Errorf("Unable to open wal file: %s", err) + } + defer f.Close() + + var safeOffset int64 = 0 + var currOffset int64 = 0 + + reader := bufio.NewReader(f) + for { + lenBuf := make([]byte, 4) + _, err = io.ReadFull(reader, lenBuf) + if err == io.EOF { + return nil + } + if err != nil { + log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) + f.Truncate(safeOffset) + f.Sync() + return nil + } + currOffset += 4 + + length := binary.LittleEndian.Uint32(lenBuf) // todo check max entry size + total := int(length) + 8 // length + data + checksum + buf := make([]byte, total) + copy(buf[:4], lenBuf) + _, err = io.ReadFull(reader, buf[4:]) + if err != nil { + log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) + f.Truncate(safeOffset) + f.Sync() + return nil + } + currOffset += int64(length) + 4 + checksum := buf[(4 + length):] + checksumFound := binary.LittleEndian.Uint32(checksum) + checksumCalc := crc32.ChecksumIEEE(buf[:(length + 4)]) + if checksumFound != checksumCalc { + log.Printf("wal corrupted, checksum mismatch, truncating till last safe read offset %d\n", safeOffset) + f.Truncate(safeOffset) + f.Sync() + return nil + } + safeOffset = currOffset + } +} + func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { return func(yield func(string, error) bool) { f, err := os.Open(w.file.Name()) // create a separate fd - if err != nil { yield("", err); return } + if err != nil { + yield("", err) + return + } defer f.Close() - _, err = f.Seek(fromOffset, io.SeekStart) - if err != nil { yield("", err); return } + if _, err := f.Seek(fromOffset, io.SeekStart); err != nil { + yield("", err) + return + } + + reader := bufio.NewReader(f) + for { + lenBuf := make([]byte, 4) + _, err := io.ReadFull(reader, lenBuf) + if err == io.EOF { + return + } + if err != nil { + yield("", err) + return + } + + length := binary.LittleEndian.Uint32(lenBuf) + total := int(length) + 4 // data + checksum + buf := make([]byte, total) - scanner := bufio.NewScanner(f) - for scanner.Scan() { - if !yield(scanner.Text(), nil) { return } + _, err = io.ReadFull(reader, buf) + if err == io.EOF { + yield("", io.ErrUnexpectedEOF) + return + } + if err != nil { + yield("", err) + return + } + data := buf[:length] + // no need to validate checksum here as we always replay wal after validation + if !yield(string(data), nil) { + return + } } - if err := scanner.Err(); err != nil { yield("", err) } } } From a4a933975617b3d3353532011a1632ebfbece04f Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 25 Apr 2026 12:42:26 +0530 Subject: [PATCH 06/14] refactor if conditions --- internal/server/tcp_server.go | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/internal/server/tcp_server.go b/internal/server/tcp_server.go index 7576de1..8ed0f15 100644 --- a/internal/server/tcp_server.go +++ b/internal/server/tcp_server.go @@ -26,11 +26,12 @@ func NewTCPServer(addr string, store *store.Store, wal *wal.Wal) *TCPServer { } func (s *TCPServer) Start(ctx context.Context) error { - err := s.wal.FixIfBroken() - if err != nil { + if err := s.wal.FixIfBroken(); err != nil { + return err + } + if err := replayWal(s.store, s.wal, 0); err != nil { return err } - err = replayWal(s.store, s.wal, 0) ln, err := net.Listen("tcp", s.addr) if err != nil { return err From d02ff7dff04c0aa42e034e1fbc876c62f3e468c4 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 25 Apr 2026 12:57:08 +0530 Subject: [PATCH 07/14] make fixIfBroken private to wal and call it on new wal creation --- internal/server/tcp_server.go | 3 --- internal/wal/wal.go | 20 ++++++++++++++------ 2 files changed, 14 insertions(+), 9 deletions(-) diff --git a/internal/server/tcp_server.go b/internal/server/tcp_server.go index 8ed0f15..f4e8a3b 100644 --- a/internal/server/tcp_server.go +++ b/internal/server/tcp_server.go @@ -26,9 +26,6 @@ func NewTCPServer(addr string, store *store.Store, wal *wal.Wal) *TCPServer { } func (s *TCPServer) Start(ctx context.Context) error { - if err := s.wal.FixIfBroken(); err != nil { - return err - } if err := replayWal(s.store, s.wal, 0); err != nil { return err } diff --git a/internal/wal/wal.go b/internal/wal/wal.go index e15c251..e2d2b65 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -34,6 +34,9 @@ type Wal struct { } func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, fatalErrChan chan error) (*Wal, error) { + if err := fixIfBroken(filepath); err != nil { + return nil, err + } file, err := os.OpenFile(filepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { return nil, err @@ -143,8 +146,8 @@ func scheduleFsyncEverySec(w *Wal, ctx context.Context) { } } -func (w *Wal) FixIfBroken() error { - f, err := os.OpenFile(w.file.Name(), os.O_RDWR, 0644) // create a separate fd +func fixIfBroken(filepath string) error { + f, err := os.OpenFile(filepath, os.O_RDWR, 0644) // create a separate fd if err != nil { return fmt.Errorf("Unable to open wal file: %s", err) } @@ -155,6 +158,7 @@ func (w *Wal) FixIfBroken() error { reader := bufio.NewReader(f) for { + // read length lenBuf := make([]byte, 4) _, err = io.ReadFull(reader, lenBuf) if err == io.EOF { @@ -167,8 +171,9 @@ func (w *Wal) FixIfBroken() error { return nil } currOffset += 4 - length := binary.LittleEndian.Uint32(lenBuf) // todo check max entry size + + // read full wal entry ([length | data | checksum]) total := int(length) + 8 // length + data + checksum buf := make([]byte, total) copy(buf[:4], lenBuf) @@ -180,6 +185,7 @@ func (w *Wal) FixIfBroken() error { return nil } currOffset += int64(length) + 4 + // validate checksum checksum := buf[(4 + length):] checksumFound := binary.LittleEndian.Uint32(checksum) checksumCalc := crc32.ChecksumIEEE(buf[:(length + 4)]) @@ -209,6 +215,7 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { reader := bufio.NewReader(f) for { + // read length lenBuf := make([]byte, 4) _, err := io.ReadFull(reader, lenBuf) if err == io.EOF { @@ -218,11 +225,11 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { yield("", err) return } - length := binary.LittleEndian.Uint32(lenBuf) + + //read wal entry data with checksum total := int(length) + 4 // data + checksum buf := make([]byte, total) - _, err = io.ReadFull(reader, buf) if err == io.EOF { yield("", io.ErrUnexpectedEOF) @@ -233,7 +240,8 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { return } data := buf[:length] - // no need to validate checksum here as we always replay wal after validation + + // no need to validate checksum here as we will always replay wal after validation (wal.FixIfBroken) if !yield(string(data), nil) { return } From c9492317e5c7e6075f94b8cd0c73a2a0deb7a020 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 25 Apr 2026 13:03:05 +0530 Subject: [PATCH 08/14] gits --- internal/wal/wal.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index e2d2b65..998c5c6 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -241,7 +241,7 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { } data := buf[:length] - // no need to validate checksum here as we will always replay wal after validation (wal.FixIfBroken) + // no need to validate checksum here as we will always replay wal after validation (wal.fixIfBroken) if !yield(string(data), nil) { return } From 97f76d1e591d6e91ebb786f7c9022503977fe4e0 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 25 Apr 2026 13:03:32 +0530 Subject: [PATCH 09/14] update comments --- internal/wal/wal.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 998c5c6..5ffdd6e 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -241,7 +241,7 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { } data := buf[:length] - // no need to validate checksum here as we will always replay wal after validation (wal.fixIfBroken) + // no need to validate checksum here as we will always replay wal after validation, see [fixIfBroken] if !yield(string(data), nil) { return } From fe194328f1ac834d3b866b17085fb8d5a1132a26 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 25 Apr 2026 13:07:33 +0530 Subject: [PATCH 10/14] lock wal mutex on iteration --- internal/wal/wal.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 5ffdd6e..6160ddd 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -201,6 +201,8 @@ func fixIfBroken(filepath string) error { func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { return func(yield func(string, error) bool) { + w.mu.Lock() + defer w.mu.Unlock() f, err := os.Open(w.file.Name()) // create a separate fd if err != nil { yield("", err) From 63969a7157d5f0c69429ba7114f50d6264c707e8 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sat, 25 Apr 2026 13:55:46 +0530 Subject: [PATCH 11/14] make replay and repair in one pas --- internal/wal/wal.go | 106 ++++++++++++++++---------------------------- 1 file changed, 37 insertions(+), 69 deletions(-) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 6160ddd..7b40521 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -34,9 +34,6 @@ type Wal struct { } func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, fatalErrChan chan error) (*Wal, error) { - if err := fixIfBroken(filepath); err != nil { - return nil, err - } file, err := os.OpenFile(filepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { return nil, err @@ -146,64 +143,12 @@ func scheduleFsyncEverySec(w *Wal, ctx context.Context) { } } -func fixIfBroken(filepath string) error { - f, err := os.OpenFile(filepath, os.O_RDWR, 0644) // create a separate fd - if err != nil { - return fmt.Errorf("Unable to open wal file: %s", err) - } - defer f.Close() - - var safeOffset int64 = 0 - var currOffset int64 = 0 - - reader := bufio.NewReader(f) - for { - // read length - lenBuf := make([]byte, 4) - _, err = io.ReadFull(reader, lenBuf) - if err == io.EOF { - return nil - } - if err != nil { - log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) - f.Truncate(safeOffset) - f.Sync() - return nil - } - currOffset += 4 - length := binary.LittleEndian.Uint32(lenBuf) // todo check max entry size - - // read full wal entry ([length | data | checksum]) - total := int(length) + 8 // length + data + checksum - buf := make([]byte, total) - copy(buf[:4], lenBuf) - _, err = io.ReadFull(reader, buf[4:]) - if err != nil { - log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) - f.Truncate(safeOffset) - f.Sync() - return nil - } - currOffset += int64(length) + 4 - // validate checksum - checksum := buf[(4 + length):] - checksumFound := binary.LittleEndian.Uint32(checksum) - checksumCalc := crc32.ChecksumIEEE(buf[:(length + 4)]) - if checksumFound != checksumCalc { - log.Printf("wal corrupted, checksum mismatch, truncating till last safe read offset %d\n", safeOffset) - f.Truncate(safeOffset) - f.Sync() - return nil - } - safeOffset = currOffset - } -} - func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { + // the iterator also validates the wal file and truncates till last read safe offset return func(yield func(string, error) bool) { w.mu.Lock() defer w.mu.Unlock() - f, err := os.Open(w.file.Name()) // create a separate fd + f, err := os.OpenFile(w.file.Name(), os.O_RDWR, 0644) // create a separate fd if err != nil { yield("", err) return @@ -215,36 +160,59 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { return } + var safeOffset int64 = 0 + var currOffset int64 = 0 + reader := bufio.NewReader(f) for { // read length lenBuf := make([]byte, 4) _, err := io.ReadFull(reader, lenBuf) if err == io.EOF { + log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) + f.Truncate(safeOffset) + f.Sync() + yield("", err) return } if err != nil { - yield("", err) + log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) + f.Truncate(safeOffset) + f.Sync() + yield("", err) return } length := binary.LittleEndian.Uint32(lenBuf) + currOffset += 4 - //read wal entry data with checksum - total := int(length) + 4 // data + checksum + // read full wal entry ([length | data | checksum]) + total := int(length) + 8 // length + data + checksum buf := make([]byte, total) - _, err = io.ReadFull(reader, buf) - if err == io.EOF { - yield("", io.ErrUnexpectedEOF) - return - } + copy(buf[:4], lenBuf) + _, err = io.ReadFull(reader, buf[4:]) if err != nil { - yield("", err) + log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) + f.Truncate(safeOffset) + f.Sync() + yield("", err) return } - data := buf[:length] + data := buf[4:(length + 4)] + currOffset += int64(length) + 4 + + // validate checksum + checksum := buf[(4 + length):] + checksumFound := binary.LittleEndian.Uint32(checksum) + checksumCalc := crc32.ChecksumIEEE(buf[:(length + 4)]) + if checksumFound != checksumCalc { + log.Printf("wal corrupted, checksum mismatch, truncating till last safe read offset %d\n", safeOffset) + f.Truncate(safeOffset) + f.Sync() + return + } + safeOffset = currOffset - // no need to validate checksum here as we will always replay wal after validation, see [fixIfBroken] - if !yield(string(data), nil) { + if !yield(string(data), nil) { return } } From 9f880adf5349da610ccc0f433a6250a2bd00c9f3 Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sun, 26 Apr 2026 09:30:28 +0530 Subject: [PATCH 12/14] refactor and cleanup wal --- internal/wal/wal.go | 107 +++++++++++++++++++++++--------------------- 1 file changed, 57 insertions(+), 50 deletions(-) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 7b40521..fe9498f 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -33,6 +33,8 @@ type Wal struct { fatalErrChan chan error } +var byteOrder binary.ByteOrder = binary.LittleEndian + func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, fatalErrChan chan error) (*Wal, error) { file, err := os.OpenFile(filepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { @@ -53,49 +55,6 @@ func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, f return w, nil } -func (w *Wal) Append(cmd ... string) (int, error) { - w.mu.Lock() - defer w.mu.Unlock() - if !w.closed { - // wal entry format - [length: 4 bytes][line: bytes][checksum: 4 bytes] - line := strings.Join(cmd, " ") - - lineBytes := []byte(line) - length := len(lineBytes) - entry := make([]byte, length + 8) - - binary.LittleEndian.PutUint32(entry[:4], uint32(length)) - copy(entry[4:], lineBytes) - - checksum := crc32.ChecksumIEEE(entry[:(length + 4)]) - binary.LittleEndian.PutUint32(entry[(length + 4):], checksum) - - if w.fsyncStrategy != ALWAYS { - bytes, err := w.asyncBuffer.Write(entry) - return bytes, err - } - - bytes, err := w.file.Write(entry) - if err == nil { - err = w.file.Sync() - } - if err != nil { - log.Printf("wal append failure: %s, stopping the server", err) - w.fatalErrChan <- err - } - return bytes, err - } - return 0, fmt.Errorf("unable to append, wal already closed") -} - -func (w *Wal) Close() error { - w.mu.Lock() - defer w.mu.Unlock() - - w.closed = true - return w.file.Close() -} - func scheduleFsyncEverySec(w *Wal, ctx context.Context) { ticker := time.NewTicker(1 * time.Second) defer ticker.Stop() @@ -143,6 +102,50 @@ func scheduleFsyncEverySec(w *Wal, ctx context.Context) { } } +func (w *Wal) Append(cmd ... string) (int, error) { + w.mu.Lock() + defer w.mu.Unlock() + if !w.closed { + // wal entry format - [length: 4 bytes][line: bytes][checksum: 4 bytes] + line := strings.Join(cmd, " ") + entry := getWalEntry(line) + + if w.fsyncStrategy != ALWAYS { + bytes, err := w.asyncBuffer.Write(entry) + return bytes, err + } + + bytes, err := w.file.Write(entry) + if err == nil { + err = w.file.Sync() + } + if err != nil { + log.Printf("wal append failure: %s, stopping the server", err) + w.fatalErrChan <- err + } + return bytes, err + } + return 0, fmt.Errorf("unable to append, wal already closed") +} + +func getWalEntry(line string) []byte { + lineBytes := []byte(line) + length := len(lineBytes) + entry := make([]byte, length + 8) + + byteOrder.PutUint32(entry[:4], uint32(length)) + copy(entry[4:], lineBytes) + + checksum := calcChecksum(entry[:(length + 4)]) + byteOrder.PutUint32(entry[(length + 4):], checksum) + + return entry +} + +func calcChecksum(bytes []byte) uint32 { + return crc32.ChecksumIEEE(bytes) +} + func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { // the iterator also validates the wal file and truncates till last read safe offset return func(yield func(string, error) bool) { @@ -169,10 +172,6 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { lenBuf := make([]byte, 4) _, err := io.ReadFull(reader, lenBuf) if err == io.EOF { - log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) - f.Truncate(safeOffset) - f.Sync() - yield("", err) return } if err != nil { @@ -182,7 +181,7 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { yield("", err) return } - length := binary.LittleEndian.Uint32(lenBuf) + length := byteOrder.Uint32(lenBuf) currOffset += 4 // read full wal entry ([length | data | checksum]) @@ -202,8 +201,8 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { // validate checksum checksum := buf[(4 + length):] - checksumFound := binary.LittleEndian.Uint32(checksum) - checksumCalc := crc32.ChecksumIEEE(buf[:(length + 4)]) + checksumFound := byteOrder.Uint32(checksum) + checksumCalc := calcChecksum(buf[:(length + 4)]) if checksumFound != checksumCalc { log.Printf("wal corrupted, checksum mismatch, truncating till last safe read offset %d\n", safeOffset) f.Truncate(safeOffset) @@ -218,3 +217,11 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { } } } + +func (w *Wal) Close() error { + w.mu.Lock() + defer w.mu.Unlock() + + w.closed = true + return w.file.Close() +} From f6224835a84f84ebd4142c1933997bf896316b7c Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sun, 26 Apr 2026 09:33:39 +0530 Subject: [PATCH 13/14] add comments --- internal/wal/wal.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/internal/wal/wal.go b/internal/wal/wal.go index fe9498f..590b82c 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -106,7 +106,6 @@ func (w *Wal) Append(cmd ... string) (int, error) { w.mu.Lock() defer w.mu.Unlock() if !w.closed { - // wal entry format - [length: 4 bytes][line: bytes][checksum: 4 bytes] line := strings.Join(cmd, " ") entry := getWalEntry(line) @@ -129,6 +128,8 @@ func (w *Wal) Append(cmd ... string) (int, error) { } func getWalEntry(line string) []byte { + // wal entry format - [length: 4 bytes][line: bytes][checksum: 4 bytes] + lineBytes := []byte(line) length := len(lineBytes) entry := make([]byte, length + 8) From f793206fd6f04e0724556dcde0a9e9463620103d Mon Sep 17 00:00:00 2001 From: SantanuKar43 Date: Sun, 26 Apr 2026 10:36:39 +0530 Subject: [PATCH 14/14] use reusable buffers for wal entries and iteration --- cmd/server/main.go | 3 ++- internal/wal/wal.go | 50 ++++++++++++++++++++++++++++----------------- 2 files changed, 33 insertions(+), 20 deletions(-) diff --git a/cmd/server/main.go b/cmd/server/main.go index 960aa5d..8273cb4 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -14,6 +14,7 @@ import ( const port = ":9000" // TODO configs const walFilePath = "bin/wal.log" // TODO configs const fsyncStrategy = wal.EVERY_SEC // TODO configs +const maxWalEntryLength = 1 << 30 // TODO configs func main() { ctx, cancel := context.WithCancel(context.Background()) @@ -24,7 +25,7 @@ func main() { go handleFatalErrors(fatalErrChan, cancel) store := store.NewStore() - wal, err := wal.NewWal(walFilePath, fsyncStrategy, ctx, fatalErrChan) + wal, err := wal.NewWal(walFilePath, fsyncStrategy, maxWalEntryLength, ctx, fatalErrChan) if err != nil { log.Fatal("unable to open wal file: ", err) } diff --git a/internal/wal/wal.go b/internal/wal/wal.go index 590b82c..56f1750 100644 --- a/internal/wal/wal.go +++ b/internal/wal/wal.go @@ -30,12 +30,14 @@ type Wal struct { closed bool fsyncStrategy FsyncStrategy asyncBuffer *bytes.Buffer + entryBuf []byte fatalErrChan chan error + maxEntryLen uint32 } var byteOrder binary.ByteOrder = binary.LittleEndian -func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, fatalErrChan chan error) (*Wal, error) { +func NewWal(filepath string, fsyncStrategy FsyncStrategy, maxWalEntryLength uint32, ctx context.Context, fatalErrChan chan error) (*Wal, error) { file, err := os.OpenFile(filepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { return nil, err @@ -46,6 +48,8 @@ func NewWal(filepath string, fsyncStrategy FsyncStrategy, ctx context.Context, f closed: false, fsyncStrategy: fsyncStrategy, fatalErrChan: fatalErrChan, + maxEntryLen: maxWalEntryLength, + entryBuf: make([]byte, maxWalEntryLength), } if fsyncStrategy != ALWAYS { w.asyncBuffer = new(bytes.Buffer) @@ -107,7 +111,10 @@ func (w *Wal) Append(cmd ... string) (int, error) { defer w.mu.Unlock() if !w.closed { line := strings.Join(cmd, " ") - entry := getWalEntry(line) + entry, err := w.getWalEntry(line) + if err != nil { + return 0, err + } if w.fsyncStrategy != ALWAYS { bytes, err := w.asyncBuffer.Write(entry) @@ -127,20 +134,21 @@ func (w *Wal) Append(cmd ... string) (int, error) { return 0, fmt.Errorf("unable to append, wal already closed") } -func getWalEntry(line string) []byte { +func (w *Wal) getWalEntry(line string) ([]byte, error) { // wal entry format - [length: 4 bytes][line: bytes][checksum: 4 bytes] - lineBytes := []byte(line) length := len(lineBytes) - entry := make([]byte, length + 8) + if length > int(w.maxEntryLen) { + return nil, fmt.Errorf("line length > maxLength, can't append") + } - byteOrder.PutUint32(entry[:4], uint32(length)) - copy(entry[4:], lineBytes) + byteOrder.PutUint32(w.entryBuf[:4], uint32(length)) + copy(w.entryBuf[4:(length + 4)], lineBytes) - checksum := calcChecksum(entry[:(length + 4)]) - byteOrder.PutUint32(entry[(length + 4):], checksum) + checksum := calcChecksum(w.entryBuf[:(length + 4)]) + byteOrder.PutUint32(w.entryBuf[(length + 4):(length + 8)], checksum) - return entry + return w.entryBuf[:(length + 8)], nil } func calcChecksum(bytes []byte) uint32 { @@ -167,11 +175,12 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { var safeOffset int64 = 0 var currOffset int64 = 0 + buf := make([]byte, w.maxEntryLen) reader := bufio.NewReader(f) + for { // read length - lenBuf := make([]byte, 4) - _, err := io.ReadFull(reader, lenBuf) + _, err := io.ReadFull(reader, buf[:4]) if err == io.EOF { return } @@ -182,14 +191,17 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { yield("", err) return } - length := byteOrder.Uint32(lenBuf) + length := byteOrder.Uint32(buf[:4]) currOffset += 4 + if length > w.maxEntryLen { + log.Printf("wal corrupted, length found too big, truncating till last safe read offset %d\n", safeOffset) + f.Truncate(safeOffset) + f.Sync() + return + } // read full wal entry ([length | data | checksum]) - total := int(length) + 8 // length + data + checksum - buf := make([]byte, total) - copy(buf[:4], lenBuf) - _, err = io.ReadFull(reader, buf[4:]) + _, err = io.ReadFull(reader, buf[4:(8 + length)]) if err != nil { log.Printf("wal corrupted, err %s, truncating till last safe read offset %d\n", err, safeOffset) f.Truncate(safeOffset) @@ -201,9 +213,9 @@ func (w *Wal) WALIterator(fromOffset int64) iter.Seq2[string, error] { currOffset += int64(length) + 4 // validate checksum - checksum := buf[(4 + length):] + checksum := buf[(4 + length):(8 + length)] checksumFound := byteOrder.Uint32(checksum) - checksumCalc := calcChecksum(buf[:(length + 4)]) + checksumCalc := calcChecksum(buf[:(4 + length)]) if checksumFound != checksumCalc { log.Printf("wal corrupted, checksum mismatch, truncating till last safe read offset %d\n", safeOffset) f.Truncate(safeOffset)