wp
This commit is contained in:
@@ -39,6 +39,7 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type FreeList interface {
|
type FreeList interface {
|
||||||
|
// використовується в allocPage
|
||||||
ReservePage() uint32
|
ReservePage() uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
21
atree/io.go
21
atree/io.go
@@ -8,7 +8,7 @@ import (
|
|||||||
"math"
|
"math"
|
||||||
"os"
|
"os"
|
||||||
|
|
||||||
diploma "gordenko.dev/dima/qb"
|
"gordenko.dev/dima/qb"
|
||||||
"gordenko.dev/dima/qb/atree/redo"
|
"gordenko.dev/dima/qb/atree/redo"
|
||||||
"gordenko.dev/dima/qb/bin"
|
"gordenko.dev/dima/qb/bin"
|
||||||
)
|
)
|
||||||
@@ -74,8 +74,8 @@ func (s *Atree) releaseIndexPage(pageNo uint32) {
|
|||||||
p.ReferenceCount--
|
p.ReferenceCount--
|
||||||
return
|
return
|
||||||
} else {
|
} else {
|
||||||
diploma.Abort(
|
qb.Abort(
|
||||||
diploma.ReferenceCountBug,
|
qb.ReferenceCountBug,
|
||||||
fmt.Errorf("call releaseIndexPage on page %d with reference count = %d",
|
fmt.Errorf("call releaseIndexPage on page %d with reference count = %d",
|
||||||
pageNo, p.ReferenceCount),
|
pageNo, p.ReferenceCount),
|
||||||
)
|
)
|
||||||
@@ -91,15 +91,12 @@ func (s *Atree) allocIndexPage() AllocatedPage {
|
|||||||
)
|
)
|
||||||
|
|
||||||
allocated.PageNo = s.indexFreelist.ReservePage()
|
allocated.PageNo = s.indexFreelist.ReservePage()
|
||||||
|
s.mutex.Lock()
|
||||||
if allocated.PageNo > 0 {
|
if allocated.PageNo > 0 {
|
||||||
allocated.IsReused = true
|
allocated.IsReused = true
|
||||||
|
|
||||||
s.mutex.Lock()
|
|
||||||
} else {
|
} else {
|
||||||
s.mutex.Lock()
|
|
||||||
if s.allocatedIndexPagesQty == math.MaxUint32 {
|
if s.allocatedIndexPagesQty == math.MaxUint32 {
|
||||||
diploma.Abort(diploma.MaxAtreeSizeExceeded,
|
qb.Abort(qb.MaxAtreeSizeExceeded, errors.New("no space in Atree index"))
|
||||||
errors.New("no space in Atree index"))
|
|
||||||
}
|
}
|
||||||
s.allocatedIndexPagesQty++
|
s.allocatedIndexPagesQty++
|
||||||
allocated.PageNo = s.allocatedIndexPagesQty
|
allocated.PageNo = s.allocatedIndexPagesQty
|
||||||
@@ -163,8 +160,8 @@ func (s *Atree) releaseDataPage(pageNo uint32) {
|
|||||||
p.ReferenceCount--
|
p.ReferenceCount--
|
||||||
return
|
return
|
||||||
} else {
|
} else {
|
||||||
diploma.Abort(
|
qb.Abort(
|
||||||
diploma.ReferenceCountBug,
|
qb.ReferenceCountBug,
|
||||||
fmt.Errorf("call releaseDataPage on page %d with reference count = %d",
|
fmt.Errorf("call releaseDataPage on page %d with reference count = %d",
|
||||||
pageNo, p.ReferenceCount),
|
pageNo, p.ReferenceCount),
|
||||||
)
|
)
|
||||||
@@ -186,7 +183,7 @@ func (s *Atree) allocDataPage() AllocatedPage {
|
|||||||
} else {
|
} else {
|
||||||
s.mutex.Lock()
|
s.mutex.Lock()
|
||||||
if s.allocatedDataPagesQty == math.MaxUint32 {
|
if s.allocatedDataPagesQty == math.MaxUint32 {
|
||||||
diploma.Abort(diploma.MaxAtreeSizeExceeded,
|
qb.Abort(qb.MaxAtreeSizeExceeded,
|
||||||
errors.New("no space in Atree index"))
|
errors.New("no space in Atree index"))
|
||||||
}
|
}
|
||||||
s.allocatedDataPagesQty++
|
s.allocatedDataPagesQty++
|
||||||
@@ -303,7 +300,7 @@ func (s *Atree) pageWriter() {
|
|||||||
case <-s.writeSignalCh:
|
case <-s.writeSignalCh:
|
||||||
err := s.writeTasks()
|
err := s.writeTasks()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
diploma.Abort(diploma.WriteToAtreeFailed, err)
|
qb.Abort(qb.WriteToAtreeFailed, err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -431,12 +431,12 @@ func (s *Database) replayChangesRecord(untyped any) error {
|
|||||||
metric.LastPageNo = rec.DataPageNo
|
metric.LastPageNo = rec.DataPageNo
|
||||||
// delete free pages
|
// delete free pages
|
||||||
if rec.IsDataPageReused {
|
if rec.IsDataPageReused {
|
||||||
s.dataFreeList.DeletePages([]uint32{
|
s.dataFreeList.DeleteReservedPages([]uint32{
|
||||||
rec.DataPageNo,
|
rec.DataPageNo,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
if len(rec.ReusedIndexPages) > 0 {
|
if len(rec.ReusedIndexPages) > 0 {
|
||||||
s.indexFreeList.DeletePages(rec.ReusedIndexPages)
|
s.indexFreeList.DeleteReservedPages(rec.ReusedIndexPages)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -535,13 +535,13 @@ func (s *Database) appendMeasureAfterOverflow(extended txlog.AppendedMeasureWith
|
|||||||
metric.LastPageNo = rec.DataPageNo
|
metric.LastPageNo = rec.DataPageNo
|
||||||
|
|
||||||
if rec.IsDataPageReused {
|
if rec.IsDataPageReused {
|
||||||
s.dataFreeList.DeletePages([]uint32{
|
s.dataFreeList.DeleteReservedPages([]uint32{
|
||||||
rec.DataPageNo,
|
rec.DataPageNo,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(rec.ReusedIndexPages) > 0 {
|
if len(rec.ReusedIndexPages) > 0 {
|
||||||
s.indexFreeList.DeletePages(rec.ReusedIndexPages)
|
s.indexFreeList.DeleteReservedPages(rec.ReusedIndexPages)
|
||||||
}
|
}
|
||||||
|
|
||||||
if !extended.HoldLock {
|
if !extended.HoldLock {
|
||||||
|
|||||||
@@ -1,29 +1,46 @@
|
|||||||
package freelist
|
package freelist
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"slices"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/RoaringBitmap/roaring/v2"
|
"gordenko.dev/dima/qb/bin"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
pageSize = 1024
|
||||||
|
crcSize = 4
|
||||||
|
ptrSize = 4
|
||||||
|
pointersOnPage = (pageSize - crcSize) / ptrSize
|
||||||
|
flushTreshold = 1000 //= int(float(pageSize/ptrSize) * 1.3) // кількість вільних сторінок коли вже треба писати на диск
|
||||||
)
|
)
|
||||||
|
|
||||||
type FreeList struct {
|
type FreeList struct {
|
||||||
mutex sync.Mutex
|
mutex sync.Mutex
|
||||||
free *roaring.Bitmap
|
file *os.File
|
||||||
reserved *roaring.Bitmap
|
pages int
|
||||||
|
//free *roaring.Bitmap
|
||||||
|
//reserved *roaring.Bitmap
|
||||||
|
free []uint32
|
||||||
|
reserved []uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
func New() *FreeList {
|
func New() *FreeList {
|
||||||
return &FreeList{
|
return &FreeList{
|
||||||
free: roaring.New(),
|
//free: roaring.New(),
|
||||||
reserved: roaring.New(),
|
//reserved: roaring.New(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *FreeList) Restore(serialized []byte) error {
|
func (s *FreeList) Restore(serialized []byte) error {
|
||||||
err := s.free.UnmarshalBinary(serialized)
|
if (len(serialized) % ptrSize) != 0 {
|
||||||
if err != nil {
|
return fmt.Errorf("wrong size")
|
||||||
return fmt.Errorf("UnmarshalBinary: %s", err)
|
}
|
||||||
|
for i := 0; i < len(serialized); i += ptrSize {
|
||||||
|
s.free = append(s.free, bin.GetUint32(serialized[i:]))
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -33,10 +50,60 @@ func (s *FreeList) AddPages(pageNumbers []uint32) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
s.mutex.Lock()
|
s.mutex.Lock()
|
||||||
s.free.AddMany(pageNumbers)
|
s.free = append(s.free, pageNumbers...)
|
||||||
|
if len(s.free) > flushTreshold {
|
||||||
|
err := s.save()
|
||||||
|
if err != nil {
|
||||||
|
// ABORT
|
||||||
|
}
|
||||||
|
}
|
||||||
s.mutex.Unlock()
|
s.mutex.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *FreeList) save() error {
|
||||||
|
buf := make([]byte, pageSize)
|
||||||
|
//w := bytes.NewBuffer(nil)
|
||||||
|
//s.mutex.Lock()
|
||||||
|
i := crcSize
|
||||||
|
for _, pageNo := range s.free[:pointersOnPage] {
|
||||||
|
bin.PutUint32(buf[i:], pageNo)
|
||||||
|
i += ptrSize
|
||||||
|
}
|
||||||
|
// fix checksum
|
||||||
|
fileOffset := int64(s.pages-1) * pageSize
|
||||||
|
n, err := s.file.WriteAt(buf, fileOffset)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if n != pageSize {
|
||||||
|
return fmt.Errorf("written size %d bytes not equal page size %d", n, pageSize)
|
||||||
|
}
|
||||||
|
s.free = s.free[pointersOnPage:]
|
||||||
|
s.pages++
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *FreeList) loadLastPage() error {
|
||||||
|
var (
|
||||||
|
buf = make([]byte, pageSize)
|
||||||
|
offset = int64(s.pages-1) * pageSize
|
||||||
|
)
|
||||||
|
n, err := s.file.ReadAt(buf, offset)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("file.Seek: %s", err)
|
||||||
|
}
|
||||||
|
if n != s.pages {
|
||||||
|
return fmt.Errorf("read size %d bytes not equal page size %d", n, pageSize)
|
||||||
|
}
|
||||||
|
// check crc
|
||||||
|
//var pageNumbers []uint32
|
||||||
|
for i := crcSize; i < len(buf); i += ptrSize {
|
||||||
|
s.free = append(s.free, bin.GetUint32(buf[i:]))
|
||||||
|
}
|
||||||
|
s.pages--
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// ReserveDataPage - аллокатор резервирует страницу, но не удаляет до визова
|
// ReserveDataPage - аллокатор резервирует страницу, но не удаляет до визова
|
||||||
// DeleteFromFree, ибо транзакция может не завершится, а между віделением страници
|
// DeleteFromFree, ибо транзакция может не завершится, а между віделением страници
|
||||||
// и падением транзакции - будет создан init файл.
|
// и падением транзакции - будет создан init файл.
|
||||||
@@ -44,29 +111,46 @@ func (s *FreeList) ReservePage() (pageNo uint32) {
|
|||||||
s.mutex.Lock()
|
s.mutex.Lock()
|
||||||
defer s.mutex.Unlock()
|
defer s.mutex.Unlock()
|
||||||
|
|
||||||
if s.free.IsEmpty() {
|
if len(s.free) == 0 {
|
||||||
|
if s.pages == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
pageNo = s.free.Minimum()
|
if err := s.loadLastPage(); err != nil {
|
||||||
s.free.Remove(pageNo)
|
// ABORT
|
||||||
s.reserved.Add(pageNo)
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
lastIdx := len(s.free)
|
||||||
|
pageNo = s.free[lastIdx]
|
||||||
|
s.free = s.free[:lastIdx]
|
||||||
|
s.reserved = append(s.reserved, pageNo)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// Удаляет ранее зарезервированные страницы
|
// Удаляет ранее зарезервированные страницы
|
||||||
func (s *FreeList) DeletePages(pageNumbers []uint32) {
|
func (s *FreeList) DeleteReservedPages(pageNumbers []uint32) {
|
||||||
|
var cleared []uint32
|
||||||
s.mutex.Lock()
|
s.mutex.Lock()
|
||||||
for _, pageNo := range pageNumbers {
|
for _, reservedPageNo := range s.reserved {
|
||||||
s.reserved.Remove(pageNo)
|
if slices.Contains(pageNumbers, reservedPageNo) {
|
||||||
s.free.Remove(pageNo) // прокрута TransactionLog
|
// delete
|
||||||
|
} else {
|
||||||
|
cleared = append(cleared, reservedPageNo)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
s.reserved = cleared
|
||||||
s.mutex.Unlock()
|
s.mutex.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *FreeList) Serialize() ([]byte, error) {
|
func (s *FreeList) Serialize() ([]byte, error) {
|
||||||
|
w := bytes.NewBuffer(nil)
|
||||||
s.mutex.Lock()
|
s.mutex.Lock()
|
||||||
defer s.mutex.Unlock()
|
for _, pageNo := range s.free {
|
||||||
tmp := roaring.Or(s.free, s.reserved)
|
bin.WriteUint32(w, pageNo)
|
||||||
tmp.RunOptimize()
|
}
|
||||||
return tmp.ToBytes()
|
for _, pageNo := range s.reserved {
|
||||||
|
bin.WriteUint32(w, pageNo)
|
||||||
|
}
|
||||||
|
s.mutex.Unlock()
|
||||||
|
return w.Bytes(), nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
package diploma
|
package qb
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
Reference in New Issue
Block a user