From e7effb07ac36e1a9926598c8dfeaaa3d96043c6a Mon Sep 17 00:00:00 2001 From: pokeeffe-molecula Date: Mon, 27 Mar 2023 18:47:26 -0500 Subject: [PATCH] Now handling schema in it's own btree --- index.go | 3 +- sql3/parser/astdatatype.go | 25 +++ tstore/btree.go | 294 ++++++++++++++++++++++----------- tstore/btree_test.go | 4 +- tstore/btreenode.go | 52 ++++++ tstore/btreetypes.go | 21 +++ wireprotocol/wireprimitives.go | 42 +++++ 7 files changed, 345 insertions(+), 96 deletions(-) diff --git a/index.go b/index.go index 308278f4c..57732a2a8 100644 --- a/index.go +++ b/index.go @@ -116,7 +116,8 @@ func (i *Index) DataframesPath() string { // TStorePath returns the path of the t-store files specific to an index func (i *Index) TStorePath() string { - return filepath.Join(i.path, TStoreDir) + path := backendsDir + sep + TStoreDir + return filepath.Join(i.path, path) } // Name returns name of the index. diff --git a/sql3/parser/astdatatype.go b/sql3/parser/astdatatype.go index 87ee634cf..ff0ebf6dd 100644 --- a/sql3/parser/astdatatype.go +++ b/sql3/parser/astdatatype.go @@ -56,6 +56,7 @@ func (*DataTypeStringSet) exprDataType() {} func (*DataTypeStringSetQuantum) exprDataType() {} func (*DataTypeTimestamp) exprDataType() {} func (*DataTypeVarchar) exprDataType() {} +func (*DataTypeVarbinary) exprDataType() {} type DataTypeVoid struct { } @@ -228,6 +229,30 @@ func (d *DataTypeVarchar) TypeInfo() map[string]interface{} { } } +type DataTypeVarbinary struct { + Length int64 +} + +func NewDataTypeVarbinary(length int64) *DataTypeVarbinary { + return &DataTypeVarbinary{ + Length: length, + } +} + +func (d *DataTypeVarbinary) BaseTypeName() string { + return "varbinary" // dax.BaseTypeVarchar +} + +func (d *DataTypeVarbinary) TypeDescription() string { + return fmt.Sprintf("%s(%d)", "varbinary" /*dax.BaseTypeVarchar*/, d.Length) +} + +func (d *DataTypeVarbinary) TypeInfo() map[string]interface{} { + return map[string]interface{}{ + "length": d.Length, + } +} + type DataTypeID struct { } diff --git a/tstore/btree.go b/tstore/btree.go index 494c9543a..cc8e24159 100644 --- a/tstore/btree.go +++ b/tstore/btree.go @@ -67,16 +67,20 @@ import ( // ▶ do a test on concurrent inserts // ▶ latch buffer I/Os +const ( + SLOT_DATA_ROOT = 1 + SLOT_SCHEMA_ROOT = 2 +) + // BTree represents a b+tree structure used for storing and // retrieving tuple data for a given shard in a table // // BTrees for FeatureBase t-store data contain a header page -// TODO (pok) - the following decribes aspirational state, vs. actual state right now -// - slot 1 of the header page will be a pointer to the root page for the schema data +// - slot 1 will be the pointer to the root page of the tuple data that is being stored in +// this b-tree +// - slot 2 of the header page will be a pointer to the root page for the schema data // the schema is stored as a tuple of schema version (int), and schema ([]byte) // like any other tuple -// - slot 2 will be the pointer to the root page of the tuple data that is being stored in -// this b-tree type BTree struct { mu sync.RWMutex schema types.Schema @@ -88,6 +92,7 @@ type BTree struct { objectID int32 shard int32 rootNode bufferpool.PageID + schemaNode bufferpool.PageID bufferpool *bufferpool.BufferPool // debug @@ -157,29 +162,98 @@ func NewBTree(maxKeySize int, objectID int32, shard int32, schema types.Schema, defer headerNode.releaseWriteLatch() defer tree.unpin(headerNode) + slotCount := headerNode.page.ReadSlotCount() + if slotCount == 0 { + panic("inavlid slot count") + } + + // get the root node for the data part of the b-tree slot := headerNode.page.ReadPageSlot(0) - // no need to protect updating this with a RWMutex yet 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) - // see if the headerPage is in overflow - nextHeader := headerNode.page.ReadNextPointer() - if nextHeader.Page != bufferpool.INVALID_PAGE { - // right now overflow is an error - return nil, errors.Errorf("headerPage overflow") + // 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) } - // get the slot count from the headerPage - slotCount := headerNode.page.ReadSlotCount() - if slotCount > 1 { - // we have schema versions - slot = headerNode.page.ReadPageSlot(slotCount - 1) - pl := slot.KeyPayload(headerNode.page) - lpl := slot.LeafPayload(headerNode.page) + // get the last schema version - latestVersion := int(pl.KeyAsInt(headerNode.page)) + schemaSchema := types.Schema{ + &types.PlannerColumn{ + ColumnName: "schema", + Type: parser.NewDataTypeVarbinary(16384), + }, + } - b := lpl.ValueAsBytes(headerNode.page) + 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 + // 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) @@ -204,84 +278,37 @@ func NewBTree(maxKeySize int, objectID int32, shard int32, schema types.Schema, } } if newSchema { + // TODO(pok) add another schema version return nil, errors.Errorf("schema mismatch") } tree.schemaVersion = latestVersion tree.schema = schema - } else { - // no schema versions so write first one - b, err := wireprotocol.WriteSchema(schema) - if err != nil { - return nil, err - } - err = tree.writeLeafEntryInSlot(headerNode, 1, Int(1).Bytes(), b) - if err != nil { - return nil, err - } - headerNode.page.WriteSlotCount(2) - - tree.schemaVersion = 1 - tree.schema = schema - tree.bufferpool.FlushPage(headerNode.page.ID()) } + tree.bufferpool.FlushPage(headerNode.page.ID()) return tree, nil } -// Latching for Search --> start at root with a read latch and go down; repeatedly, -// ▶ acquire read latch on child -// ▶ then unlatch parent +// GetIterator constructs an iterator on the b+tree so that range scans may be performed or an error +func (b *BTree) GetIterator(from Sortable, to Sortable, reverse bool) (*BTreeNodeIterator, error) { + n, err := b.fetchNode(b.rootNode) + if err != nil { + return nil, err + } + n.takeReadLatch() + iter := b.getIterator(n, from, to, reverse, b.schema) + return iter, nil +} // Search finds a key k in the b+tree and returns the matching tuple. If no tuple is // found, nil is returned -// TODO(pok) - remove the first parameter for the public method -// TODO(pok) - add ability to return an error -func (b *BTree) Search(currentNode *BTreeNode, k Sortable) (Sortable, *BTreeTuple) { - if currentNode != nil { - if currentNode.isLeaf() { - defer currentNode.releaseReadLatch() - defer b.unpin(currentNode) - // search for the key - i, found := currentNode.findKey(k) - if found { - slot := currentNode.page.ReadPageSlot(int16(i)) - lpl := slot.LeafPayload(currentNode.page) - rdr := lpl.GetPayloadReader(currentNode.page) - payload := make([]byte, rdr.PayloadTotalLength) - copy(payload, rdr.PayloadChunkBytes) - bytesReceived := rdr.PayloadChunkLength - if rdr.Flags == 1 { - nextPtr := rdr.OverflowPtr - for nextPtr != bufferpool.INVALID_PAGE { - onode, _ := b.fetchNode(bufferpool.PageID{ObjectID: b.objectID, Shard: b.shard, Page: nextPtr}) - onode.takeReadLatch() - defer onode.releaseReadLatch() - defer b.unpin(onode) - - // read the overflow bytes - clen, cbytes := onode.page.ReadLeafPagePayloadBytes(bufferpool.PAGE_SLOTS_START_OFFSET) - copy(payload[bytesReceived:], cbytes) - bytesReceived += clen - - nextPtr = onode.page.ReadNextPointer().Page - } - } - return k, NewBTreeTupleFromBytes(payload, b.schema) - } - return nil, nil - } else { - nodePtr := b.findNextPointer(currentNode, k) - node, _ := b.fetchNode(nodePtr) - node.takeReadLatch() - currentNode.releaseReadLatch() - b.unpin(currentNode) - return b.Search(node, k) - } - } else { - n, _ := b.fetchNode(b.rootNode) - n.takeReadLatch() - return b.Search(n, k) +func (b *BTree) Search(currentNode *BTreeNode, k Sortable) (Sortable, *BTreeTuple, error) { + n, err := b.fetchNode(b.rootNode) + if err != nil { + return nil, nil, err } + n.takeReadLatch() + return b.search(n, k) } // Insert inserts a tuple into the b+tree. The key is assumed to be in the first column of the tuple. @@ -289,12 +316,58 @@ func (b *BTree) Search(currentNode *BTreeNode, k Sortable) (Sortable, *BTreeTupl // // We use an optimistic latching model for inserts. (see insertNonFull() for more details) func (b *BTree) Insert(tup *BTreeTuple) error { - // go get the root page from the buffer pool + // go get the root data page from the buffer pool node, err := b.fetchNode(b.rootNode) if err != nil { return err } + return b.insert(SLOT_DATA_ROOT, node, tup) +} +// private methods + +func (b *BTree) getIterator(node *BTreeNode, from Sortable, to Sortable, reverse bool, schema types.Schema) *BTreeNodeIterator { + if reverse { + // go to the right + if node.isLeaf() { + return NewBTreeNodeIterator(b, node, reverse, schema) + } else { + panic("implement me") + } + } else { + // go to the left + panic("implement me") + } +} + +// Latching for Search --> start at root with a read latch and go down; repeatedly, +// ▶ acquire read latch on child +// ▶ then unlatch parent + +func (b *BTree) search(currentNode *BTreeNode, 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 nil, nil, nil + } else { + nodePtr := b.findNextPointer(currentNode, k) + node, err := b.fetchNode(nodePtr) + if err != nil { + return nil, nil, err + } + node.takeReadLatch() + currentNode.releaseReadLatch() + b.unpin(currentNode) + return b.search(node, k) + } +} + +func (b *BTree) insert(rootSlot int, node *BTreeNode, tup *BTreeTuple) error { key := tup.keyValue() // special handling for the case where root is a leaf @@ -350,7 +423,7 @@ func (b *BTree) Insert(tup *BTreeTuple) error { } // this is the root node splitting so handle that... - err = b.handleRootNodeSplit(pivot, lhsPtr, rhsPtr) + err = b.handleRootNodeSplit(rootSlot, pivot, lhsPtr, rhsPtr) if err != nil { return err } @@ -374,7 +447,37 @@ func (b *BTree) Insert(tup *BTreeTuple) error { } } -// private methods +func (b *BTree) getTuple(node *BTreeNode, slotNumber int, schema types.Schema) (Sortable, *BTreeTuple, error) { + slot := node.page.ReadPageSlot(int16(slotNumber)) + kpl := slot.KeyPayload(node.page) + lpl := slot.LeafPayload(node.page) + rdr := lpl.GetPayloadReader(node.page) + payload := make([]byte, rdr.PayloadTotalLength) + copy(payload, rdr.PayloadChunkBytes) + bytesReceived := rdr.PayloadChunkLength + if rdr.Flags == 1 { + nextPtr := rdr.OverflowPtr + for nextPtr != bufferpool.INVALID_PAGE { + onode, err := b.fetchNode(bufferpool.PageID{ObjectID: b.objectID, Shard: b.shard, Page: nextPtr}) + if err != nil { + return nil, nil, err + } + onode.takeReadLatch() + // read the overflow bytes + clen, cbytes := onode.page.ReadLeafPagePayloadBytes(bufferpool.PAGE_SLOTS_START_OFFSET) + copy(payload[bytesReceived:], cbytes) + bytesReceived += clen + + nextPtr = onode.page.ReadNextPointer().Page + + // release the latch and unpin + onode.releaseReadLatch() + b.unpin(onode) + } + } + return Int(kpl.KeyAsInt(node.page)), NewBTreeTupleFromBytes(payload, schema), nil + +} // fetchNode gets the page specified by pageID from the buffer pool and wraps // it in a BTreeNode. The page is pinned and read to use. It is the callers @@ -545,7 +648,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(newRootNode bufferpool.PageID) error { +func (b *BTree) setRootNode(rootSlot int, newRootNode bufferpool.PageID) error { b.mu.Lock() defer b.mu.Unlock() @@ -558,8 +661,8 @@ func (b *BTree) setRootNode(newRootNode bufferpool.PageID) error { defer headerNode.releaseWriteLatch() defer b.unpin(headerNode) - // root page pointer is in slot 0 - slot := headerNode.page.ReadPageSlot(0) + // root page pointer is in rootSlot - 1 + slot := headerNode.page.ReadPageSlot(int16(rootSlot)) ipl := slot.InternalPayload(headerNode.page) @@ -568,13 +671,18 @@ func (b *BTree) setRootNode(newRootNode bufferpool.PageID) error { b.bufferpool.FlushPage(headerNode.page.ID()) - b.rootNode = newRootNode + switch rootSlot { + case SLOT_DATA_ROOT: + b.rootNode = newRootNode + case SLOT_SCHEMA_ROOT: + b.schemaNode = newRootNode + } return nil } // 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(pivot Sortable, lhsPtr bufferpool.PageID, rhsPtr bufferpool.PageID) error { +func (b *BTree) handleRootNodeSplit(rootSlot int, pivot Sortable, lhsPtr bufferpool.PageID, rhsPtr bufferpool.PageID) error { // create a new root node newRoot, err := b.newInternal() if err != nil { @@ -590,7 +698,7 @@ func (b *BTree) handleRootNodeSplit(pivot Sortable, lhsPtr bufferpool.PageID, rh // set the next ptr to point to newNode newRoot.page.WriteNextPointer(rhsPtr) - b.setRootNode(newRoot.page.ID()) + b.setRootNode(rootSlot, newRoot.page.ID()) return nil } diff --git a/tstore/btree_test.go b/tstore/btree_test.go index 229e91944..445f5b98f 100644 --- a/tstore/btree_test.go +++ b/tstore/btree_test.go @@ -72,7 +72,7 @@ func TestAddItemsToBTreeAndValidate(t *testing.T) { } } - key, tuple := b.Search(nil, Int(33)) + key, tuple, _ := b.Search(nil, Int(33)) fmt.Printf("%v, %v\n\n", key, tuple) @@ -157,7 +157,7 @@ func TestAddItemsToBTreeAndValidate_VeryWide(t *testing.T) { fmt.Printf("inserted %d rows in %v\n", numRecs, duration) start = time.Now() - key, tuple := b.Search(nil, Int(524)) + key, tuple, _ := b.Search(nil, Int(524)) duration = time.Since(start) vals := "[" diff --git a/tstore/btreenode.go b/tstore/btreenode.go index a899af0dc..2506024b5 100644 --- a/tstore/btreenode.go +++ b/tstore/btreenode.go @@ -4,6 +4,7 @@ package tstore import ( "github.com/featurebasedb/featurebase/v3/bufferpool" + "github.com/featurebasedb/featurebase/v3/sql3/planner/types" ) type BTreeNode struct { @@ -89,3 +90,54 @@ func (n *BTreeNode) findNextPointer(key Sortable, objectID int32, shard int32) ( ipl := slot.InternalPayload(n.page) return ipl.ValueAsPagePointer(n.page), nil } + +type BTreeNodeIterator struct { + tree *BTree + node *BTreeNode + schema types.Schema + reverse bool + cursor int16 +} + +func NewBTreeNodeIterator(tree *BTree, initialNode *BTreeNode, reverse bool, schema types.Schema) *BTreeNodeIterator { + return &BTreeNodeIterator{ + tree: tree, + node: initialNode, + schema: schema, + reverse: reverse, + cursor: 0, + } +} + +func (i *BTreeNodeIterator) init() { + if i.reverse { + slotCount := i.node.slotCount() + if slotCount > 0 { + i.cursor = int16(slotCount) + } else { + i.cursor = 0 + } + } else { + panic("implement me") + } +} + +func (i *BTreeNodeIterator) Next() (Sortable, *BTreeTuple, error) { + if i.cursor == 0 { + i.init() + } + if i.reverse { + if i.cursor == 0 { + return nil, nil, nil + } + ci := i.cursor + i.cursor -= 1 + return i.tree.getTuple(i.node, int(ci-1), i.schema) + } else { + panic("implement me") + } +} + +func (i *BTreeNodeIterator) Dispose() { + i.node.releaseReadLatch() +} diff --git a/tstore/btreetypes.go b/tstore/btreetypes.go index f445d92cd..f5c5e0151 100644 --- a/tstore/btreetypes.go +++ b/tstore/btreetypes.go @@ -46,6 +46,13 @@ func NewBTreeTupleFromBytes(b []byte, schema types.Schema) *BTreeTuple { bvalue := make([]byte, l) binary.Read(rdr, binary.BigEndian, &bvalue) t.Tuple[i] = string(bvalue) + + case *parser.DataTypeVarbinary: + var l int32 + binary.Read(rdr, binary.BigEndian, &l) + bvalue := make([]byte, l) + binary.Read(rdr, binary.BigEndian, &bvalue) + t.Tuple[i] = []byte(bvalue) default: panic("unexpected type") } @@ -77,6 +84,20 @@ func (b *BTreeTuple) Bytes() ([]byte, error) { valueBuf.Write(b) valueBuf.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) + } default: return []byte{}, errors.Errorf("unexpected type '%T'", ty) } diff --git a/wireprotocol/wireprimitives.go b/wireprotocol/wireprimitives.go index eb96d0116..8591cfaac 100644 --- a/wireprotocol/wireprimitives.go +++ b/wireprotocol/wireprimitives.go @@ -37,6 +37,7 @@ const ( TYPE_STRING int8 = 0x07 TYPE_STRINGSET int8 = 0x08 TYPE_VARCHAR int8 = 0x09 + TYPE_VARBINARY int8 = 0xA ) func ExpectToken(reader io.Reader, token int16) (int16, error) { @@ -123,6 +124,10 @@ func WriteSchema(schema types.Schema) ([]byte, error) { writeInt8(writer, TYPE_VARCHAR) writeInt32(writer, int32(ty.Length)) + case *parser.DataTypeVarbinary: + writeInt8(writer, TYPE_VARBINARY) + writeInt32(writer, int32(ty.Length)) + default: return []byte{}, errors.Errorf("unexpected type '%T'", s.Type) } @@ -201,6 +206,14 @@ func ReadSchema(reader io.Reader) (types.Schema, error) { return nil, err } dataType = parser.NewDataTypeVarchar(int64(length)) + + case TYPE_VARBINARY: + var length int32 + err = binary.Read(reader, binary.BigEndian, &length) + if err != nil { + return nil, err + } + dataType = parser.NewDataTypeVarbinary(int64(length)) } schema = append(schema, &types.PlannerColumn{ @@ -372,6 +385,18 @@ func WriteRow(row types.Row, schema types.Schema) ([]byte, error) { writer.WriteString(v) } + case *parser.DataTypeVarbinary: + if val == nil { + writeInt32(writer, 0) + } else { + v, ok := row[i].([]byte) + if !ok { + return []byte{}, errors.Errorf("unexpected type '%T'", row[i]) + } + writeInt32(writer, int32(len(v))) + writer.Write(v) + } + default: return []byte{}, errors.Errorf("unexpected type '%T'", s.Type) } @@ -534,6 +559,23 @@ func ReadRow(reader io.Reader, schema types.Schema) (types.Row, error) { row[idx] = string(bvalue) } + case *parser.DataTypeVarbinary: + var len int32 + err := binary.Read(reader, binary.BigEndian, &len) + if err != nil { + return nil, err + } + if len == 0 { + row[idx] = nil + } else { + bvalue := make([]byte, len) + err = binary.Read(reader, binary.BigEndian, &bvalue) + if err != nil { + return nil, err + } + row[idx] = bvalue + } + default: return nil, errors.Errorf("unexpected type '%T'", s.Type) }