null inserts, better tuple payload

This commit is contained in:
pokeeffe-molecula 2023-03-30 11:38:30 -05:00
parent 64e864d928
commit 791aa6096c
7 changed files with 563 additions and 405 deletions

View file

@ -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()

View file

@ -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
}

View file

@ -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,

View file

@ -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

View file

@ -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)
}

109
tstore/btreeschema.go Normal file
View file

@ -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
}

View file

@ -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
}