From adb69d1b6d72aaed5b20a09b6fb85d3f8d0c3a77 Mon Sep 17 00:00:00 2001 From: AltSumpreme Date: Fri, 18 Jul 2025 00:04:41 +0530 Subject: [PATCH 1/2] feat(pager):switching from json to binary storage and major refactor made in the page and storage - Switched from json to binary storage and stored the types in the meta json file - added tests for the read and write --- .gitignore | 2 + db/database.go | 2 +- db/types.go | 10 +- engine.go | 11 +- pager/pager.go | 232 ++++++++++++++++++++++++++++++++++------- repl/repl.go | 2 +- storage/jsonstorage.go | 71 ------------- storage/storage.go | 140 +++++++++++++++++++++++++ tests/write_test.go | 110 +++++++++++++++++++ 9 files changed, 459 insertions(+), 121 deletions(-) delete mode 100644 storage/jsonstorage.go create mode 100644 storage/storage.go create mode 100644 tests/write_test.go diff --git a/.gitignore b/.gitignore index 7553af7..7d823be 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,4 @@ *.log *.json +./db + diff --git a/db/database.go b/db/database.go index be5d797..2f3ba4c 100644 --- a/db/database.go +++ b/db/database.go @@ -6,7 +6,7 @@ import ( ) type Database struct { - Tables map[string]*Table `json:"tables"` + Tables map[string]*Table } func NewDatabase() *Database { diff --git a/db/types.go b/db/types.go index 4b0ba2c..1738e07 100644 --- a/db/types.go +++ b/db/types.go @@ -8,14 +8,14 @@ const ( ) type Column struct { - Name string `json:"name"` - Type FieldType `json:"type"` + Name string + Type FieldType } type Row map[string]interface{} type Table struct { - Name string `json:"name"` - Columns []Column `json:"columns"` - Rows []Row `json:"rows"` + Name string + Columns []Column + Rows []Row } diff --git a/engine.go b/engine.go index f0a2bac..d8b1192 100644 --- a/engine.go +++ b/engine.go @@ -9,13 +9,13 @@ import ( type Engine struct { DB *db.Database - Pager *pager.Pager + Pager *pager.Page } func NewEngine() (*Engine, error) { var database *db.Database - if _, err := os.Stat(storage.Path); err == nil { + if _, err := os.Stat(storage.DBDir); err == nil { database, err = storage.LoadFromDisk() if err != nil { return nil, err @@ -25,11 +25,11 @@ func NewEngine() (*Engine, error) { } else { return nil, err } - pagerFile, err := os.OpenFile(storage.Path, os.O_RDWR|os.O_CREATE, 0666) + err := os.MkdirAll(storage.DBDir, 0775) if err != nil { return nil, err } - pgr := pager.PagerInit(pagerFile) + pgr := pager.NewPage() if pgr == nil { return nil, err } @@ -46,5 +46,6 @@ func (engine *Engine) Close() error { if err := storage.SaveToDisk(engine.DB); err != nil { return err } - return engine.Pager.Close() + + return nil } diff --git a/pager/pager.go b/pager/pager.go index efd25cc..10a2343 100644 --- a/pager/pager.go +++ b/pager/pager.go @@ -1,64 +1,220 @@ package pager -import "os" +import ( + "bytes" + "encoding/binary" + "errors" + "fmt" + "io" + "pebbledb/db" +) -const PageSize = 4096 // 4KB page size +const ( + PageSize = 4096 + PageHeaderSize = 16 + ItemIDSize = 6 + MaxItemsPerPage = 128 + SpecialSpaceSize = 0 // Reserved for future user for indexing metadata + DataRegionSize = PageSize - PageHeaderSize - MaxItemsPerPage*ItemIDSize - SpecialSpaceSize +) -type Page struct { - ID int - Data [PageSize]byte +type PageHeader struct { + LSN uint64 // Future WAL support + NumItems uint16 + PdLower uint16 + PdUpper uint16 +} + +type ItemID struct { + Offset uint16 + Length uint16 + DeletedFlag uint8 + _ [1]byte // Padding to align to 6 bytes } -type Pager struct { - file *os.File - pageSize int - nextPageID int + +type Page struct { + Header PageHeader + Items [MaxItemsPerPage]ItemID + Data [DataRegionSize]byte } -func PagerInit(file *os.File) *Pager { - return &Pager{ - file: file, - pageSize: PageSize, - nextPageID: 0, +func NewPage() *Page { + return &Page{ + Header: PageHeader{ + PdLower: PageHeaderSize, + PdUpper: DataRegionSize, + }, + Items: [MaxItemsPerPage]ItemID{}, + Data: [DataRegionSize]byte{}, } } -func (p *Pager) NewPage() *Page { - page := &Page{ - ID: p.nextPageID, - Data: [PageSize]byte{}, +func (p *Page) InsertTuple(record []byte) (int, error) { + if len(record) > DataRegionSize { + return -1, errors.New("record too large") + } + if len(record) == 0 { + return -1, errors.New("record cannot be empty") + } + if p.Header.NumItems >= MaxItemsPerPage { + return -1, errors.New("page is full") + } + + // 🧠 Calculate available free space + freeSpace := int(p.Header.PdUpper) - int(p.Header.PdLower) + requiredSpace := len(record) + ItemIDSize + if requiredSpace > freeSpace { + return -1, errors.New("not enough space in page") + } + + // 👇 Insert the record at the top + newUpper := p.Header.PdUpper - uint16(len(record)) + copy(p.Data[newUpper:], record) + + // 👇 Add the metadata for this tuple + slot := int(p.Header.NumItems) + p.Items[slot] = ItemID{ + Offset: newUpper, + Length: uint16(len(record)), + DeletedFlag: 1, } - p.nextPageID++ - return page + + p.Header.PdUpper = newUpper + p.Header.PdLower += ItemIDSize + p.Header.NumItems++ + + return slot, nil } -func (p *Pager) WritePage(page *Page) error { - offset := int64(page.ID) * PageSize - _, err := p.file.WriteAt(page.Data[:], offset) - if err != nil { - return err +func (p *Page) ReadTuple(slot int) ([]byte, error) { + + if slot < 0 || slot >= int(p.Header.NumItems) { + return nil, errors.New("invalid slot number") } + items := p.Items[slot] + if items.DeletedFlag == 0 { + return nil, errors.New("tuple has been deleted") + } + data := make([]byte, items.Length) + copy(data, p.Data[items.Offset:items.Offset+items.Length]) + return data, nil +} +func (p *Page) DeleteTuple(slot int) error { + if slot < 0 || slot >= int(p.Header.NumItems) { + return errors.New("invalid slot number") + } + items := p.Items[slot] + if items.DeletedFlag == 0 { + return errors.New("tuple has already been deleted") + } + p.Items[slot].DeletedFlag = 0 return nil } +func SerializeRow(row db.Row, columns []db.Column) ([]byte, error) { + var buf bytes.Buffer + for _, col := range columns { + val := row[col.Name] + switch col.Type { + case db.TypeInt: + if v, ok := val.(int); ok { + if err := binary.Write(&buf, binary.LittleEndian, int32(v)); err != nil { + return nil, err + } + + } + case db.TypeString: + if str, ok := val.(string); ok { + if err := binary.Write(&buf, binary.LittleEndian, uint16(len(str))); err != nil { + return nil, err + } + if _, err := buf.Write([]byte(str)); err != nil { + return nil, err + } + } -func (p *Pager) ReadPage(pageID int) (*Page, error) { - page := &Page{ - ID: pageID, - Data: [PageSize]byte{}, + } } + return buf.Bytes(), nil +} - offset := int64(pageID) * PageSize - _, err := p.file.ReadAt(page.Data[:], offset) - if err != nil { - return nil, err +func DeserializeRow(data []byte, columns []db.Column) (db.Row, error) { + row := make(db.Row) + buf := bytes.NewBuffer(data) + + for _, col := range columns { + switch col.Type { + case db.TypeInt: + var val int32 + if err := binary.Read(buf, binary.LittleEndian, &val); err != nil { + return nil, err + } + row[col.Name] = int(val) + case db.TypeString: + var length uint16 + if err := binary.Read(buf, binary.LittleEndian, &length); err != nil { + return nil, err + } + strData := make([]byte, length) + if _, err := io.ReadFull(buf, strData); err != nil { + return nil, err + } + row[col.Name] = string(strData) + } } + fmt.Printf("%v\n", row) + return row, nil +} - return page, nil +func SerializePage(page *Page) []byte { + buf := make([]byte, PageSize) + writer := bytes.NewBuffer(buf[:0]) + + if err := binary.Write(writer, binary.LittleEndian, &page.Header); err != nil { + panic("failed to write page header: " + err.Error()) + } + + for i := 0; i < MaxItemsPerPage; i++ { + if err := binary.Write(writer, binary.LittleEndian, &page.Items[i]); err != nil { + panic("failed to write page item: " + err.Error()) + } + } + + if _, err := writer.Write(page.Data[:]); err != nil { + panic("failed to write page data: " + err.Error()) + } + final := writer.Bytes() + if len(final) < PageSize { + padding := make([]byte, PageSize-len(final)) + final = append(final, padding...) + } + + return final } -func (p *Pager) Close() error { - if p.file != nil { - return p.file.Close() +func DeserializePage(buf []byte) (*Page, error) { + if len(buf) < PageSize { + return nil, errors.New("buffer too small to be a valid page") } - return nil + + page := &Page{} + reader := bytes.NewReader(buf) + + if err := binary.Read(reader, binary.LittleEndian, &page.Header); err != nil { + return nil, err + } + + for i := 0; i < MaxItemsPerPage; i++ { + var items ItemID + if err := binary.Read(reader, binary.LittleEndian, &items); err != nil { + + return nil, err + } + page.Items[i] = items + } + + if _, err := reader.Read(page.Data[:]); err != nil { + return nil, err + } + return page, nil } diff --git a/repl/repl.go b/repl/repl.go index 1a7f9d1..1cef9c9 100644 --- a/repl/repl.go +++ b/repl/repl.go @@ -16,7 +16,7 @@ func ReplInit(database *db.Database) { fmt.Println("Welcome to PebbleDB ! Type 'Exit' to quit.") for { - fmt.Println("> ") + fmt.Print("> ") input, _ := reader.ReadString('\n') input = strings.TrimSpace(input) if input == "" { diff --git a/storage/jsonstorage.go b/storage/jsonstorage.go deleted file mode 100644 index 56a2250..0000000 --- a/storage/jsonstorage.go +++ /dev/null @@ -1,71 +0,0 @@ -package storage - -import ( - "encoding/json" - "fmt" - "os" - "pebbledb/db" -) - -const Path = "./db.json" - -func SaveToDisk(database *db.Database) error { - file, err := os.Create(Path) - - if err != nil { - return err - } - defer file.Close() - - encoder := json.NewEncoder(file) - encoder.SetIndent("", " ") - return encoder.Encode(database) -} -func LoadFromDisk() (*db.Database, error) { - data, err := os.ReadFile(Path) - if err != nil { - return nil, fmt.Errorf("failed to read storage file: %w", err) - } - - if len(data) == 0 { - fmt.Println("Storage file is empty. Initializing fresh database.") - return db.NewDatabase(), nil - } - - var dbState db.Database - if err := json.Unmarshal(data, &dbState); err != nil { - return nil, fmt.Errorf("failed to parse JSON: %w", err) - } - if dbState.Tables == nil { - fmt.Println("Warning: No tables loaded from disk.") - } - - if len(dbState.Tables) == 0 { - fmt.Println("Warning: No tables loaded from disk.") - } - for _, table := range dbState.Tables { - normalizeRows(table) - } - - return &dbState, nil -} - -func normalizeRows(table *db.Table) { - for i, row := range table.Rows { - newRow := make(db.Row) - for _, col := range table.Columns { - val := row[col.Name] - switch col.Type { - case db.TypeInt: - if f, ok := val.(float64); ok { - newRow[col.Name] = int(f) - } else { - newRow[col.Name] = val - } - default: - newRow[col.Name] = val - } - } - table.Rows[i] = newRow - } -} diff --git a/storage/storage.go b/storage/storage.go new file mode 100644 index 0000000..1a65236 --- /dev/null +++ b/storage/storage.go @@ -0,0 +1,140 @@ +package storage + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "pebbledb/db" + "pebbledb/pager" +) + +const DBDir = "./db" + +func SaveToDisk(database *db.Database) error { + if err := os.MkdirAll(DBDir, 0775); err != nil { + return err + } + + for tableName, table := range database.Tables { + metaPath := filepath.Join(DBDir, tableName+".meta.json") + metaFile, err := os.Create(metaPath) + if err != nil { + return fmt.Errorf("create schema file for %s: %w", tableName, err) + } + if err := json.NewEncoder(metaFile).Encode(table.Columns); err != nil { + return fmt.Errorf("encode schema for %s: %w", tableName, err) + } + metaFile.Close() + dataPath := filepath.Join(DBDir, tableName+".db") + dataFile, err := os.Create(dataPath) + if err != nil { + return fmt.Errorf("create data file for %s: %w", tableName, err) + } + + page := pager.NewPage() + for _, row := range table.Rows { + fmt.Printf("Serializing row: %v\n", row) + serialized, err := pager.SerializeRow(row, table.Columns) + fmt.Printf("Serialized row: %v\n", serialized) + if err != nil { + return fmt.Errorf("serialize row: %w", err) + } + if _, err := page.InsertTuple(serialized); err != nil { + return fmt.Errorf("insert tuple: %w", err) + } + } + + if _, err := dataFile.Write(pager.SerializePage(page)); err != nil { + return fmt.Errorf("write page for %s: %w", tableName, err) + } + dataFile.Close() + } + + return nil +} + +func LoadSchemaFromDisk(tableName string) ([]db.Column, error) { + metafile := filepath.Join(DBDir, tableName+".meta.json") + if _, err := os.Stat(metafile); os.IsNotExist(err) { + return nil, fmt.Errorf("schema file for table %s does not exist", tableName) + } + file, err := os.Open(metafile) + if err != nil { + return nil, fmt.Errorf("failed to open schema file for table %s: %w", tableName, err) + } + defer file.Close() + + var columnDef []db.Column + if err := json.NewDecoder(file).Decode(&columnDef); err != nil { + return nil, fmt.Errorf("failed to decode schema for table %s: %w", tableName, err) + } + return columnDef, nil +} +func LoadFromDisk() (*db.Database, error) { + files, err := os.ReadDir(DBDir) + if err != nil { + return nil, fmt.Errorf("failed to read directory %s: %w", DBDir, err) + } + database := db.NewDatabase() + for _, file := range files { + if file.IsDir() { + continue + } + filePath := fmt.Sprintf(DBDir + "/" + file.Name()) + f, err := os.Open(filePath) + if err != nil { + return nil, fmt.Errorf("failed to open file %s: %w", filePath, err) + } + + buf := make([]byte, pager.PageSize) + _, err = f.Read(buf) + f.Close() + if err != nil { + return nil, fmt.Errorf("failed to read file %s: %w", filePath, err) + } + + page, err := pager.DeserializePage(buf) + + if err != nil { + return nil, fmt.Errorf("failed to deserialize page from file %s: %w", filePath, err) + } + if filepath.Ext(filePath) != ".db" { + continue + } + tableName := file.Name() + tableName = tableName[:len(tableName)-3] + columnDef, err := LoadSchemaFromDisk(tableName) + if err != nil { + return nil, fmt.Errorf("table %s does not exist in database: %w", tableName, err) + } + table := &db.Table{ + Name: tableName, + Columns: columnDef, + Rows: []db.Row{}, + } + for i := 0; i < int(page.Header.NumItems); i++ { + items := page.Items[i] + if items.DeletedFlag == 0 || items.Length == 0 { + continue + } + + tupleData, err := page.ReadTuple(i) + + if err != nil { + return nil, fmt.Errorf("failed to read tuple from page: %w", err) + } + row, err := pager.DeserializeRow(tupleData, table.Columns) + if err != nil { + return nil, fmt.Errorf("failed to deserialize row: %w", err) + } + fmt.Printf("Deserialized row: %v\n", row) + table.Rows = append(table.Rows, row) + + } + database.Tables[tableName] = table + + } + return database, nil + +} diff --git a/tests/write_test.go b/tests/write_test.go new file mode 100644 index 0000000..6589e5b --- /dev/null +++ b/tests/write_test.go @@ -0,0 +1,110 @@ +package tests + +import ( + "fmt" + "os" + "path/filepath" + "pebbledb/db" + "pebbledb/pager" + "pebbledb/storage" + "testing" +) + +func TestWrite(t *testing.T) { + _ = os.MkdirAll(storage.DBDir, 0775) + + table := &db.Table{ + Name: "users", + Columns: []db.Column{ + {Name: "id", Type: db.TypeInt}, + {Name: "name", Type: db.TypeString}, + }, + Rows: []db.Row{ + {"id": 42, "name": "Reuben"}, + {"id": 43, "name": "Alice"}, + {"id": 44, "name": "Bob"}, + {"id": 45, "name": "Charlie"}, + {"id": 46, "name": "Diana"}, + {"id": 47, "name": "Eve"}, + {"id": 48, "name": "Frank"}, + }, + } + + filePath := filepath.Join(storage.DBDir, "users.db") + file, err := os.Create(filePath) + if err != nil { + t.Fatalf("failed to create file: %v", err) + } + defer file.Close() + + page := pager.NewPage() + for i, row := range table.Rows { + data, err := pager.SerializeRow(row, table.Columns) + if err != nil { + t.Fatalf("failed to serialize row %d: %v", i, err) + } + _, err = page.InsertTuple(data) + if err != nil { + t.Fatalf("failed to insert tuple %d: %v", i, err) + } + } + + pageBytes := pager.SerializePage(page) + if _, err := file.Write(pageBytes); err != nil { + t.Fatalf("failed to write page to file: %v", err) + } + if err := file.Sync(); err != nil { + t.Fatalf("failed to sync file: %v", err) + } +} + +func TestReadFromDisk(t *testing.T) { + filePath := filepath.Join(storage.DBDir, "users.db") + + f, err := os.Open(filePath) + if err != nil { + t.Fatalf("Failed to open file: %v", err) + } + defer f.Close() + + buf := make([]byte, pager.PageSize) + _, err = f.Read(buf) + if err != nil { + t.Fatalf("Failed to read file: %v", err) + } + + page, err := pager.DeserializePage(buf) + if err != nil { + t.Fatalf("Failed to deserialize page: %v", err) + } + + // Hardcoded columns + columns := []db.Column{ + {Name: "id", Type: db.TypeInt}, + {Name: "name", Type: db.TypeString}, + } + + fmt.Println("ROWS FOUND:") + fmt.Printf("Header.NumItems: %d\n", page.Header.NumItems) + + for i := 0; i < int(page.Header.NumItems); i++ { + item := page.Items[i] + if item.DeletedFlag == 0 || item.Length == 0 { + continue + } + + tupleData, err := page.ReadTuple(i) + if err != nil { + t.Errorf("Failed to read tuple %d: %v", i, err) + continue + } + + row, err := pager.DeserializeRow(tupleData, columns) + if err != nil { + t.Errorf("Failed to deserialize tuple %d: %v", i, err) + continue + } + + fmt.Printf("Row %d: %+v\n", i, row) + } +} From 00bcf47d08269bcdb019111b387ae439f717bc83 Mon Sep 17 00:00:00 2001 From: AltSumpreme Date: Mon, 21 Jul 2025 22:43:38 +0530 Subject: [PATCH 2/2] feat(storage): added multi page storage for tables - Added multipage storage for tables - made the storage code more modular - added tests for multipage write --- pager/pager.go | 7 +- storage/Diskloader.go | 85 ++++++++++++++++++++++++ storage/SavetoDisk.go | 71 ++++++++++++++++++++ storage/schemaloader.go | 27 ++++++++ storage/storage.go | 140 ---------------------------------------- tests/write_test.go | 114 +++++++++++--------------------- 6 files changed, 221 insertions(+), 223 deletions(-) create mode 100644 storage/Diskloader.go create mode 100644 storage/SavetoDisk.go create mode 100644 storage/schemaloader.go delete mode 100644 storage/storage.go diff --git a/pager/pager.go b/pager/pager.go index 10a2343..7b60b96 100644 --- a/pager/pager.go +++ b/pager/pager.go @@ -57,26 +57,23 @@ func (p *Page) InsertTuple(record []byte) (int, error) { return -1, errors.New("record cannot be empty") } if p.Header.NumItems >= MaxItemsPerPage { - return -1, errors.New("page is full") + return -2, errors.New("page is full") } - // 🧠 Calculate available free space freeSpace := int(p.Header.PdUpper) - int(p.Header.PdLower) requiredSpace := len(record) + ItemIDSize if requiredSpace > freeSpace { return -1, errors.New("not enough space in page") } - // 👇 Insert the record at the top newUpper := p.Header.PdUpper - uint16(len(record)) copy(p.Data[newUpper:], record) - // 👇 Add the metadata for this tuple slot := int(p.Header.NumItems) p.Items[slot] = ItemID{ Offset: newUpper, Length: uint16(len(record)), - DeletedFlag: 1, + DeletedFlag: 1, // Mark as not deleted } p.Header.PdUpper = newUpper diff --git a/storage/Diskloader.go b/storage/Diskloader.go new file mode 100644 index 0000000..55d00db --- /dev/null +++ b/storage/Diskloader.go @@ -0,0 +1,85 @@ +package storage + +import ( + "fmt" + "os" + "path/filepath" + "pebbledb/db" + "pebbledb/pager" + "strings" +) + +func LoadFromDisk() (*db.Database, error) { + files, err := os.ReadDir(DBDir) + var tableName string + if err != nil { + return nil, fmt.Errorf("failed to read directory %s: %w", DBDir, err) + } + tablesMap := map[string][]string{} + database := db.NewDatabase() + for _, file := range files { + if file.IsDir() { + continue + } + if filepath.Ext(file.Name()) == ".db" { + base_file := strings.TrimSuffix(file.Name(), filepath.Ext(file.Name())) + underscoreIndex := strings.LastIndex(base_file, "_") + if underscoreIndex == -1 { + return nil, fmt.Errorf("invalid file name %s, expected format __.db", file.Name()) + } + tableName = base_file[:underscoreIndex] + tablesMap[tableName] = append(tablesMap[tableName], file.Name()) + + } + } + + for tableName, pageFile := range tablesMap { + columnDef, err := LoadSchemaFromDisk(tableName) + if err != nil { + return nil, fmt.Errorf("failed to load schema for table %s: %w", tableName, err) + } + table := db.Table{ + Name: tableName, + Columns: columnDef, + Rows: []db.Row{}, + } + + for _, files := range pageFile { + path := filepath.Join(DBDir, files) + file, err := os.Open(path) + if err != nil { + return nil, fmt.Errorf("failed to open file %s: %w", path, err) + } + buf := make([]byte, pager.PageSize) + _, err = file.Read(buf) + if err != nil { + return nil, fmt.Errorf("failed to read file %s: %w", path, err) + } + + page, err := pager.DeserializePage(buf) + if err != nil { + return nil, fmt.Errorf("failed to deserialize page from file %s: %w", path, err) + } + for i := 0; i < pager.MaxItemsPerPage; i++ { + if page.Items[i].DeletedFlag == 0 || page.Items[i].Length == 0 { + continue + } + + tupleData, err := page.ReadTuple(i) + if err != nil { + return nil, fmt.Errorf("failed to read tuple %d from page: %w", i, err) + } + deserialized, err := pager.DeserializeRow(tupleData, columnDef) + if err != nil { + return nil, fmt.Errorf("failed to deserialize row from tuple %d: %w", i, err) + } + table.Rows = append(table.Rows, deserialized) + + } + + } + database.Tables[tableName] = &table + + } + return database, nil +} diff --git a/storage/SavetoDisk.go b/storage/SavetoDisk.go new file mode 100644 index 0000000..06471b2 --- /dev/null +++ b/storage/SavetoDisk.go @@ -0,0 +1,71 @@ +package storage + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "pebbledb/db" + "pebbledb/pager" +) + +const DBDir = "./db" + +func SaveToDisk(database *db.Database) error { + if err := os.MkdirAll(DBDir, 0775); err != nil { + return err + } + + for tableName, table := range database.Tables { + metaPath := filepath.Join(DBDir, tableName+".meta.json") + metaFile, err := os.Create(metaPath) + if err != nil { + return fmt.Errorf("create schema file for %s: %w", tableName, err) + } + if err := json.NewEncoder(metaFile).Encode(table.Columns); err != nil { + return fmt.Errorf("encode schema for %s: %w", tableName, err) + } + metaFile.Close() + + page_index := 0 + + currentPage := pager.NewPage() + for _, row := range table.Rows { + fmt.Printf("Serializing row: %v\n", row) + serialized, err := pager.SerializeRow(row, table.Columns) + fmt.Printf("Serialized row: %v\n", serialized) + if err != nil { + return fmt.Errorf("serialize row: %w", err) + } + if slot, err := currentPage.InsertTuple(serialized); err != nil { + if slot == -2 { + pageBytes := pager.SerializePage(currentPage) + pageFileName := fmt.Sprintf("%s_%d.db", tableName, page_index) + filePath := filepath.Join(DBDir, pageFileName) + if err := os.WriteFile(filePath, pageBytes, 0664); err != nil { + return fmt.Errorf("write page %d for %s: %w", page_index, tableName, err) + } + } else { + return fmt.Errorf("insert tuple: %w", err) + } + page_index++ + + currentPage = pager.NewPage() + if _, err := currentPage.InsertTuple(serialized); err != nil { + return fmt.Errorf("insert tuple in new page: %w", err) + } + } + if len(currentPage.Data) > 0 { + pageBytes := pager.SerializePage(currentPage) + pageFileName := fmt.Sprintf("%s_%d.db", tableName, page_index) + filePath := filepath.Join(DBDir, pageFileName) + if err := os.WriteFile(filePath, pageBytes, 0664); err != nil { + return fmt.Errorf("write final page for %s: %w", tableName, err) + } + } + + } + + } + return nil +} diff --git a/storage/schemaloader.go b/storage/schemaloader.go new file mode 100644 index 0000000..831e4ca --- /dev/null +++ b/storage/schemaloader.go @@ -0,0 +1,27 @@ +package storage + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "pebbledb/db" +) + +func LoadSchemaFromDisk(tableName string) ([]db.Column, error) { + metafile := filepath.Join(DBDir, tableName+".meta.json") + if _, err := os.Stat(metafile); os.IsNotExist(err) { + return nil, fmt.Errorf("schema file for table %s does not exist", tableName) + } + file, err := os.Open(metafile) + if err != nil { + return nil, fmt.Errorf("failed to open schema file for table %s: %w", tableName, err) + } + defer file.Close() + + var columnDef []db.Column + if err := json.NewDecoder(file).Decode(&columnDef); err != nil { + return nil, fmt.Errorf("failed to decode schema for table %s: %w", tableName, err) + } + return columnDef, nil +} diff --git a/storage/storage.go b/storage/storage.go deleted file mode 100644 index 1a65236..0000000 --- a/storage/storage.go +++ /dev/null @@ -1,140 +0,0 @@ -package storage - -import ( - "encoding/json" - "fmt" - "os" - "path/filepath" - "pebbledb/db" - "pebbledb/pager" -) - -const DBDir = "./db" - -func SaveToDisk(database *db.Database) error { - if err := os.MkdirAll(DBDir, 0775); err != nil { - return err - } - - for tableName, table := range database.Tables { - metaPath := filepath.Join(DBDir, tableName+".meta.json") - metaFile, err := os.Create(metaPath) - if err != nil { - return fmt.Errorf("create schema file for %s: %w", tableName, err) - } - if err := json.NewEncoder(metaFile).Encode(table.Columns); err != nil { - return fmt.Errorf("encode schema for %s: %w", tableName, err) - } - metaFile.Close() - dataPath := filepath.Join(DBDir, tableName+".db") - dataFile, err := os.Create(dataPath) - if err != nil { - return fmt.Errorf("create data file for %s: %w", tableName, err) - } - - page := pager.NewPage() - for _, row := range table.Rows { - fmt.Printf("Serializing row: %v\n", row) - serialized, err := pager.SerializeRow(row, table.Columns) - fmt.Printf("Serialized row: %v\n", serialized) - if err != nil { - return fmt.Errorf("serialize row: %w", err) - } - if _, err := page.InsertTuple(serialized); err != nil { - return fmt.Errorf("insert tuple: %w", err) - } - } - - if _, err := dataFile.Write(pager.SerializePage(page)); err != nil { - return fmt.Errorf("write page for %s: %w", tableName, err) - } - dataFile.Close() - } - - return nil -} - -func LoadSchemaFromDisk(tableName string) ([]db.Column, error) { - metafile := filepath.Join(DBDir, tableName+".meta.json") - if _, err := os.Stat(metafile); os.IsNotExist(err) { - return nil, fmt.Errorf("schema file for table %s does not exist", tableName) - } - file, err := os.Open(metafile) - if err != nil { - return nil, fmt.Errorf("failed to open schema file for table %s: %w", tableName, err) - } - defer file.Close() - - var columnDef []db.Column - if err := json.NewDecoder(file).Decode(&columnDef); err != nil { - return nil, fmt.Errorf("failed to decode schema for table %s: %w", tableName, err) - } - return columnDef, nil -} -func LoadFromDisk() (*db.Database, error) { - files, err := os.ReadDir(DBDir) - if err != nil { - return nil, fmt.Errorf("failed to read directory %s: %w", DBDir, err) - } - database := db.NewDatabase() - for _, file := range files { - if file.IsDir() { - continue - } - filePath := fmt.Sprintf(DBDir + "/" + file.Name()) - f, err := os.Open(filePath) - if err != nil { - return nil, fmt.Errorf("failed to open file %s: %w", filePath, err) - } - - buf := make([]byte, pager.PageSize) - _, err = f.Read(buf) - f.Close() - if err != nil { - return nil, fmt.Errorf("failed to read file %s: %w", filePath, err) - } - - page, err := pager.DeserializePage(buf) - - if err != nil { - return nil, fmt.Errorf("failed to deserialize page from file %s: %w", filePath, err) - } - if filepath.Ext(filePath) != ".db" { - continue - } - tableName := file.Name() - tableName = tableName[:len(tableName)-3] - columnDef, err := LoadSchemaFromDisk(tableName) - if err != nil { - return nil, fmt.Errorf("table %s does not exist in database: %w", tableName, err) - } - table := &db.Table{ - Name: tableName, - Columns: columnDef, - Rows: []db.Row{}, - } - for i := 0; i < int(page.Header.NumItems); i++ { - items := page.Items[i] - if items.DeletedFlag == 0 || items.Length == 0 { - continue - } - - tupleData, err := page.ReadTuple(i) - - if err != nil { - return nil, fmt.Errorf("failed to read tuple from page: %w", err) - } - row, err := pager.DeserializeRow(tupleData, table.Columns) - if err != nil { - return nil, fmt.Errorf("failed to deserialize row: %w", err) - } - fmt.Printf("Deserialized row: %v\n", row) - table.Rows = append(table.Rows, row) - - } - database.Tables[tableName] = table - - } - return database, nil - -} diff --git a/tests/write_test.go b/tests/write_test.go index 6589e5b..c31499e 100644 --- a/tests/write_test.go +++ b/tests/write_test.go @@ -3,108 +3,66 @@ package tests import ( "fmt" "os" - "path/filepath" "pebbledb/db" - "pebbledb/pager" "pebbledb/storage" + "strconv" "testing" ) -func TestWrite(t *testing.T) { +func TestWriteAndReadMultiPage(t *testing.T) { + _ = os.RemoveAll(storage.DBDir) _ = os.MkdirAll(storage.DBDir, 0775) - table := &db.Table{ - Name: "users", - Columns: []db.Column{ - {Name: "id", Type: db.TypeInt}, - {Name: "name", Type: db.TypeString}, - }, - Rows: []db.Row{ - {"id": 42, "name": "Reuben"}, - {"id": 43, "name": "Alice"}, - {"id": 44, "name": "Bob"}, - {"id": 45, "name": "Charlie"}, - {"id": 46, "name": "Diana"}, - {"id": 47, "name": "Eve"}, - {"id": 48, "name": "Frank"}, - }, - } - - filePath := filepath.Join(storage.DBDir, "users.db") - file, err := os.Create(filePath) + // Initialize and create the table + database := db.NewDatabase() + err := database.CreateTable("users", []db.Column{ + {Name: "id", Type: db.TypeInt}, + {Name: "name", Type: db.TypeString}, + }) if err != nil { - t.Fatalf("failed to create file: %v", err) + t.Fatalf("Failed to create table: %v", err) } - defer file.Close() - page := pager.NewPage() - for i, row := range table.Rows { - data, err := pager.SerializeRow(row, table.Columns) - if err != nil { - t.Fatalf("failed to serialize row %d: %v", i, err) - } - _, err = page.InsertTuple(data) + // Insert 300 rows + for i := 0; i < 300; i++ { + err := database.InsertValue("users", []string{ + strconv.Itoa(i), + "user" + strconv.Itoa(i), + }) if err != nil { - t.Fatalf("failed to insert tuple %d: %v", i, err) + t.Fatalf("Insert failed at row %d: %v", i, err) } } - pageBytes := pager.SerializePage(page) - if _, err := file.Write(pageBytes); err != nil { - t.Fatalf("failed to write page to file: %v", err) + // Save database to disk + if err := storage.SaveToDisk(database); err != nil { + t.Fatalf("Failed to save to disk: %v", err) } - if err := file.Sync(); err != nil { - t.Fatalf("failed to sync file: %v", err) - } -} - -func TestReadFromDisk(t *testing.T) { - filePath := filepath.Join(storage.DBDir, "users.db") - f, err := os.Open(filePath) + // Load the database from disk + loadedDB, err := storage.LoadFromDisk() if err != nil { - t.Fatalf("Failed to open file: %v", err) + t.Fatalf("Failed to load from disk: %v", err) } - defer f.Close() - buf := make([]byte, pager.PageSize) - _, err = f.Read(buf) + // Validate: ensure all 300 rows are retrievable and correct + rows, err := loadedDB.SelectAll("users") if err != nil { - t.Fatalf("Failed to read file: %v", err) + t.Fatalf("Select failed: %v", err) } - page, err := pager.DeserializePage(buf) - if err != nil { - t.Fatalf("Failed to deserialize page: %v", err) + if len(rows) != 300 { + t.Fatalf("Expected 300 rows, got %d", len(rows)) } - // Hardcoded columns - columns := []db.Column{ - {Name: "id", Type: db.TypeInt}, - {Name: "name", Type: db.TypeString}, - } - - fmt.Println("ROWS FOUND:") - fmt.Printf("Header.NumItems: %d\n", page.Header.NumItems) - - for i := 0; i < int(page.Header.NumItems); i++ { - item := page.Items[i] - if item.DeletedFlag == 0 || item.Length == 0 { - continue - } - - tupleData, err := page.ReadTuple(i) - if err != nil { - t.Errorf("Failed to read tuple %d: %v", i, err) - continue + for i := 0; i < 300; i++ { + expectedName := "user" + strconv.Itoa(i) + if rows[i]["name"] != expectedName { + t.Errorf("Row %d: expected name '%s', got '%v'", i, expectedName, rows[i]["name"]) } - - row, err := pager.DeserializeRow(tupleData, columns) - if err != nil { - t.Errorf("Failed to deserialize tuple %d: %v", i, err) - continue - } - - fmt.Printf("Row %d: %+v\n", i, row) } + + // Lookup test + specificRow := rows[257] // Random row for validation + fmt.Printf("Looked up row 257 -> ID: %v, Name: %v\n", specificRow["id"], specificRow["name"]) }