From 791aa6096c97a85815d6d5d7f8eb18fb88a6ca53 Mon Sep 17 00:00:00 2001 From: pokeeffe-molecula Date: Thu, 30 Mar 2023 11:38:30 -0500 Subject: [PATCH] null inserts, better tuple payload --- bufferpool/ondiskdiskmanager.go | 126 +++++++----- bufferpool/page.go | 58 ++++-- holder.go | 6 +- tstore/btree.go | 345 ++++++++------------------------ tstore/btree_test.go | 118 ++++++++--- tstore/btreeschema.go | 109 ++++++++++ tstore/btreetypes.go | 206 +++++++++++++++---- 7 files changed, 563 insertions(+), 405 deletions(-) create mode 100644 tstore/btreeschema.go diff --git a/bufferpool/ondiskdiskmanager.go b/bufferpool/ondiskdiskmanager.go index eb12e1dcd..e1678c9c0 100644 --- a/bufferpool/ondiskdiskmanager.go +++ b/bufferpool/ondiskdiskmanager.go @@ -14,22 +14,22 @@ type FileShardID struct { Shard int32 } -// OnDiskDiskManager is a on disk implementation for a DiskManager interface -type OnDiskDiskManager struct { +// TupleStoreDiskManager is a on disk implementation for a DiskManager interface +type TupleStoreDiskManager struct { mu sync.Mutex files map[FileShardID]*os.File } -// NewInMemDiskSpillingDiskManager returns a in-memory version of disk manager -func NewOnDiskDiskManager() *OnDiskDiskManager { - dm := &OnDiskDiskManager{ +// NewTupleStoreDiskManager returns a disk manager for t-store btrees +func NewTupleStoreDiskManager() *TupleStoreDiskManager { + dm := &TupleStoreDiskManager{ files: make(map[FileShardID]*os.File), } return dm } -func (d *OnDiskDiskManager) CreateOrOpenShard(objectId int32, shard int32, dataFile string) error { +func (d *TupleStoreDiskManager) CreateOrOpenShard(objectId int32, shard int32, dataFile string) error { // serialize access here d.mu.Lock() defer d.mu.Unlock() @@ -53,7 +53,7 @@ func (d *OnDiskDiskManager) CreateOrOpenShard(objectId int32, shard int32, dataF return fmt.Errorf("open file: %w", err) } // write the root page - err = d.writeRootPage(fd, objectId, shard) + err = d.bootstrapTStoreFile(fd, objectId, shard) if err != nil { return err } @@ -69,7 +69,37 @@ func (d *OnDiskDiskManager) CreateOrOpenShard(objectId int32, shard int32, dataF return nil } -func (d *OnDiskDiskManager) writeRootPage(fd *os.File, objectId int32, shard int32) error { +func (d *TupleStoreDiskManager) writePageIDToSlot(page *Page, slotNumber int16, key byte, ptr int64) { + // get the free space offset + freeSpaceOffset := page.ReadFreeSpaceOffset() + keyBytes := []byte{0, 0, 0, key} + // build a chunk + chunk := InternalPageChunk{ + KeyLength: int16(len(keyBytes)), + KeyBytes: keyBytes, + PtrValue: ptr, + } + // compute the new free space offset + freeSpaceOffset -= int16(chunk.Length()) + page.WriteInternalPageChunk(freeSpaceOffset, chunk) + // update the free space offset + page.WriteFreeSpaceOffset(int16(freeSpaceOffset)) + + // make a slot + slot := PageSlot{ + PayloadOffset: freeSpaceOffset, + } + // write the slot + page.WritePageSlot(slotNumber, slot) + + // update the slot count + slotCount := page.ReadSlotCount() + slotCount += 1 + page.WriteSlotCount(slotCount) +} + +func (d *TupleStoreDiskManager) bootstrapTStoreFile(fd *os.File, objectId int32, shard int32) error { + // make a new header page headerPage := NewPage(PageID{objectId, shard, 0}, 0) headerPage.WritePageNumber(0) headerPage.WriteFreeSpaceOffset(int16(PAGE_SIZE)) @@ -77,61 +107,53 @@ func (d *OnDiskDiskManager) writeRootPage(fd *os.File, objectId int32, shard int headerPage.WritePrevPointer(PageID{0, 0, INVALID_PAGE}) headerPage.WritePageType(PAGE_TYPE_BTREE_HEADER) - // write a slot that points to the initial root page - rootPageID := PageID{objectId, shard, 1} + // write the pages ids for the data and schema btrees into the right slots + dataRootPageID := PageID{objectId, shard, 1} + schemaRootPageID := PageID{objectId, shard, 2} - // get the free space offset - freeSpaceOffset := headerPage.ReadFreeSpaceOffset() + d.writePageIDToSlot(headerPage, 0, 0, dataRootPageID.Page) + d.writePageIDToSlot(headerPage, 1, 1, schemaRootPageID.Page) - keyBytes := []byte{0, 0, 0, 0} + fileOffset := int64(0) - // build a chunk - chunk := InternalPageChunk{ - KeyLength: int16(len(keyBytes)), - KeyBytes: keyBytes, - PtrValue: 1, - } - - // compute the new free space offset - freeSpaceOffset -= int16(chunk.Length()) - - headerPage.WriteInternalPageChunk(freeSpaceOffset, chunk) - - // update the free space offset - headerPage.WriteFreeSpaceOffset(int16(freeSpaceOffset)) - - // make a slot - slot := PageSlot{ - PayloadOffset: freeSpaceOffset, - } - // write the slot - headerPage.WritePageSlot(0, slot) - - // update the slot count - headerPage.WriteSlotCount(int16(1)) - - _, err := fd.WriteAt(headerPage.data[:], 0) + // write the header page + _, err := fd.WriteAt(headerPage.data[:], fileOffset) if err != nil { return err } - rootPage := NewPage(rootPageID, 0) - rootPage.WritePageNumber(1) - rootPage.WriteFreeSpaceOffset(int16(PAGE_SIZE)) - rootPage.WriteNextPointer(PageID{0, 0, INVALID_PAGE}) - rootPage.WritePrevPointer(PageID{0, 0, INVALID_PAGE}) - rootPage.WritePageType(PAGE_TYPE_BTREE_LEAF) + // write the data root page + fileOffset += PAGE_SIZE + dataRootPage := NewPage(dataRootPageID, 0) + dataRootPage.WritePageNumber(dataRootPageID.Page) + dataRootPage.WriteFreeSpaceOffset(int16(PAGE_SIZE)) + dataRootPage.WriteNextPointer(PageID{0, 0, INVALID_PAGE}) + dataRootPage.WritePrevPointer(PageID{0, 0, INVALID_PAGE}) + dataRootPage.WritePageType(PAGE_TYPE_BTREE_LEAF) - _, err = fd.WriteAt(rootPage.data[:], PAGE_SIZE) + _, err = fd.WriteAt(dataRootPage.data[:], fileOffset) if err != nil { return err } + // write the schema root page + fileOffset += PAGE_SIZE + schemaRootPage := NewPage(schemaRootPageID, 0) + schemaRootPage.WritePageNumber(dataRootPageID.Page) + schemaRootPage.WriteFreeSpaceOffset(int16(PAGE_SIZE)) + schemaRootPage.WriteNextPointer(PageID{0, 0, INVALID_PAGE}) + schemaRootPage.WritePrevPointer(PageID{0, 0, INVALID_PAGE}) + schemaRootPage.WritePageType(PAGE_TYPE_BTREE_LEAF) + + _, err = fd.WriteAt(schemaRootPage.data[:], fileOffset) + if err != nil { + return err + } return nil } // ReadPage reads a page from disk -func (d *OnDiskDiskManager) ReadPage(pageID PageID) (*Page, error) { +func (d *TupleStoreDiskManager) ReadPage(pageID PageID) (*Page, error) { d.mu.Lock() defer d.mu.Unlock() @@ -173,7 +195,7 @@ func (d *OnDiskDiskManager) ReadPage(pageID PageID) (*Page, error) { } // WritePage writes a page in memory to pages -func (d *OnDiskDiskManager) WritePage(page *Page) error { +func (d *TupleStoreDiskManager) WritePage(page *Page) error { d.mu.Lock() defer d.mu.Unlock() @@ -212,7 +234,7 @@ func (d *OnDiskDiskManager) WritePage(page *Page) error { } // AllocatePage allocates a page and returns the page number -func (d *OnDiskDiskManager) AllocatePage(objectId int32, shard int32) (PageID, error) { +func (d *TupleStoreDiskManager) AllocatePage(objectId int32, shard int32) (PageID, error) { d.mu.Lock() defer d.mu.Unlock() @@ -242,7 +264,7 @@ func (d *OnDiskDiskManager) AllocatePage(objectId int32, shard int32) (PageID, e } // DeallocatePage removes page from disk -func (d *OnDiskDiskManager) DeallocatePage(pageID PageID) error { +func (d *TupleStoreDiskManager) DeallocatePage(pageID PageID) error { d.mu.Lock() defer d.mu.Unlock() @@ -250,12 +272,12 @@ func (d *OnDiskDiskManager) DeallocatePage(pageID PageID) error { return nil } -func (d *OnDiskDiskManager) FileSize(objectId int32, shard int32) int64 { +func (d *TupleStoreDiskManager) FileSize(objectId int32, shard int32) int64 { // TODO(pok) return correct file size return 0 } -func (d *OnDiskDiskManager) Close() { +func (d *TupleStoreDiskManager) Close() { d.mu.Lock() defer d.mu.Unlock() diff --git a/bufferpool/page.go b/bufferpool/page.go index 9cb25a974..2b06fb9ee 100644 --- a/bufferpool/page.go +++ b/bufferpool/page.go @@ -65,10 +65,11 @@ const PAGE_TYPE_HASH_TABLE = 12 // == row payload bytes == // writeTID (int64) // schemaVersion (int16) -// redoPtr (int64) -// fieldOffsets (one for each field, int16, FF is null) +// versionsPtr (int64) +// flags (int8) flags that include a deletion marker +// fieldOffsets (one for each field, int32, FF is null) // offsets point to: -// *fieldData (one for each field) +// fieldData (one for each field) // valueLen (int32) only used for variable length types // valueBytes @@ -385,28 +386,38 @@ func (pg *Page) Dump(label string) { fmt.Printf("%sKEYS: -->\n", fmt.Sprintf("%*s", indent, "")) indent += 4 - // get the keys off the page - keys := make([]int, 0) - pointers := make([]PageID, 0) - iter := NewPageSlotIterator(pg, 0) - for { - ps := iter.Next() - if ps == nil { - break + switch pageType { + case PAGE_TYPE_BTREE_LEAF: + // get the keys off the page + keys := make([]int, 0) + iter := NewPageSlotIterator(pg, 0) + for { + ps := iter.Next() + if ps == nil { + break + } + pl := ps.KeyPayload(pg) + keys = append(keys, int(pl.KeyAsInt(pg))) } - pl := ps.KeyPayload(pg) - keys = append(keys, int(pl.KeyAsInt(pg))) - if pageType == PAGE_TYPE_BTREE_INTERNAL { - ipl := ps.InternalPayload(pg) - pointers = append(pointers, ipl.ValueAsPagePointer(pg)) - } - } - if pageType == PAGE_TYPE_BTREE_LEAF { for _, key := range keys { fmt.Printf("%s(%d)\n", fmt.Sprintf("%*s", indent, ""), key) } - } else { + case PAGE_TYPE_BTREE_INTERNAL, PAGE_TYPE_BTREE_HEADER: + // get the keys off the page + keys := make([]int, 0) + pointers := make([]PageID, 0) + iter := NewPageSlotIterator(pg, 0) + for { + ps := iter.Next() + if ps == nil { + break + } + pl := ps.KeyPayload(pg) + keys = append(keys, int(pl.KeyAsInt(pg))) + ipl := ps.InternalPayload(pg) + pointers = append(pointers, ipl.ValueAsPagePointer(pg)) + } for idx, key := range keys { ptr := pointers[idx] fmt.Printf("%s(%d, %d)\n", fmt.Sprintf("%*s", indent, ""), key, ptr) @@ -414,7 +425,6 @@ func (pg *Page) Dump(label string) { ptr := pg.ReadNextPointer() fmt.Printf("%s(-->, %d)\n", fmt.Sprintf("%*s", indent, ""), ptr) } - } type KeyPayload struct { @@ -490,7 +500,7 @@ func (l *LeafPayload) ValueLength(page *Page) int32 { return valueLen } -// this wil fail in overflow +// only call this if you are sure this there is no overflow func (l *LeafPayload) ValueAsBytes(page *Page) []byte { offset := l.valueOffset(page) valueLen := int32(binary.BigEndian.Uint32(page.data[offset:])) @@ -535,6 +545,10 @@ func (l *LeafPayload) GetPayloadReader(page *Page) LeafPagePayLoadReader { } } +func (l *LeafPayload) IsVisibleToTID(page *Page, tid int64) bool { + return true +} + type PageSlot struct { PayloadOffset int16 } diff --git a/holder.go b/holder.go index 856b1259c..9016b9c73 100644 --- a/holder.go +++ b/holder.go @@ -151,7 +151,7 @@ type Holder struct { // t-store tstorepool *bufferpool.BufferPool - tstoredisk *bufferpool.OnDiskDiskManager + tstoredisk *bufferpool.TupleStoreDiskManager } // HolderOpts holds information about the holder which other things might want @@ -264,7 +264,7 @@ type HolderConfig struct { CacheFlushInterval time.Duration Logger logger.Logger TStoreBufferPool *bufferpool.BufferPool - TStoreDiskManager *bufferpool.OnDiskDiskManager + TStoreDiskManager *bufferpool.TupleStoreDiskManager StorageConfig *storage.Config RBFConfig *rbfcfg.Config @@ -277,7 +277,7 @@ type HolderConfig struct { // need to override these; that's usually handled by server options // such as OptServerOpenTranslateStore. func DefaultHolderConfig() *HolderConfig { - dm := bufferpool.NewOnDiskDiskManager() + dm := bufferpool.NewTupleStoreDiskManager() return &HolderConfig{ PartitionN: disco.DefaultPartitionN, OpenTranslateStore: OpenInMemTranslateStore, diff --git a/tstore/btree.go b/tstore/btree.go index ea012cddc..ec6e6008f 100644 --- a/tstore/btree.go +++ b/tstore/btree.go @@ -3,36 +3,19 @@ package tstore import ( - "bytes" "fmt" "sync" "github.com/featurebasedb/featurebase/v3/bufferpool" "github.com/featurebasedb/featurebase/v3/sql3/parser" "github.com/featurebasedb/featurebase/v3/sql3/planner/types" - "github.com/featurebasedb/featurebase/v3/wireprotocol" "github.com/pkg/errors" ) -// TODO(pok) -// ▶ on open index run the aries recovery - -// ▶ put the insert into the split routine, so we don't have the do the insert after the fact -// ▶ add a BTreeFile struct to wrap n (data, schema) BTrees -// ▶ move bootstrap into BTreeFile ? - -// - [x] Tuple format -// writeTID (int64) -// schemaVersion (int16) (will we ever have > 64K of schema versions??) -// versionsPtr (int64) ptr to old versions page for this tuple -// fieldOffsets (one for each field, int16, 0xFF is null) -// offsets point to: -// fieldData (one for each field) -// [valueLen (int32)] optional only used for variable length types (right now varchar) -// valueBytes +// TODO... // ▶ WAL -// ▶ implement log buffers for wal files +// ▶ implement log buffers for wal files in bufferpool // in the buffer pool keep some buffers // ▶ implement writing of log // log every change to a page, before we write to the page @@ -41,25 +24,35 @@ import ( // - TID // - record type // - PageID -// - old data (undo) (+ offset and len) +// - old data (undo) (+ offset and len) where applicable // - new data (redo) (+ offset and len) +// ▶ on open featurebase index run aries recovery (https://blog.acolyer.org/2016/01/08/aries/) + // ▶ lazy writer on the buffer pool // ▶ implement a checkpoint that scans the pool and writes out dirty pages periodically // ▶ MVCC versioning -// ▶ handle nulls on insert +// ▶ handle updates (delete and an insert) + +// ▶ handle deletes + // ▶ backup/restore -// later +// Later.... +// ▶ put the insert into the split routine, so we don't have the do the insert after the fact // ▶ do a test on concurrent inserts -// ▶ latch buffer I/Os +// ▶ smarter latching of buffer I/Os (currently serialized) + +type TID int64 const ( - SLOT_DATA_ROOT = 1 - SLOT_SCHEMA_ROOT = 2 + // the slot on the header page that points to the root page for data + HEADER_SLOT_DATA_ROOT = 0 + // the slot on the header page that points to the root page for schema versions + HEADER_SLOT_SCHEMA_ROOT = 1 ) // BTree represents a b+tree structure used for storing and @@ -84,9 +77,6 @@ type BTree struct { rootNode bufferpool.PageID schemaNode bufferpool.PageID bufferpool *bufferpool.BufferPool - - // debug - // pinnedPages map[bufferpool.PageID]bufferpool.PageID } // NewBTree creates a b+tree for a schema object denoted by objectID, and a shard denoted by shard. The maxKeySize specifies the @@ -103,21 +93,24 @@ func NewBTree(maxKeySize int, objectID int32, shard int32, schema types.Schema, sizeWithoutHeader := bufferpool.PAGE_SIZE - bufferpool.PAGE_SLOTS_START_OFFSET - // for internal pages chunk size is: + // for internal pages size on page per key is: // keyLength 2 // keyBytes maxKeySize // ptrValue 8 - chunkSize := int64(2 + maxKeySize + 8 + bufferpool.PAGE_SLOT_LENGTH) + chunkSize := int64( /* keyLength */ 2 + maxKeySize + /* ptrValue */ 8 + bufferpool.PAGE_SLOT_LENGTH) keysPerInternalPage := sizeWithoutHeader / (bufferpool.PAGE_SLOT_LENGTH + chunkSize) - // for leaf pages chunk size is: + // for leaf pages size on page (worst case) is: // keyLength 2 // keyBytes maxKeySize - // rowPayloadLen 4 + // flags 1 + // overflowPtr 8 + // rowPayloadTotalLen 4 + // rowPayloadChunkLen 2 // rowPayloadBytes [payload length] // get the payload length from the schema - payLoadLength := 0 + payLoadLength := /* writeTID */ 8 + /* schemaVersion */ 2 + /* versionsPtr */ 8 + /* flags */ 1 for _, s := range schema { switch ty := s.Type.(type) { case *parser.DataTypeVarchar: @@ -128,7 +121,7 @@ func NewBTree(maxKeySize int, objectID int32, shard int32, schema types.Schema, } } - chunkSize = int64(2 + maxKeySize + 1 + 2 + payLoadLength) + chunkSize = int64( /* keyLength*/ 2 + maxKeySize + /* flags */ 1 + /* overflowPtr */ 8 + /* rowPayloadTotalLen */ 8 + /* rowPayloadChunkLen */ 2 + payLoadLength) keysPerLeafPage := int64(4) if payLoadLength <= bufferpool.MAX_PAYLOAD_CHUNK_SIZE_ON_PAGE { keysPerLeafPage = sizeWithoutHeader / (bufferpool.PAGE_SLOT_LENGTH + chunkSize) @@ -140,143 +133,33 @@ func NewBTree(maxKeySize int, objectID int32, shard int32, schema types.Schema, bufferpool: bpool, keysPerLeafPage: keysPerLeafPage, keysPerInternalPage: keysPerInternalPage, - // debug - // pinnedPages: make(map[bufferpool.PageID]bufferpool.PageID), } headerNode, err := tree.fetchNode(bufferpool.PageID{ObjectID: objectID, Shard: shard, Page: 0}) if err != nil { return nil, err } - headerNode.takeWriteLatch() - defer headerNode.releaseWriteLatch() + headerNode.takeReadLatch() + defer headerNode.releaseReadLatch() defer tree.unpin(headerNode) + // sanity check slotCount := headerNode.page.ReadSlotCount() - if slotCount == 0 { + if slotCount != 2 { panic("inavlid slot count") } // get the root node for the data part of the b-tree - slot := headerNode.page.ReadPageSlot(0) + slot := headerNode.page.ReadPageSlot(HEADER_SLOT_DATA_ROOT) ipl := slot.InternalPayload(headerNode.page) - // no need to protect updating rootNode with a RWMutex while we are still creating tree.rootNode = ipl.ValueAsPagePointer(headerNode.page) - // make the root for the schema part of the b-tree if it does not exist - if slotCount == 1 { - schemaNode, err := tree.newLeaf() - if err != nil { - return nil, err - } - keyBytes := []byte{0, 0, 0, 1} - tree.writeInternalEntryInSlot(headerNode, 1, keyBytes, schemaNode.page.ID().Page) - headerNode.page.WriteSlotCount(2) - tree.unpin(schemaNode) - tree.schemaNode = schemaNode.page.ID() - } else { - slot := headerNode.page.ReadPageSlot(1) - ipl := slot.InternalPayload(headerNode.page) - tree.schemaNode = ipl.ValueAsPagePointer(headerNode.page) - } + slot = headerNode.page.ReadPageSlot(HEADER_SLOT_SCHEMA_ROOT) + ipl = slot.InternalPayload(headerNode.page) + tree.schemaNode = ipl.ValueAsPagePointer(headerNode.page) - // get the last schema version + tree.resolveSchema(schema) - schemaSchema := types.Schema{ - &types.PlannerColumn{ - ColumnName: "schema", - Type: parser.NewDataTypeVarbinary(16384), - }, - } - - insertSchemaSchema := types.Schema{ - &types.PlannerColumn{ - ColumnName: "_id", - Type: parser.NewDataTypeID(), - }, - &types.PlannerColumn{ - ColumnName: "schema", - Type: parser.NewDataTypeVarbinary(16384), - }, - } - - schemaNode, err := tree.fetchNode(tree.schemaNode) - if err != nil { - return nil, err - } - schemaNode.takeReadLatch() - defer schemaNode.releaseAnyLatch() - defer tree.unpin(schemaNode) - - // get an iterator on the schema and get the last row - iter := tree.getIterator(schemaNode, nil, nil, true, schemaSchema) - sk, ss, err := iter.Next() - if err != nil { - return nil, err - } - iter.Dispose() - - // if there isn't any schema, insert this one as the last - if sk == nil { - // make a tuple - b, err := wireprotocol.WriteSchema(schema) - if err != nil { - return nil, err - } - - tup := &BTreeTuple{ - TupleSchema: insertSchemaSchema, - Tuple: types.Row{ - int64(1), - b, - }, - } - - err = tree.insert(SLOT_SCHEMA_ROOT, schemaNode, tup) - if err != nil { - return nil, err - } - tree.schemaVersion = 1 - tree.schema = schema - // TODO(pok) - can take this out once we have logging and checkpoints - tree.bufferpool.FlushPage(schemaNode.page.ID()) - } else { - // if there is one, see if it is the same as this one, if it is not, insert this one as the latest - latestVersion := int(sk.(Int)) - - b := ss.Tuple[0].([]byte) - - rdr := bytes.NewReader(b) - _, err = wireprotocol.ExpectToken(rdr, wireprotocol.TOKEN_SCHEMA_INFO) - if err != nil { - return nil, err - } - - s, err := wireprotocol.ReadSchema(rdr) - if err != nil { - return nil, err - } - - newSchema := false - if len(s) != len(schema) { - newSchema = true - } else { - for i, c := range s { - if schema[i].ColumnName != c.ColumnName || schema[i].Type.BaseTypeName() != c.Type.BaseTypeName() { - newSchema = true - break - } - } - } - if newSchema { - // TODO(pok) add another schema version - return nil, errors.Errorf("schema mismatch") - } - tree.schemaVersion = latestVersion - tree.schema = schema - - } - tree.bufferpool.FlushPage(headerNode.page.ID()) return tree, nil } @@ -293,13 +176,14 @@ func (b *BTree) GetIterator(from Sortable, to Sortable, reverse bool) (*BTreeNod // Search finds a key k in the b+tree and returns the matching tuple. If no tuple is // found, nil is returned -func (b *BTree) Search(currentNode *BTreeNode, k Sortable) (Sortable, *BTreeTuple, error) { +func (b *BTree) Search(k Sortable) (Sortable, *BTreeTuple, error) { + // go get the root data page from the buffer pool n, err := b.fetchNode(b.rootNode) if err != nil { return nil, nil, err } n.takeReadLatch() - return b.search(n, k) + return b.search(n, b.schema, k) } // Insert inserts a tuple into the b+tree. The key is assumed to be in the first column of the tuple. @@ -312,9 +196,10 @@ func (b *BTree) Insert(tup *BTreeTuple) error { if err != nil { return err } - return b.insert(SLOT_DATA_ROOT, node, tup) + return b.insert(HEADER_SLOT_DATA_ROOT, node, tup, b.schema, b.schemaVersion) } +// ==================================================================================================== // private methods func (b *BTree) getIterator(node *BTreeNode, from Sortable, to Sortable, reverse bool, schema types.Schema) *BTreeNodeIterator { @@ -335,14 +220,14 @@ func (b *BTree) getIterator(node *BTreeNode, from Sortable, to Sortable, reverse // ▶ acquire read latch on child // ▶ then unlatch parent -func (b *BTree) search(currentNode *BTreeNode, k Sortable) (Sortable, *BTreeTuple, error) { +func (b *BTree) search(currentNode *BTreeNode, schema types.Schema, k Sortable) (Sortable, *BTreeTuple, error) { if currentNode.isLeaf() { defer currentNode.releaseReadLatch() defer b.unpin(currentNode) // search for the key i, found := currentNode.findKey(k) if found { - return b.getTuple(currentNode, i, b.schema) + return b.getTuple(currentNode, i, schema) } return nil, nil, nil } else { @@ -354,11 +239,11 @@ func (b *BTree) search(currentNode *BTreeNode, k Sortable) (Sortable, *BTreeTupl node.takeReadLatch() currentNode.releaseReadLatch() b.unpin(currentNode) - return b.search(node, k) + return b.search(node, schema, k) } } -func (b *BTree) insert(rootSlot int, node *BTreeNode, tup *BTreeTuple) error { +func (b *BTree) insert(headerPageRootSlot int, node *BTreeNode, tup *BTreeTuple, schema types.Schema, schemaVersion int) error { key := tup.keyValue() // special handling for the case where root is a leaf @@ -408,20 +293,20 @@ func (b *BTree) insert(rootSlot int, node *BTreeNode, tup *BTreeTuple) error { } // do the insert into the node - err = b.insertNonFull(n, key, tup, forceExclusive) + err = b.insertNonFull(n, key, tup, forceExclusive, schema, schemaVersion) if err != nil { return err } // this is the root node splitting so handle that... - err = b.handleRootNodeSplit(rootSlot, pivot, lhsPtr, rhsPtr) + err = b.handleRootNodeSplit(headerPageRootSlot, pivot, lhsPtr, rhsPtr) if err != nil { return err } return nil } else { - err := b.insertNonFull(node, key, tup, forceExclusive) + err := b.insertNonFull(node, key, tup, forceExclusive, schema, schemaVersion) if err != nil { if err == ErrNeedsExclusive { // release the read latch & take a write latch @@ -438,14 +323,25 @@ func (b *BTree) insert(rootSlot int, node *BTreeNode, tup *BTreeTuple) error { } } +// getTuple reads a tuple off a page (and overflow pages). It will read the header of the tuple and decide (TODO) whether +// the tuple should be visible based on transaction ids or whether it is deleted func (b *BTree) getTuple(node *BTreeNode, slotNumber int, schema types.Schema) (Sortable, *BTreeTuple, error) { + // get the slot slot := node.page.ReadPageSlot(int16(slotNumber)) + + // get the chunk off this leaf page kpl := slot.KeyPayload(node.page) lpl := slot.LeafPayload(node.page) rdr := lpl.GetPayloadReader(node.page) + + _ /*tupleHdr*/ = NewBTreeTupleHeaderFromBytes(rdr.PayloadChunkBytes) + + // TODO(pok) see if this tuple is visible to this TID, and is not deleted etc. + + // make a buffer to fit the payload payload := make([]byte, rdr.PayloadTotalLength) copy(payload, rdr.PayloadChunkBytes) - bytesReceived := rdr.PayloadChunkLength + bytesReceived := int32(rdr.PayloadChunkLength) if rdr.Flags == 1 { nextPtr := rdr.OverflowPtr for nextPtr != bufferpool.INVALID_PAGE { @@ -457,7 +353,7 @@ func (b *BTree) getTuple(node *BTreeNode, slotNumber int, schema types.Schema) ( // read the overflow bytes clen, cbytes := onode.page.ReadLeafPagePayloadBytes(bufferpool.PAGE_SLOTS_START_OFFSET) copy(payload[bytesReceived:], cbytes) - bytesReceived += clen + bytesReceived += int32(clen) nextPtr = onode.page.ReadNextPointer().Page @@ -467,7 +363,6 @@ func (b *BTree) getTuple(node *BTreeNode, slotNumber int, schema types.Schema) ( } } return Int(kpl.KeyAsInt(node.page)), NewBTreeTupleFromBytes(payload, schema), nil - } // fetchNode gets the page specified by pageID from the buffer pool and wraps @@ -481,27 +376,11 @@ func (b *BTree) fetchNode(pageID bufferpool.PageID) (*BTreeNode, error) { node := &BTreeNode{ page: page, } - //debug - // if page.ID().Page == 200 { - // fmt.Printf("here\n") - // } - // _, ok := b.pinnedPages[pageID] - // if ok { - // fmt.Printf("pinning page (already pinned) %v\n", pageID) - // } else { - // fmt.Printf("pinning page %v\n", pageID) - // b.pinnedPages[pageID] = pageID - // } - //--- return node, nil } // unpin is a convenience method to unpin the page wrapped by node func (b *BTree) unpin(node *BTreeNode) error { - // debug - // fmt.Printf("unpinning page %v\n", node.page.ID()) - // delete(b.pinnedPages, node.page.ID()) - //--- return b.bufferpool.UnpinPage(node.page.ID()) } @@ -517,17 +396,6 @@ func (b *BTree) newOverflow() (*BTreeNode, error) { node := &BTreeNode{ page: page, } - - //debug - // _, ok := b.pinnedPages[page.ID()] - // if ok { - // fmt.Printf("pinning page (already pinned) %v\n", page.ID()) - // } else { - // fmt.Printf("pinning page %v\n", page.ID()) - // b.pinnedPages[page.ID()] = page.ID() - // } - //--- - return node, nil } @@ -543,15 +411,6 @@ func (b *BTree) newLeaf() (*BTreeNode, error) { node := &BTreeNode{ page: page, } - //debug - // _, ok := b.pinnedPages[page.ID()] - // if ok { - // fmt.Printf("pinning page (already pinned) %v\n", page.ID()) - // } else { - // fmt.Printf("pinning page %v\n", page.ID()) - // b.pinnedPages[page.ID()] = page.ID() - // } - //--- return node, nil } @@ -567,15 +426,6 @@ func (b *BTree) newInternal() (*BTreeNode, error) { node := &BTreeNode{ page: page, } - //debug - // _, ok := b.pinnedPages[page.ID()] - // if ok { - // fmt.Printf("pinning page (already pinned) %v\n", page.ID()) - // } else { - // fmt.Printf("pinning page %v\n", page.ID()) - // b.pinnedPages[page.ID()] = page.ID() - // } - //--- return node, nil } @@ -639,7 +489,7 @@ func (b *BTree) findNextPointer(node *BTreeNode, key Sortable) bufferpool.PageID // setRootNode writes the new root node page id into the appropriate slot // in the header page of the b+tree, flushes the page and then updates the // rootNode member on the BTree struct instance pointed to by b -func (b *BTree) setRootNode(rootSlot int, newRootNode bufferpool.PageID) error { +func (b *BTree) setRootNode(headerPageRootSlot int, newRootNode bufferpool.PageID) error { b.mu.Lock() defer b.mu.Unlock() @@ -652,20 +502,22 @@ func (b *BTree) setRootNode(rootSlot int, newRootNode bufferpool.PageID) error { defer headerNode.releaseWriteLatch() defer b.unpin(headerNode) - // root page pointer is in rootSlot - 1 - slot := headerNode.page.ReadPageSlot(int16(rootSlot)) + // root page pointer is on the header page in headerPageRootSlot + slot := headerNode.page.ReadPageSlot(int16(headerPageRootSlot)) ipl := slot.InternalPayload(headerNode.page) // update the root page pointer ipl.PutPagePointer(headerNode.page, newRootNode) + // TODO(pok) remove this once we have WAL b.bufferpool.FlushPage(headerNode.page.ID()) - switch rootSlot { - case SLOT_DATA_ROOT: + // depending on which slot we have, update the cached page ptrs for root pages + switch headerPageRootSlot { + case HEADER_SLOT_DATA_ROOT: b.rootNode = newRootNode - case SLOT_SCHEMA_ROOT: + case HEADER_SLOT_SCHEMA_ROOT: b.schemaNode = newRootNode } return nil @@ -673,7 +525,7 @@ func (b *BTree) setRootNode(rootSlot int, newRootNode bufferpool.PageID) error { // handleRootNodeSplit creates a new root node, inserts the seperator key pivot to point to the lhsPtr and // the next pointer to point to the rhsPtr -func (b *BTree) handleRootNodeSplit(rootSlot int, pivot Sortable, lhsPtr bufferpool.PageID, rhsPtr bufferpool.PageID) error { +func (b *BTree) handleRootNodeSplit(headerPageRootSlot int, pivot Sortable, lhsPtr bufferpool.PageID, rhsPtr bufferpool.PageID) error { // create a new root node newRoot, err := b.newInternal() if err != nil { @@ -689,7 +541,7 @@ func (b *BTree) handleRootNodeSplit(rootSlot int, pivot Sortable, lhsPtr bufferp // set the next ptr to point to newNode newRoot.page.WriteNextPointer(rhsPtr) - b.setRootNode(rootSlot, newRoot.page.ID()) + b.setRootNode(headerPageRootSlot, newRoot.page.ID()) return nil } @@ -754,10 +606,6 @@ func (b *BTree) compactInternalPage(node *BTreeNode) error { // 2. iterate the slots on this page and copy data // 3. put the scratch page back over this one - // debug - // node.page.Dump("pre compact") - // -- - scratchPage := b.bufferpool.ScratchPage() scratchPage.WritePageType(bufferpool.PAGE_TYPE_BTREE_INTERNAL) freeSpaceOffset := scratchPage.ReadFreeSpaceOffset() @@ -794,10 +642,6 @@ func (b *BTree) compactInternalPage(node *BTreeNode) error { scratchPage.CopyPageTo(node.page) - // debug - // node.page.Dump("post compact") - // --- - return nil } @@ -834,9 +678,6 @@ func (b *BTree) writeLeafEntryInSlot(node *BTreeNode, slotNumber int16, keyBytes if err != nil { return err } - // debug - // node.page.Dump("after compact") - // --- // check again if (onPageSize + bufferpool.PAGE_SLOT_LENGTH) > int32(node.page.FreeSpaceOnPage()) { // this shouldn't happen - if it does we have a logic error @@ -932,7 +773,7 @@ func (b *BTree) writeLeafEntryInSlot(node *BTreeNode, slotNumber int16, keyBytes } // insertLeafEntryAt inserts a tuple at the slot denoted by keyPosition -func (b *BTree) insertLeafEntryAt(node *BTreeNode, keyPosition int, key Sortable, tup *BTreeTuple) error { +func (b *BTree) insertLeafEntryAt(node *BTreeNode, keyPosition int, key Sortable, tup *BTreeTuple, schema types.Schema, schemaVersion int) error { if node.latchState() != bufferpool.Write { panic("unexpected latch state") } @@ -941,7 +782,9 @@ func (b *BTree) insertLeafEntryAt(node *BTreeNode, keyPosition int, key Sortable slotCount := int(node.page.ReadSlotCount()) // put the payload together - valueData, err := tup.Bytes() + // TODO(pok) add TID and deal with existing rows (i.e. moving old versions to overflow pages) + // TODO(pok) also should writers do the garbage collection for old rows that from commited transactions? + valueData, err := tup.Bytes(0, schema, schemaVersion, bufferpool.INVALID_PAGE, false) if err != nil { return err } @@ -954,7 +797,6 @@ func (b *BTree) insertLeafEntryAt(node *BTreeNode, keyPosition int, key Sortable // update the slot count slotCount++ node.page.WriteSlotCount(int16(slotCount)) - return nil } @@ -974,9 +816,6 @@ func (b *BTree) writeInternalEntryInSlot(node *BTreeNode, slotNumber int16, keyB if err != nil { return err } - // debug - // node.page.Dump("after compact") - // --- // check again if (onPageSize + bufferpool.PAGE_SLOT_LENGTH) > int32(node.page.FreeSpaceOnPage()) { // this shouldn't happen - if it does we have a logic error @@ -1064,15 +903,8 @@ func (b *BTree) splitNode(nodeToSplit *BTreeNode) (*BTreeNode, Sortable, *BTreeN // • A safe node is one that will not split or merge when updated. // ▶ Not full (on insertion) -func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, forceExclusive bool) error { +func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, forceExclusive bool, schema types.Schema, schemaVersion int) error { if node.isLeaf() { - // debug - // fmt.Printf("leaf insert on page %v (%v)\n", node.page.ID(), tup) - - // if key == Int(289) { - // node.page.Dump("200") - // } - // --- if node.latchState() != bufferpool.Write { panic("unexpected latch state") @@ -1086,16 +918,11 @@ func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, fo return errors.Errorf("key violation") } - err := b.insertLeafEntryAt(node, i, key, tup) + err := b.insertLeafEntryAt(node, i, key, tup, schema, schemaVersion) if err != nil { return err } - // debug - // if node.page.ID().Page == 200 { - // node.page.Dump("200") - // } - // return nil } else { // its an internal node so follow the pointers @@ -1108,8 +935,6 @@ func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, fo return err } - // fmt.Printf("internal node search on page %v: key=%v, childptr=%v\n", node.page.ID(), key, childPtr) - // we need to latch child node if forceExclusive { childNode.takeWriteLatch() @@ -1169,11 +994,11 @@ func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, fo if key.Less(pivot) { rhs.releaseWriteLatch() b.unpin(rhs) - return b.insertNonFull(lhs, key, tup, forceExclusive) + return b.insertNonFull(lhs, key, tup, forceExclusive, schema, schemaVersion) } else { lhs.releaseWriteLatch() b.unpin(lhs) - return b.insertNonFull(rhs, key, tup, forceExclusive) + return b.insertNonFull(rhs, key, tup, forceExclusive, schema, schemaVersion) } } else { // given the child node was not full, we can unlatch the parent here @@ -1181,7 +1006,7 @@ func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, fo b.unpin(node) // now do the insert on the child node - return b.insertNonFull(childNode, key, tup, forceExclusive) + return b.insertNonFull(childNode, key, tup, forceExclusive, schema, schemaVersion) } } } @@ -1189,7 +1014,6 @@ func (b *BTree) insertNonFull(node *BTreeNode, key Sortable, tup *BTreeTuple, fo // splitLeafNode handles splitting of a leaf node. It returns the original node, the seperator key on // which the split happended and the new node, or an error. func (b *BTree) splitLeafNode(nodeToSplit *BTreeNode) (*BTreeNode, Sortable, *BTreeNode, error) { - // fmt.Printf("leaf split on page %v\n", nodeToSplit.page.ID()) // where we are going to split slotCount := int(nodeToSplit.page.ReadSlotCount()) @@ -1248,21 +1072,12 @@ func (b *BTree) splitLeafNode(nodeToSplit *BTreeNode) (*BTreeNode, Sortable, *BT lhs.page.WriteNextPointer(rhs.page.ID()) rhs.page.WritePrevPointer(lhs.page.ID()) - // debug - // if lhs.page.ID().Page == 200 { - // lhs.page.Dump("200:lhs") - // rhs.page.Dump("200:rhs") - // } - // --- - return lhs, seperationKey, rhs, nil } // splitInternalNode handles splitting of an internal node. It returns the original node, the seperator key on // which the split happended and the new node, or an error. func (b *BTree) splitInternalNode(nodeToSplit *BTreeNode) (*BTreeNode, Sortable, *BTreeNode, error) { - // fmt.Printf("internal split on page %v\n", nodeToSplit.page.ID()) - // where we are going to split slotCount := int(nodeToSplit.page.ReadSlotCount()) splitPoint := slotCount / 2 diff --git a/tstore/btree_test.go b/tstore/btree_test.go index 07ddda329..86c7f7d6b 100644 --- a/tstore/btree_test.go +++ b/tstore/btree_test.go @@ -13,7 +13,7 @@ import ( ) func TestAddItemsToBTreeAndValidate(t *testing.T) { - diskManager := bufferpool.NewOnDiskDiskManager() + diskManager := bufferpool.NewTupleStoreDiskManager() objectId := int32(1) shard := int32(0) @@ -77,24 +77,7 @@ func TestAddItemsToBTreeAndValidate(t *testing.T) { duration := time.Since(start) fmt.Printf("inserted %d rows in %v\n", 3000, duration) - start = time.Now() - key, tuple, _ := b.Search(nil, Int(33)) - duration = time.Since(start) - - vals := "[" - for i, v := range tuple.Tuple { - if i > 10 { - vals += "..." - break - } - if i != 0 { - vals += ", " - } - vals += fmt.Sprintf("%v", v) - } - vals += "]" - - fmt.Printf("retrieved key %v, tuple (%d columns), %s in %v\n", key, len(tuple.TupleSchema), vals, duration) + getKey(b, Int(33)) fmt.Printf("\n\n") @@ -102,7 +85,7 @@ func TestAddItemsToBTreeAndValidate(t *testing.T) { } func TestAddItemsToBTreeAndValidate_VeryWide(t *testing.T) { - diskManager := bufferpool.NewOnDiskDiskManager() + diskManager := bufferpool.NewTupleStoreDiskManager() objectId := int32(1) shard := int32(0) @@ -178,9 +161,95 @@ func TestAddItemsToBTreeAndValidate_VeryWide(t *testing.T) { duration := time.Since(start) fmt.Printf("inserted %d rows in %v\n", numRecs, duration) - start = time.Now() - key, tuple, _ := b.Search(nil, Int(524)) - duration = time.Since(start) + getKey(b, Int(524)) + + fmt.Printf("\n\n") + + b.Dump(0) +} + +func TestAddItemsWithNullsToBTreeAndValidate(t *testing.T) { + diskManager := bufferpool.NewTupleStoreDiskManager() + + objectId := int32(1) + shard := int32(0) + dataFile := fmt.Sprintf("ts-shard.%04d", shard) + + os.Remove(dataFile) + + diskManager.CreateOrOpenShard(objectId, shard, dataFile) + + bufferPool := bufferpool.NewBufferPool(100, diskManager) + + tableSchema := types.Schema{ + &types.PlannerColumn{ + ColumnName: "vtest", + Type: parser.NewDataTypeVarchar(50), + }, + } + + b, err := NewBTree(8, objectId, shard, tableSchema, bufferPool) + if err != nil { + t.Fatal(err) + } + + rowSchema := types.Schema{ + &types.PlannerColumn{ + ColumnName: "_id", + Type: parser.NewDataTypeID(), + }, + &types.PlannerColumn{ + ColumnName: "vtest", + Type: parser.NewDataTypeVarchar(50), + }, + } + + inserts := make([]int, 0) + for i := 1; i <= 50; i++ { + inserts = append(inserts, i) + } + rand.Shuffle(len(inserts), func(i, j int) { inserts[i], inserts[j] = inserts[j], inserts[i] }) + + rr := make(types.Row, 2) + + start := time.Now() + for _, i := range inserts { + rr[0] = int64(i) + if i%2 == 0 { + rr[1] = fmt.Sprintf("This is a test of things %d", i) + } else { + rr[1] = nil + } + + tup := &BTreeTuple{ + TupleSchema: rowSchema, + Tuple: rr, + } + + // fmt.Printf("%v", tup) + + err = b.Insert(tup) + if err != nil { + t.Fatal(err) + } + } + + duration := time.Since(start) + fmt.Printf("inserted %d rows in %v\n", 50, duration) + + getKey(b, Int(32)) + fmt.Printf("\n\n") + + getKey(b, Int(33)) + fmt.Printf("\n\n") + + b.Dump(0) +} + +func getKey(b *BTree, k Sortable) { + start := time.Now() + key, tuple, _ := b.Search(k) + duration := time.Since(start) vals := "[" for i, v := range tuple.Tuple { @@ -197,7 +266,4 @@ func TestAddItemsToBTreeAndValidate_VeryWide(t *testing.T) { fmt.Printf("retrieved key %v, tuple (%d columns), %s in %v\n", key, len(tuple.TupleSchema), vals, duration) - fmt.Printf("\n\n") - - b.Dump(0) } diff --git a/tstore/btreeschema.go b/tstore/btreeschema.go new file mode 100644 index 000000000..928e47aba --- /dev/null +++ b/tstore/btreeschema.go @@ -0,0 +1,109 @@ +// Copyright 2023 Molecula Corp. All rights reserved. + +package tstore + +import ( + "bytes" + + "github.com/featurebasedb/featurebase/v3/sql3/parser" + "github.com/featurebasedb/featurebase/v3/sql3/planner/types" + "github.com/featurebasedb/featurebase/v3/wireprotocol" + "github.com/pkg/errors" +) + +func (b *BTree) resolveSchema(schema types.Schema) error { + + // get the last schema version + schemaSchema := types.Schema{ + &types.PlannerColumn{ + ColumnName: "schema", + Type: parser.NewDataTypeVarbinary(16384), + }, + } + + insertSchemaSchema := types.Schema{ + &types.PlannerColumn{ + ColumnName: "_id", + Type: parser.NewDataTypeID(), + }, + &types.PlannerColumn{ + ColumnName: "schema", + Type: parser.NewDataTypeVarbinary(16384), + }, + } + + schemaNode, err := b.fetchNode(b.schemaNode) + if err != nil { + return err + } + schemaNode.takeReadLatch() + defer schemaNode.releaseAnyLatch() + defer b.unpin(schemaNode) + + // get an iterator on the schema and get the last row + iter := b.getIterator(schemaNode, nil, nil, true, schemaSchema) + sk, ss, err := iter.Next() + if err != nil { + return err + } + iter.Dispose() + + // if there isn't any schema, insert this one as the last + if sk == nil { + // make a tuple + schemaBytes, err := wireprotocol.WriteSchema(schema) + if err != nil { + return err + } + + tup := &BTreeTuple{ + TupleSchema: insertSchemaSchema, + Tuple: types.Row{ + int64(1), + schemaBytes, + }, + } + + err = b.insert(HEADER_SLOT_SCHEMA_ROOT, schemaNode, tup, schemaSchema, 1) + if err != nil { + return err + } + b.schemaVersion = 1 + b.schema = schema + } else { + // if there is one, see if it is the same as this one, if it is not, insert this one as the latest + latestVersion := int(sk.(Int)) + + schemaBytes := ss.Tuple[0].([]byte) + + rdr := bytes.NewReader(schemaBytes) + _, err = wireprotocol.ExpectToken(rdr, wireprotocol.TOKEN_SCHEMA_INFO) + if err != nil { + return err + } + + s, err := wireprotocol.ReadSchema(rdr) + if err != nil { + return err + } + + newSchema := false + if len(s) != len(schema) { + newSchema = true + } else { + for i, c := range s { + if schema[i].ColumnName != c.ColumnName || schema[i].Type.BaseTypeName() != c.Type.BaseTypeName() { + newSchema = true + break + } + } + } + if newSchema { + // TODO(pok) add another schema version + return errors.Errorf("schema mismatch") + } + b.schemaVersion = latestVersion + b.schema = schema + } + return nil +} diff --git a/tstore/btreetypes.go b/tstore/btreetypes.go index f5c5e0151..54049dd79 100644 --- a/tstore/btreetypes.go +++ b/tstore/btreetypes.go @@ -5,6 +5,7 @@ package tstore import ( "bytes" "encoding/binary" + "io" "github.com/pkg/errors" @@ -18,9 +19,41 @@ const ( var ErrNeedsExclusive = errors.New("ErrNeedsExclusive") +type BTreeTupleHeader struct { + writeTID TID + schemaVersion int16 + versionsPtr int64 + flags int8 +} + +func NewBTreeTupleHeaderFromBytes(b []byte) *BTreeTupleHeader { + rdr := bytes.NewReader(b) + + var writeTID int64 + binary.Read(rdr, binary.BigEndian, &writeTID) + + var schemaVersion int16 + binary.Read(rdr, binary.BigEndian, &schemaVersion) + + var versionsPtr int64 + binary.Read(rdr, binary.BigEndian, &versionsPtr) + + var flags int8 + binary.Read(rdr, binary.BigEndian, &flags) + + return &BTreeTupleHeader{ + writeTID: TID(writeTID), + schemaVersion: schemaVersion, + versionsPtr: versionsPtr, + flags: flags, + } +} + type BTreeTuple struct { TupleSchema types.Schema Tuple types.Row + + tupleSchemaIndex map[string]int } func (t *BTreeTuple) keyValue() Sortable { @@ -31,6 +64,21 @@ func (t *BTreeTuple) keyValue() Sortable { return Int(kv) } +func (t *BTreeTuple) column(name string) (int, *types.PlannerColumn) { + if t.tupleSchemaIndex == nil { + t.tupleSchemaIndex = make(map[string]int) + for i, s := range t.TupleSchema { + t.tupleSchemaIndex[s.ColumnName] = i + } + } + idx, ok := t.tupleSchemaIndex[name] + if !ok { + return -1, nil + } + ps := t.TupleSchema[idx] + return idx, ps +} + func NewBTreeTupleFromBytes(b []byte, schema types.Schema) *BTreeTuple { t := &BTreeTuple{ TupleSchema: schema, @@ -38,21 +86,45 @@ func NewBTreeTupleFromBytes(b []byte, schema types.Schema) *BTreeTuple { } rdr := bytes.NewReader(b) + var writeTID int64 + binary.Read(rdr, binary.BigEndian, &writeTID) + + var schemaVersion int16 + binary.Read(rdr, binary.BigEndian, &schemaVersion) + + var versionsPtr int64 + binary.Read(rdr, binary.BigEndian, &versionsPtr) + + var flags int8 + binary.Read(rdr, binary.BigEndian, &flags) + + dataRdr := bytes.NewReader(b) + for i, s := range schema { + // read the offset + var fieldOffset uint32 + binary.Read(rdr, binary.BigEndian, &fieldOffset) + if fieldOffset == 0xFFFFFFFF { + t.Tuple[i] = nil + continue + } + switch s.Type.(type) { case *parser.DataTypeVarchar: + dataRdr.Seek(int64(fieldOffset), io.SeekStart) var l int32 - binary.Read(rdr, binary.BigEndian, &l) + binary.Read(dataRdr, binary.BigEndian, &l) bvalue := make([]byte, l) - binary.Read(rdr, binary.BigEndian, &bvalue) + binary.Read(dataRdr, binary.BigEndian, &bvalue) t.Tuple[i] = string(bvalue) case *parser.DataTypeVarbinary: + dataRdr.Seek(int64(fieldOffset), io.SeekStart) var l int32 - binary.Read(rdr, binary.BigEndian, &l) + binary.Read(dataRdr, binary.BigEndian, &l) bvalue := make([]byte, l) - binary.Read(rdr, binary.BigEndian, &bvalue) - t.Tuple[i] = []byte(bvalue) + binary.Read(dataRdr, binary.BigEndian, &bvalue) + t.Tuple[i] = bvalue default: panic("unexpected type") } @@ -60,48 +132,108 @@ func NewBTreeTupleFromBytes(b []byte, schema types.Schema) *BTreeTuple { return t } -func (b *BTreeTuple) Bytes() ([]byte, error) { +// tuple format +// writeTID (int64) +// schemaVersion (int16) (will we ever have > 64K of schema versions??) +// versionsPtr (int64) ptr to old versions page for this tuple +// flags (int8) flags that include a deletion marker +// fieldOffsets (one for each field, int32, 0xFFFFFFFF is null) +// offsets point to: +// fieldData (one for each field) +// [valueLen (int32)] optional only used for variable length types (right now varchar) +// valueBytes + +func (b *BTreeTuple) Bytes(tid TID, schema types.Schema, schemaVersion int, versionsPtr int64, forDeletion bool) ([]byte, error) { var valueBuf bytes.Buffer - for cidx, c := range b.TupleSchema { - // skip the key column - if cidx == 0 { + + hdrBytes := make([]byte, 8+2+8+1) + offset := 0 + + // write TID + binary.BigEndian.PutUint64(hdrBytes[offset:], uint64(tid)) + offset += 8 + + // schemaVersion + binary.BigEndian.PutUint16(hdrBytes[offset:], uint16(schemaVersion)) + offset += 2 + + // versionsPtr + binary.BigEndian.PutUint64(hdrBytes[offset:], uint64(versionsPtr)) + offset += 8 + + // flags + hdrBytes[offset] = 0 + offset += 1 + + // write header + valueBuf.Write(hdrBytes) + + // hang on to a null value + nullValue := make([]byte, 4) + binary.BigEndian.PutUint32(nullValue, uint32(0xFFFFFFFF)) + + // initialize offset to be header + field array + offset += len(schema) * 4 + + // now write offsets and field data + var fieldData bytes.Buffer + for _, sc := range schema { + cidx, tsc := b.column(sc.ColumnName) + // if the schema column is not in the schema being used for insert we write a null + if tsc == nil { + valueBuf.Write(nullValue) continue } + rd := b.Tuple[cidx] - switch ty := c.Type.(type) { + // if we get an explict null write it + if rd == nil { + valueBuf.Write(nullValue) + continue + } + + switch tsc.Type.(type) { case *parser.DataTypeVarchar: - if rd == nil { - b := []byte{0, 0, 0, 0} - valueBuf.Write(b) - } else { - data, ok := rd.(string) - if !ok { - return []byte{}, errors.Errorf("unexpected type conversion '%T'", rd) - } - b := make([]byte, 4) - binary.BigEndian.PutUint32(b, uint32(len(data))) - valueBuf.Write(b) - valueBuf.WriteString(data) - } + // write field offset + b := make([]byte, 4) + binary.BigEndian.PutUint32(b, uint32(offset)) + valueBuf.Write(b) + + // write field data + data := rd.(string) + b = make([]byte, 4) + l := len(data) + offset += 4 + l + binary.BigEndian.PutUint32(b, uint32(l)) + fieldData.Write(b) + fieldData.WriteString(data) + case *parser.DataTypeVarbinary: - if rd == nil { - b := []byte{0, 0, 0, 0} - valueBuf.Write(b) - } else { - data, ok := rd.([]byte) - if !ok { - return []byte{}, errors.Errorf("unexpected type conversion '%T'", rd) - } - b := make([]byte, 4) - binary.BigEndian.PutUint32(b, uint32(len(data))) - valueBuf.Write(b) - valueBuf.Write(data) - } + // write field offset + b := make([]byte, 4) + binary.BigEndian.PutUint32(b, uint32(offset)) + valueBuf.Write(b) + + // write field data + data := rd.([]byte) + b = make([]byte, 4) + l := len(data) + offset += 4 + l + binary.BigEndian.PutUint32(b, uint32(l)) + binary.BigEndian.PutUint32(b, uint32(len(data))) + fieldData.Write(b) + fieldData.Write(data) + default: - return []byte{}, errors.Errorf("unexpected type '%T'", ty) + panic("unexpected field type") } } + + // write the field data + valueBuf.Write(fieldData.Bytes()) + + // and we're done return valueBuf.Bytes(), nil }