Compare commits
9 Commits
Author | SHA1 | Date | |
---|---|---|---|
|
91b2ba30f6 | ||
|
c2828592ac | ||
|
13e53f0c88 | ||
|
54c9e89c3e | ||
|
875957f662 | ||
|
b251368b09 | ||
|
c2a1a7f247 | ||
|
9785637b3b | ||
|
728b34b684 |
@@ -4,6 +4,7 @@ Replicated in-memory database and file store.
|
|||||||
|
|
||||||
## TODO
|
## TODO
|
||||||
|
|
||||||
|
* [ ] mdb: Tests for using `nil` snapshots ?
|
||||||
* [ ] mdb: tests for sanitize and validate functions
|
* [ ] mdb: tests for sanitize and validate functions
|
||||||
* [ ] Test: lib/wal iterator w/ corrupt file (random corruptions)
|
* [ ] Test: lib/wal iterator w/ corrupt file (random corruptions)
|
||||||
* [ ] Test: lib/wal io.go
|
* [ ] Test: lib/wal io.go
|
||||||
|
@@ -8,26 +8,27 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type Error struct {
|
type Error struct {
|
||||||
msg string
|
Code int64
|
||||||
code int64
|
Msg string
|
||||||
|
StackTrace string
|
||||||
|
|
||||||
collection string
|
collection string
|
||||||
index string
|
index string
|
||||||
stackTrace string
|
|
||||||
err error // Wrapped error
|
err error // Wrapped error
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewErr(code int64, msg string) *Error {
|
func NewErr(code int64, msg string) *Error {
|
||||||
return &Error{
|
return &Error{
|
||||||
msg: msg,
|
Code: code,
|
||||||
code: code,
|
Msg: msg,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Error) Error() string {
|
func (e *Error) Error() string {
|
||||||
if e.collection != "" || e.index != "" {
|
if e.collection != "" || e.index != "" {
|
||||||
return fmt.Sprintf(`[%d] (%s/%s) %s`, e.code, e.collection, e.index, e.msg)
|
return fmt.Sprintf(`[%d] (%s/%s) %s`, e.Code, e.collection, e.index, e.Msg)
|
||||||
} else {
|
} else {
|
||||||
return fmt.Sprintf("[%d] %s", e.code, e.msg)
|
return fmt.Sprintf("[%d] %s", e.Code, e.Msg)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -36,11 +37,15 @@ func (e *Error) Is(rhs error) bool {
|
|||||||
if !ok {
|
if !ok {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
return e.code == e2.code
|
return e.Code == e2.Code
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *Error) Unwrap() error {
|
||||||
|
return e.err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Error) WithErr(err error) *Error {
|
func (e *Error) WithErr(err error) *Error {
|
||||||
if e2, ok := err.(*Error); ok && e2.code == e.code {
|
if e2, ok := err.(*Error); ok && e2.Code == e.Code {
|
||||||
return e2
|
return e2
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -49,18 +54,11 @@ func (e *Error) WithErr(err error) *Error {
|
|||||||
return e2
|
return e2
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Error) Unwrap() error {
|
|
||||||
if e.err != nil {
|
|
||||||
return e.err
|
|
||||||
}
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
|
|
||||||
func (e *Error) WithMsg(msg string, args ...any) *Error {
|
func (e *Error) WithMsg(msg string, args ...any) *Error {
|
||||||
err := *e
|
err := *e
|
||||||
err.msg += ": " + fmt.Sprintf(msg, args...)
|
err.Msg += ": " + fmt.Sprintf(msg, args...)
|
||||||
if len(err.stackTrace) == 0 {
|
if len(err.StackTrace) == 0 {
|
||||||
err.stackTrace = string(debug.Stack())
|
err.StackTrace = string(debug.Stack())
|
||||||
}
|
}
|
||||||
return &err
|
return &err
|
||||||
}
|
}
|
||||||
@@ -78,16 +76,16 @@ func (e *Error) WithIndex(s string) *Error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (e *Error) msgTruncacted() string {
|
func (e *Error) msgTruncacted() string {
|
||||||
if len(e.msg) > 255 {
|
if len(e.Msg) > 255 {
|
||||||
return e.msg[:255]
|
return e.Msg[:255]
|
||||||
}
|
}
|
||||||
return e.msg
|
return e.Msg
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Error) Write(w io.Writer) error {
|
func (e *Error) Write(w io.Writer) error {
|
||||||
msg := e.msgTruncacted()
|
msg := e.msgTruncacted()
|
||||||
|
|
||||||
if err := binary.Write(w, binary.LittleEndian, e.code); err != nil {
|
if err := binary.Write(w, binary.LittleEndian, e.Code); err != nil {
|
||||||
return IO.WithErr(err)
|
return IO.WithErr(err)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -103,7 +101,7 @@ func (e *Error) Read(r io.Reader) error {
|
|||||||
size uint8
|
size uint8
|
||||||
)
|
)
|
||||||
|
|
||||||
if err := binary.Read(r, binary.LittleEndian, &e.code); err != nil {
|
if err := binary.Read(r, binary.LittleEndian, &e.Code); err != nil {
|
||||||
return IO.WithErr(err)
|
return IO.WithErr(err)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -116,6 +114,6 @@ func (e *Error) Read(r io.Reader) error {
|
|||||||
return IO.WithErr(err)
|
return IO.WithErr(err)
|
||||||
}
|
}
|
||||||
|
|
||||||
e.msg = string(msgBuf)
|
e.Msg = string(msgBuf)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
@@ -10,12 +10,12 @@ func FmtDetails(err error) string {
|
|||||||
|
|
||||||
var s string
|
var s string
|
||||||
if e.collection != "" || e.index != "" {
|
if e.collection != "" || e.index != "" {
|
||||||
s = fmt.Sprintf(`[%d] (%s/%s) %s`, e.code, e.collection, e.index, e.msg)
|
s = fmt.Sprintf(`[%d] (%s/%s) %s`, e.Code, e.collection, e.index, e.Msg)
|
||||||
} else {
|
} else {
|
||||||
s = fmt.Sprintf("[%d] %s", e.code, e.msg)
|
s = fmt.Sprintf("[%d] %s", e.Code, e.Msg)
|
||||||
}
|
}
|
||||||
if len(e.stackTrace) != 0 {
|
if len(e.StackTrace) != 0 {
|
||||||
s += "\n\nStack Trace:\n" + e.stackTrace + "\n"
|
s += "\n\nStack Trace:\n" + e.StackTrace + "\n"
|
||||||
}
|
}
|
||||||
|
|
||||||
return s
|
return s
|
||||||
|
@@ -12,12 +12,6 @@ const (
|
|||||||
pathStreamWAL = "stream-wal"
|
pathStreamWAL = "stream-wal"
|
||||||
)
|
)
|
||||||
|
|
||||||
func (rep *Replicator) RegisterHandlers(mux *http.ServeMux, rootPath string) {
|
|
||||||
mux.HandleFunc(path.Join(rootPath, pathGetInfo), rep.handleGetInfo)
|
|
||||||
mux.HandleFunc(path.Join(rootPath, pathSendState), rep.handleSendState)
|
|
||||||
mux.HandleFunc(path.Join(rootPath, pathStreamWAL), rep.handleStreamWAL)
|
|
||||||
}
|
|
||||||
|
|
||||||
// TODO: Remove this!
|
// TODO: Remove this!
|
||||||
func (rep *Replicator) Handle(w http.ResponseWriter, r *http.Request) {
|
func (rep *Replicator) Handle(w http.ResponseWriter, r *http.Request) {
|
||||||
// We'll handle two types of requests: HTTP GET requests for JSON, or
|
// We'll handle two types of requests: HTTP GET requests for JSON, or
|
||||||
|
@@ -15,7 +15,7 @@ func (rep *Replicator) runWALGC() {
|
|||||||
select {
|
select {
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
state := rep.getState()
|
state := rep.getState()
|
||||||
before := time.Now().Unix() - rep.conf.WALSegMaxAgeSec
|
before := time.Now().Unix() - rep.conf.WALSegGCAgeSec
|
||||||
if err := rep.wal.DeleteBefore(before, state.SeqNum); err != nil {
|
if err := rep.wal.DeleteBefore(before, state.SeqNum); err != nil {
|
||||||
log.Printf("[WAL-GC] failed to delete wal segments: %v", err)
|
log.Printf("[WAL-GC] failed to delete wal segments: %v", err)
|
||||||
}
|
}
|
||||||
|
@@ -159,6 +159,15 @@ func (c *Collection[T]) Get(tx *Snapshot, id uint64) *T {
|
|||||||
return c.ByID.Get(tx, item)
|
return c.ByID.Get(tx, item)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Collection[T]) Has(tx *Snapshot, id uint64) bool {
|
||||||
|
if tx == nil {
|
||||||
|
tx = c.db.Snapshot()
|
||||||
|
}
|
||||||
|
item := new(T)
|
||||||
|
c.setID(item, id)
|
||||||
|
return c.ByID.Has(tx, item)
|
||||||
|
}
|
||||||
|
|
||||||
func (c *Collection[T]) Insert(tx *Snapshot, userItem *T) error {
|
func (c *Collection[T]) Insert(tx *Snapshot, userItem *T) error {
|
||||||
if tx == nil {
|
if tx == nil {
|
||||||
return c.db.Update(func(tx *Snapshot) error {
|
return c.db.Update(func(tx *Snapshot) error {
|
||||||
@@ -237,6 +246,27 @@ func (c *Collection[T]) update(tx *Snapshot, userItem *T) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *Collection[T]) UpdateFunc(tx *Snapshot, id uint64, update func(item *T) error) error {
|
||||||
|
if tx == nil {
|
||||||
|
return c.db.Update(func(tx *Snapshot) error {
|
||||||
|
return c.updateFunc(tx, id, update)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
return c.updateFunc(tx, id, update)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *Collection[T]) updateFunc(tx *Snapshot, id uint64, update func(item *T) error) error {
|
||||||
|
item := c.Get(tx, id)
|
||||||
|
if item == nil {
|
||||||
|
return errs.NotFound
|
||||||
|
}
|
||||||
|
if err := update(item); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
c.setID(item, id) // Don't allow the ID to change.
|
||||||
|
return c.update(tx, item)
|
||||||
|
}
|
||||||
|
|
||||||
func (c *Collection[T]) Upsert(tx *Snapshot, item *T) error {
|
func (c *Collection[T]) Upsert(tx *Snapshot, item *T) error {
|
||||||
if tx == nil {
|
if tx == nil {
|
||||||
return c.db.Update(func(tx *Snapshot) error {
|
return c.db.Update(func(tx *Snapshot) error {
|
||||||
|
@@ -73,6 +73,10 @@ type Database struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func New(conf Config) *Database {
|
func New(conf Config) *Database {
|
||||||
|
if conf.NetTimeout <= 0 {
|
||||||
|
conf.NetTimeout = time.Minute
|
||||||
|
}
|
||||||
|
|
||||||
if conf.MaxConcurrentUpdates <= 0 {
|
if conf.MaxConcurrentUpdates <= 0 {
|
||||||
conf.MaxConcurrentUpdates = 32
|
conf.MaxConcurrentUpdates = 32
|
||||||
}
|
}
|
||||||
|
@@ -51,6 +51,10 @@ func (f *freeList) Push(pages ...uint64) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (f *freeList) SetNextPage(nextPage uint64) {
|
||||||
|
f.nextPage = nextPage
|
||||||
|
}
|
||||||
|
|
||||||
func (f *freeList) Pop(count int, out []uint64) []uint64 {
|
func (f *freeList) Pop(count int, out []uint64) []uint64 {
|
||||||
out = out[:0]
|
out = out[:0]
|
||||||
|
|
||||||
|
@@ -13,14 +13,19 @@ type Index struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func NewIndex(f *File) (*Index, error) {
|
func NewIndex(f *File) (*Index, error) {
|
||||||
|
firstPage, err := f.pageCount()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
idx := &Index{
|
idx := &Index{
|
||||||
fList: newFreeList(0),
|
fList: newFreeList(firstPage),
|
||||||
aList: *newAllocList(),
|
aList: *newAllocList(),
|
||||||
seen: map[[2]uint64]struct{}{},
|
seen: map[[2]uint64]struct{}{},
|
||||||
mask: []bool{},
|
mask: []bool{},
|
||||||
}
|
}
|
||||||
|
|
||||||
err := f.iterate(func(pageID uint64, page dataPage) error {
|
err = f.iterate(func(pageID uint64, page dataPage) error {
|
||||||
header := page.Header()
|
header := page.Header()
|
||||||
switch header.PageType {
|
switch header.PageType {
|
||||||
case pageTypeHead:
|
case pageTypeHead:
|
||||||
|
@@ -134,6 +134,21 @@ func (pf *File) writePage(page dataPage, id uint64) error {
|
|||||||
// Reading
|
// Reading
|
||||||
// ----------------------------------------------------------------------------
|
// ----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
func (pf *File) pageCount() (uint64, error) {
|
||||||
|
fi, err := pf.f.Stat()
|
||||||
|
if err != nil {
|
||||||
|
return 0, errs.IO.WithErr(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
fileSize := fi.Size()
|
||||||
|
if fileSize%pageSize != 0 {
|
||||||
|
return 0, errs.Corrupt.WithMsg("File size isn't a multiple of page size.")
|
||||||
|
}
|
||||||
|
|
||||||
|
maxPage := uint64(fileSize / pageSize)
|
||||||
|
return maxPage, nil
|
||||||
|
}
|
||||||
|
|
||||||
func (pf *File) iterate(each func(pageID uint64, page dataPage) error) error {
|
func (pf *File) iterate(each func(pageID uint64, page dataPage) error) error {
|
||||||
pf.lock.RLock()
|
pf.lock.RLock()
|
||||||
defer pf.lock.RUnlock()
|
defer pf.lock.RUnlock()
|
||||||
|
Reference in New Issue
Block a user