Compare commits
19 Commits
ff8b87d6ea
...
v0.10.0
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0518c5dcca | ||
|
|
6b0b7408bc | ||
|
|
5ffd50bdea | ||
|
|
91b2ba30f6 | ||
|
|
c2828592ac | ||
|
|
13e53f0c88 | ||
|
|
54c9e89c3e | ||
|
|
875957f662 | ||
|
|
b251368b09 | ||
|
|
c2a1a7f247 | ||
|
|
9785637b3b | ||
|
|
728b34b684 | ||
|
|
8be663a0a0 | ||
|
|
6da018353a | ||
|
|
5528c264d3 | ||
|
|
9078335d70 | ||
|
|
3c9e5505ab | ||
| ba1990e379 | |||
| 526196ef9d |
@@ -4,6 +4,7 @@ Replicated in-memory database and file store.
|
||||
|
||||
## TODO
|
||||
|
||||
* [ ] mdb: Tests for using `nil` snapshots ?
|
||||
* [ ] mdb: tests for sanitize and validate functions
|
||||
* [ ] Test: lib/wal iterator w/ corrupt file (random corruptions)
|
||||
* [ ] Test: lib/wal io.go
|
||||
|
||||
@@ -3,10 +3,11 @@ package fstore
|
||||
import (
|
||||
"embed"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/fstore/pages"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
|
||||
"git.crumpington.com/public/jldb/fstore/pages"
|
||||
)
|
||||
|
||||
//go:embed static/*
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"io"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
package fstore
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func filesRootPath(rootDir string) string {
|
||||
|
||||
@@ -1,11 +1,12 @@
|
||||
package fstore
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/idgen"
|
||||
"log"
|
||||
"os"
|
||||
"path/filepath"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/idgen"
|
||||
)
|
||||
|
||||
func (s *Store) applyStoreFromReader(cmd command) error {
|
||||
|
||||
@@ -5,12 +5,13 @@ import (
|
||||
"errors"
|
||||
"io"
|
||||
"io/fs"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
)
|
||||
|
||||
func (s *Store) repSendState(conn net.Conn) error {
|
||||
|
||||
@@ -3,15 +3,16 @@ package fstore
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/idgen"
|
||||
"git.crumpington.com/public/jldb/lib/rep"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/idgen"
|
||||
"git.crumpington.com/public/jldb/lib/rep"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
|
||||
3
go.mod
3
go.mod
@@ -1,9 +1,10 @@
|
||||
module git.crumpington.com/public/jldb
|
||||
|
||||
go 1.21.1
|
||||
go 1.22
|
||||
|
||||
require (
|
||||
github.com/google/btree v1.1.2
|
||||
go.uber.org/goleak v1.3.0
|
||||
golang.org/x/net v0.15.0
|
||||
golang.org/x/sys v0.12.0
|
||||
)
|
||||
|
||||
10
go.sum
10
go.sum
@@ -1,6 +1,16 @@
|
||||
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
|
||||
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/google/btree v1.1.2 h1:xf4v41cLI2Z6FxbKm+8Bu+m8ifhj15JuZ9sa0jZCMUU=
|
||||
github.com/google/btree v1.1.2/go.mod h1:qOPhT0dTNdNzV6Z/lhRX0YXUafgPLFUh+gZMl761Gm4=
|
||||
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
github.com/stretchr/testify v1.8.0 h1:pSgiaMZlXftHpm5L7V1+rVB+AZJydKsMxsQBIJw4PKk=
|
||||
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
|
||||
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
|
||||
golang.org/x/net v0.15.0 h1:ugBLEUaxABaB5AJqW9enI0ACdci2RUd4eP51NTBvuJ8=
|
||||
golang.org/x/net v0.15.0/go.mod h1:idbUs1IY1+zTqbi8yxTbhexhEEk5ur9LInksu6HrEpk=
|
||||
golang.org/x/sys v0.12.0 h1:CM0HF96J0hcLAwsHPJZjfdNzs0gftsLfgKt57wWHJ0o=
|
||||
golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
|
||||
@@ -3,9 +3,10 @@ package atomicheader
|
||||
import (
|
||||
"encoding/binary"
|
||||
"hash/crc32"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"os"
|
||||
"sync"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
const (
|
||||
|
||||
@@ -8,26 +8,27 @@ import (
|
||||
)
|
||||
|
||||
type Error struct {
|
||||
msg string
|
||||
code int64
|
||||
Code int64
|
||||
Msg string
|
||||
StackTrace string
|
||||
|
||||
collection string
|
||||
index string
|
||||
stackTrace string
|
||||
err error // Wrapped error
|
||||
}
|
||||
|
||||
func NewErr(code int64, msg string) *Error {
|
||||
return &Error{
|
||||
msg: msg,
|
||||
code: code,
|
||||
Code: code,
|
||||
Msg: msg,
|
||||
}
|
||||
}
|
||||
|
||||
func (e *Error) Error() string {
|
||||
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 {
|
||||
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 {
|
||||
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 {
|
||||
if e2, ok := err.(*Error); ok && e2.code == e.code {
|
||||
if e2, ok := err.(*Error); ok && e2.Code == e.Code {
|
||||
return e2
|
||||
}
|
||||
|
||||
@@ -49,18 +54,11 @@ func (e *Error) WithErr(err error) *Error {
|
||||
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 {
|
||||
err := *e
|
||||
err.msg += ": " + fmt.Sprintf(msg, args...)
|
||||
if len(err.stackTrace) == 0 {
|
||||
err.stackTrace = string(debug.Stack())
|
||||
err.Msg += ": " + fmt.Sprintf(msg, args...)
|
||||
if len(err.StackTrace) == 0 {
|
||||
err.StackTrace = string(debug.Stack())
|
||||
}
|
||||
return &err
|
||||
}
|
||||
@@ -78,16 +76,16 @@ func (e *Error) WithIndex(s string) *Error {
|
||||
}
|
||||
|
||||
func (e *Error) msgTruncacted() string {
|
||||
if len(e.msg) > 255 {
|
||||
return e.msg[:255]
|
||||
if len(e.Msg) > 255 {
|
||||
return e.Msg[:255]
|
||||
}
|
||||
return e.msg
|
||||
return e.Msg
|
||||
}
|
||||
|
||||
func (e *Error) Write(w io.Writer) error {
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -103,7 +101,7 @@ func (e *Error) Read(r io.Reader) error {
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -116,6 +114,6 @@ func (e *Error) Read(r io.Reader) error {
|
||||
return IO.WithErr(err)
|
||||
}
|
||||
|
||||
e.msg = string(msgBuf)
|
||||
e.Msg = string(msgBuf)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -10,12 +10,12 @@ func FmtDetails(err error) string {
|
||||
|
||||
var s string
|
||||
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 {
|
||||
s = fmt.Sprintf("[%d] %s", e.code, e.msg)
|
||||
s = fmt.Sprintf("[%d] %s", e.Code, e.Msg)
|
||||
}
|
||||
if len(e.stackTrace) != 0 {
|
||||
s += "\n\nStack Trace:\n" + e.stackTrace + "\n"
|
||||
if len(e.StackTrace) != 0 {
|
||||
s += "\n\nStack Trace:\n" + e.StackTrace + "\n"
|
||||
}
|
||||
|
||||
return s
|
||||
|
||||
@@ -6,11 +6,12 @@ import (
|
||||
"crypto/tls"
|
||||
"errors"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
var ErrInvalidStatus = errors.New("invalid status")
|
||||
|
||||
@@ -1,51 +0,0 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"net"
|
||||
"path/filepath"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
|
||||
func lockFilePath(rootDir string) string {
|
||||
return filepath.Join(rootDir, "lock")
|
||||
}
|
||||
|
||||
func walRootDir(rootDir string) string {
|
||||
return filepath.Join(rootDir, "wal")
|
||||
}
|
||||
|
||||
func stateFilePath(rootDir string) string {
|
||||
return filepath.Join(rootDir, "state")
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
|
||||
func sendJSON(
|
||||
item any,
|
||||
conn net.Conn,
|
||||
timeout time.Duration,
|
||||
) error {
|
||||
|
||||
buf := bufPoolGet()
|
||||
defer bufPoolPut(buf)
|
||||
|
||||
if err := json.NewEncoder(buf).Encode(item); err != nil {
|
||||
return errs.Unexpected.WithErr(err)
|
||||
}
|
||||
|
||||
sizeBuf := make([]byte, 2)
|
||||
binary.LittleEndian.PutUint16(sizeBuf, uint16(buf.Len()))
|
||||
|
||||
conn.SetWriteDeadline(time.Now().Add(timeout))
|
||||
buffers := net.Buffers{sizeBuf, buf.Bytes()}
|
||||
if _, err := buffers.WriteTo(conn); err != nil {
|
||||
return errs.IO.WithErr(err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -1,178 +1,109 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/httpconn"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
type client struct {
|
||||
// Mutex-protected variables.
|
||||
lock sync.Mutex
|
||||
closed bool
|
||||
conn net.Conn
|
||||
client *http.Client
|
||||
|
||||
// The following are constant.
|
||||
endpoint string
|
||||
psk []byte
|
||||
pskBytes [64]byte
|
||||
timeout time.Duration
|
||||
|
||||
buf []byte
|
||||
lock sync.Mutex
|
||||
conn net.Conn
|
||||
}
|
||||
|
||||
func newClient(endpoint, psk string, timeout time.Duration) *client {
|
||||
b := make([]byte, 256)
|
||||
copy(b, []byte(psk))
|
||||
httpClient := &http.Client{
|
||||
Timeout: timeout,
|
||||
}
|
||||
|
||||
if !strings.HasSuffix(endpoint, "/") {
|
||||
endpoint += "/"
|
||||
}
|
||||
|
||||
return &client{
|
||||
client: httpClient,
|
||||
endpoint: endpoint,
|
||||
psk: b,
|
||||
pskBytes: pskToBytes(psk),
|
||||
timeout: timeout,
|
||||
}
|
||||
}
|
||||
|
||||
func (c *client) GetInfo() (info Info, err error) {
|
||||
err = c.withConn(cmdGetInfo, func(conn net.Conn) error {
|
||||
return c.recvJSON(&info, conn, c.timeout)
|
||||
})
|
||||
return info, err
|
||||
req, err := http.NewRequest(http.MethodGet, c.endpoint+pathGetInfo, nil)
|
||||
if err != nil {
|
||||
return info, errs.Unexpected.WithErr(err)
|
||||
}
|
||||
req.SetBasicAuth("", string(c.pskBytes[:]))
|
||||
|
||||
resp, err := c.client.Do(req)
|
||||
if err != nil {
|
||||
return info, errs.IO.WithErr(err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if err := json.NewDecoder(resp.Body).Decode(&info); err != nil {
|
||||
return info, errs.IO.WithErr(err)
|
||||
}
|
||||
return info, nil
|
||||
}
|
||||
|
||||
func (c *client) RecvState(recv func(net.Conn) error) error {
|
||||
return c.withConn(cmdSendState, recv)
|
||||
err := c.dialConnect(c.endpoint + pathSendState)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer c.conn.Close()
|
||||
|
||||
return recv(c.conn)
|
||||
}
|
||||
|
||||
func (c *client) StreamWAL(w *wal.WAL) error {
|
||||
return c.withConn(cmdStreamWAL, func(conn net.Conn) error {
|
||||
return w.Recv(conn, c.timeout)
|
||||
})
|
||||
err := c.dialConnect(c.endpoint + pathStreamWAL)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer c.conn.Close()
|
||||
|
||||
return w.Recv(c.conn, c.timeout)
|
||||
}
|
||||
|
||||
func (c *client) Close() {
|
||||
c.lock.Lock()
|
||||
defer c.lock.Unlock()
|
||||
c.closed = true
|
||||
|
||||
if c.conn != nil {
|
||||
c.conn.Close()
|
||||
c.conn = nil
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
|
||||
func (c *client) writeCmd(cmd byte) error {
|
||||
c.conn.SetWriteDeadline(time.Now().Add(c.timeout))
|
||||
if _, err := c.conn.Write([]byte{cmd}); err != nil {
|
||||
return errs.IO.WithErr(err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *client) dial() error {
|
||||
c.conn = nil
|
||||
|
||||
conn, err := httpconn.Dial(c.endpoint)
|
||||
func (c *client) dialConnect(endpoint string) error {
|
||||
conn, err := httpconn.Dial(endpoint)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
conn.SetWriteDeadline(time.Now().Add(c.timeout))
|
||||
if _, err := conn.Write(c.psk); err != nil {
|
||||
conn.Close()
|
||||
if _, err := conn.Write(c.pskBytes[:]); err != nil {
|
||||
return errs.IO.WithErr(err)
|
||||
}
|
||||
|
||||
c.lock.Lock()
|
||||
defer c.lock.Unlock()
|
||||
c.conn = conn
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *client) withConn(cmd byte, fn func(net.Conn) error) error {
|
||||
conn, err := c.getConn(cmd)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := fn(conn); err != nil {
|
||||
conn.Close()
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *client) getConn(cmd byte) (net.Conn, error) {
|
||||
c.lock.Lock()
|
||||
defer c.lock.Unlock()
|
||||
|
||||
if c.closed {
|
||||
return nil, errs.IO.WithErr(io.EOF)
|
||||
}
|
||||
|
||||
dialed := false
|
||||
|
||||
if c.conn == nil {
|
||||
if err := c.dial(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
dialed = true
|
||||
}
|
||||
|
||||
if err := c.writeCmd(cmd); err != nil {
|
||||
if dialed {
|
||||
c.conn = nil
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := c.dial(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := c.writeCmd(cmd); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
return c.conn, nil
|
||||
}
|
||||
|
||||
func (c *client) recvJSON(
|
||||
item any,
|
||||
conn net.Conn,
|
||||
timeout time.Duration,
|
||||
) error {
|
||||
|
||||
if cap(c.buf) < 2 {
|
||||
c.buf = make([]byte, 0, 1024)
|
||||
}
|
||||
buf := c.buf[:2]
|
||||
|
||||
conn.SetReadDeadline(time.Now().Add(timeout))
|
||||
|
||||
if _, err := io.ReadFull(conn, buf); err != nil {
|
||||
return errs.IO.WithErr(err)
|
||||
}
|
||||
|
||||
size := binary.LittleEndian.Uint16(buf)
|
||||
|
||||
if cap(buf) < int(size) {
|
||||
buf = make([]byte, size)
|
||||
c.buf = buf
|
||||
}
|
||||
buf = buf[:size]
|
||||
|
||||
if _, err := io.ReadFull(conn, buf); err != nil {
|
||||
return errs.IO.WithErr(err)
|
||||
}
|
||||
|
||||
if err := json.Unmarshal(buf, item); err != nil {
|
||||
return errs.Unexpected.WithErr(err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
54
lib/rep/http-handler-util.go
Normal file
54
lib/rep/http-handler-util.go
Normal file
@@ -0,0 +1,54 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"crypto/subtle"
|
||||
"io"
|
||||
"log"
|
||||
"net"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/httpconn"
|
||||
)
|
||||
|
||||
func (rep *Replicator) handlerLogf(pattern string, args ...any) {
|
||||
log.Printf("[HTTP-HANDLER] "+pattern, args...)
|
||||
}
|
||||
|
||||
// checkBasicAuth tests whether the caller has provided the appropriate basic auth
|
||||
// header. The caller should provide the PSK as the basic-auth password.
|
||||
func (rep *Replicator) checkBasicAuth(w http.ResponseWriter, r *http.Request) bool {
|
||||
_, pwd, _ := r.BasicAuth()
|
||||
if subtle.ConstantTimeCompare([]byte(pwd), rep.pskBytes[:]) != 1 {
|
||||
rep.handlerLogf("PSK mismatch.")
|
||||
http.Error(w, "not authorized", http.StatusUnauthorized)
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// acceptConnect accepts a CONNECT request and checks the PSK.
|
||||
func (rep *Replicator) acceptConnect(w http.ResponseWriter, r *http.Request) net.Conn {
|
||||
conn, err := httpconn.Accept(w, r)
|
||||
if err != nil {
|
||||
rep.handlerLogf("Failed to accept connection: %s", err)
|
||||
return nil
|
||||
}
|
||||
|
||||
psk := [64]byte{}
|
||||
|
||||
conn.SetReadDeadline(time.Now().Add(rep.conf.NetTimeout))
|
||||
if _, err := io.ReadFull(conn, psk[:]); err != nil {
|
||||
conn.Close()
|
||||
rep.handlerLogf("Failed to read PSK: %v", err)
|
||||
return nil
|
||||
}
|
||||
|
||||
if subtle.ConstantTimeCompare(psk[:], rep.pskBytes[:]) != 1 {
|
||||
conn.Close()
|
||||
rep.handlerLogf("PSK mismatch.")
|
||||
return nil
|
||||
}
|
||||
|
||||
return conn
|
||||
}
|
||||
@@ -1,79 +1,69 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"crypto/subtle"
|
||||
"git.crumpington.com/public/jldb/lib/httpconn"
|
||||
"log"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"time"
|
||||
"path"
|
||||
)
|
||||
|
||||
const (
|
||||
cmdGetInfo = 10
|
||||
cmdSendState = 20
|
||||
cmdStreamWAL = 30
|
||||
pathGetInfo = "get-info"
|
||||
pathSendState = "send-state"
|
||||
pathStreamWAL = "stream-wal"
|
||||
)
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func (rep *Replicator) Handle(w http.ResponseWriter, r *http.Request) {
|
||||
logf := func(pattern string, args ...any) {
|
||||
log.Printf("[HTTP-HANDLER] "+pattern, args...)
|
||||
// We'll handle two types of requests: HTTP GET requests for JSON, or
|
||||
// streaming requets for state or wall.
|
||||
|
||||
base := path.Base(r.URL.Path)
|
||||
switch base {
|
||||
case pathGetInfo:
|
||||
rep.handleGetInfo(w, r)
|
||||
case pathSendState:
|
||||
rep.handleSendState(w, r)
|
||||
case pathStreamWAL:
|
||||
rep.handleStreamWAL(w, r)
|
||||
default:
|
||||
http.Error(w, "not found", http.StatusNotFound)
|
||||
}
|
||||
}
|
||||
|
||||
func (rep *Replicator) handleGetInfo(w http.ResponseWriter, r *http.Request) {
|
||||
if !rep.checkBasicAuth(w, r) {
|
||||
return
|
||||
}
|
||||
|
||||
conn, err := httpconn.Accept(w, r)
|
||||
if err != nil {
|
||||
logf("Failed to accept connection: %s", err)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
|
||||
if err := json.NewEncoder(w).Encode(rep.Info()); err != nil {
|
||||
rep.handlerLogf("Failed to send info: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (rep *Replicator) handleSendState(w http.ResponseWriter, r *http.Request) {
|
||||
conn := rep.acceptConnect(w, r)
|
||||
if conn == nil {
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
psk := make([]byte, 256)
|
||||
|
||||
conn.SetReadDeadline(time.Now().Add(rep.conf.NetTimeout))
|
||||
if _, err := conn.Read(psk); err != nil {
|
||||
logf("Failed to read PSK: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
expected := rep.pskBytes
|
||||
if subtle.ConstantTimeCompare(expected, psk) != 1 {
|
||||
logf("PSK mismatch.")
|
||||
return
|
||||
}
|
||||
|
||||
cmd := make([]byte, 1)
|
||||
|
||||
for {
|
||||
conn.SetReadDeadline(time.Now().Add(rep.conf.NetTimeout))
|
||||
if _, err := conn.Read(cmd); err != nil {
|
||||
logf("Read failed: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
switch cmd[0] {
|
||||
|
||||
case cmdGetInfo:
|
||||
if err := sendJSON(rep.Info(), conn, rep.conf.NetTimeout); err != nil {
|
||||
logf("Failed to send info: %s", err)
|
||||
return
|
||||
}
|
||||
|
||||
case cmdSendState:
|
||||
|
||||
if err := rep.sendState(conn); err != nil {
|
||||
if !rep.stopped() {
|
||||
logf("Failed to send state: %s", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
case cmdStreamWAL:
|
||||
err := rep.wal.Send(conn, rep.conf.NetTimeout)
|
||||
if !rep.stopped() {
|
||||
logf("Failed when sending WAL: %s", err)
|
||||
}
|
||||
return
|
||||
if err := rep.sendState(conn); err != nil {
|
||||
if !rep.stopped() {
|
||||
rep.handlerLogf("Failed to send state: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (rep *Replicator) handleStreamWAL(w http.ResponseWriter, r *http.Request) {
|
||||
conn := rep.acceptConnect(w, r)
|
||||
if conn == nil {
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
err := rep.wal.Send(conn, rep.conf.NetTimeout)
|
||||
if !rep.stopped() {
|
||||
rep.handlerLogf("Failed when streaming WAL: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
11
lib/rep/main_test.go
Normal file
11
lib/rep/main_test.go
Normal file
@@ -0,0 +1,11 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"go.uber.org/goleak"
|
||||
)
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
goleak.VerifyTestMain(m)
|
||||
}
|
||||
17
lib/rep/paths.go
Normal file
17
lib/rep/paths.go
Normal file
@@ -0,0 +1,17 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
)
|
||||
|
||||
func lockFilePath(rootDir string) string {
|
||||
return filepath.Join(rootDir, "lock")
|
||||
}
|
||||
|
||||
func walRootDir(rootDir string) string {
|
||||
return filepath.Join(rootDir, "wal")
|
||||
}
|
||||
|
||||
func stateFilePath(rootDir string) string {
|
||||
return filepath.Join(rootDir, "state")
|
||||
}
|
||||
@@ -1,21 +0,0 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"sync"
|
||||
)
|
||||
|
||||
var bufPool = sync.Pool{
|
||||
New: func() any {
|
||||
return &bytes.Buffer{}
|
||||
},
|
||||
}
|
||||
|
||||
func bufPoolGet() *bytes.Buffer {
|
||||
return bufPool.Get().(*bytes.Buffer)
|
||||
}
|
||||
|
||||
func bufPoolPut(b *bytes.Buffer) {
|
||||
b.Reset()
|
||||
bufPool.Put(b)
|
||||
}
|
||||
16
lib/rep/psk.go
Normal file
16
lib/rep/psk.go
Normal file
@@ -0,0 +1,16 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
)
|
||||
|
||||
func pskToBytes(in string) [64]byte {
|
||||
b := sha256.Sum256([]byte(in))
|
||||
dst := [64]byte{}
|
||||
i := hex.Encode(dst[:], b[:])
|
||||
if i != 64 {
|
||||
panic(i)
|
||||
}
|
||||
return dst
|
||||
}
|
||||
@@ -2,9 +2,10 @@ package rep
|
||||
|
||||
import (
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"net"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func (rep *Replicator) sendState(conn net.Conn) error {
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
package rep
|
||||
|
||||
import (
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/flock"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
func (rep *Replicator) loadConfigDefaults() {
|
||||
@@ -26,9 +27,7 @@ func (rep *Replicator) loadConfigDefaults() {
|
||||
}
|
||||
|
||||
rep.conf = conf
|
||||
|
||||
rep.pskBytes = make([]byte, 256)
|
||||
copy(rep.pskBytes, []byte(conf.ReplicationPSK))
|
||||
rep.pskBytes = pskToBytes(conf.ReplicationPSK)
|
||||
}
|
||||
|
||||
func (rep *Replicator) initDirectories() error {
|
||||
|
||||
@@ -15,7 +15,7 @@ func (rep *Replicator) runWALGC() {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
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 {
|
||||
log.Printf("[WAL-GC] failed to delete wal segments: %v", err)
|
||||
}
|
||||
|
||||
@@ -30,9 +30,8 @@ func (rep *Replicator) runWALRecvrOnce() {
|
||||
log.Printf("[WAL-RECVR] "+pattern, args...)
|
||||
}
|
||||
|
||||
if err := rep.client.StreamWAL(rep.wal); err != nil {
|
||||
if !rep.stopped() {
|
||||
logf("Recv failed: %v", err)
|
||||
}
|
||||
err := rep.client.StreamWAL(rep.wal)
|
||||
if !rep.stopped() {
|
||||
logf("Recv failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,14 +2,16 @@ package rep
|
||||
|
||||
import (
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
@@ -34,8 +36,8 @@ type App struct {
|
||||
// SendState: The primary may need to send storage state to a secondary node.
|
||||
SendState func(conn net.Conn) error
|
||||
|
||||
// (1) RecvState: Secondary nodes may need to load state from the primary if the
|
||||
// WAL is too far behind.
|
||||
// (1) RecvState: Secondary nodes may need to load state from the primary if
|
||||
// the WAL is too far behind.
|
||||
RecvState func(conn net.Conn) error
|
||||
|
||||
// (2) InitStorage: Prepare application storage for possible calls to
|
||||
@@ -59,8 +61,9 @@ type Replicator struct {
|
||||
conf Config
|
||||
|
||||
lockFile *os.File
|
||||
pskBytes []byte
|
||||
wal *wal.WAL
|
||||
pskBytes [64]byte // 64 ascii characters. See pskToBytes.
|
||||
|
||||
wal *wal.WAL
|
||||
|
||||
appendNotify chan struct{}
|
||||
|
||||
@@ -92,41 +95,49 @@ func Open(app App, conf Config) (*Replicator, error) {
|
||||
rep.client = newClient(rep.conf.PrimaryEndpoint, rep.conf.ReplicationPSK, rep.conf.NetTimeout)
|
||||
|
||||
if err := rep.initDirectories(); err != nil {
|
||||
log.Printf("Failed to init directories: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.acquireLock(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to acquire lock: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.loadLocalState(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to load local state: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.openWAL(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to open WAL: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.recvStateIfNecessary(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to recv state: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.app.InitStorage(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to init storage: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.replay(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to replay: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rep.app.LoadFromStorage(); err != nil {
|
||||
rep.Close()
|
||||
log.Printf("Failed to load from storage: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -161,10 +172,6 @@ func (rep *Replicator) Primary() bool {
|
||||
return rep.conf.Primary
|
||||
}
|
||||
|
||||
// TODO: Probably remove this.
|
||||
// The caller may call Ack after Apply to acknowledge that the change has also
|
||||
// been applied to the caller's application. Alternatively, the caller may use
|
||||
// follow to apply changes to their application state.
|
||||
func (rep *Replicator) ack(seqNum, timestampMS int64) error {
|
||||
state := rep.getState()
|
||||
state.SeqNum = seqNum
|
||||
|
||||
@@ -56,6 +56,7 @@ func (h TestAppHarness) Run(t *testing.T) {
|
||||
WALSegMaxAgeSec: 1,
|
||||
WALSegGCAgeSec: 1,
|
||||
})
|
||||
defer app2.Close()
|
||||
|
||||
val.MethodByName(method.Name).Call([]reflect.Value{
|
||||
reflect.ValueOf(t),
|
||||
|
||||
@@ -5,12 +5,14 @@ import (
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"math/rand"
|
||||
"net"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
)
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
@@ -82,7 +84,7 @@ func newApp(t *testing.T, id int64, conf Config) *TestApp {
|
||||
Apply: a.apply,
|
||||
}, conf)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
t.Fatal(errs.FmtDetails(err))
|
||||
}
|
||||
|
||||
return a
|
||||
|
||||
@@ -2,8 +2,9 @@ package wal
|
||||
|
||||
import (
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"testing"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func TestCorruptWAL(t *testing.T) {
|
||||
|
||||
@@ -5,13 +5,14 @@ import (
|
||||
"encoding/binary"
|
||||
"errors"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"math/rand"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
type waLog interface {
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"errors"
|
||||
"hash/crc32"
|
||||
"io"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"encoding/binary"
|
||||
"hash/crc32"
|
||||
"io"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
|
||||
@@ -3,10 +3,11 @@ package wal
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/testutil"
|
||||
"math/rand"
|
||||
"testing"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/testutil"
|
||||
)
|
||||
|
||||
func NewRecordForTesting() Record {
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
package wal
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
type segmentIterator struct {
|
||||
|
||||
@@ -3,11 +3,12 @@ package wal
|
||||
import (
|
||||
"bufio"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
type segment struct {
|
||||
|
||||
@@ -4,11 +4,12 @@ import (
|
||||
"bytes"
|
||||
crand "crypto/rand"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func newSegmentForTesting(t *testing.T) *segment {
|
||||
|
||||
@@ -1,8 +1,9 @@
|
||||
package wal
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
type walIterator struct {
|
||||
|
||||
@@ -3,9 +3,10 @@ package wal
|
||||
import (
|
||||
"encoding/binary"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"net"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func (wal *WAL) Recv(conn net.Conn, timeout time.Duration) error {
|
||||
|
||||
@@ -2,9 +2,10 @@ package wal
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"net"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
const (
|
||||
|
||||
@@ -1,8 +1,6 @@
|
||||
package wal
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/testutil"
|
||||
"log"
|
||||
"math/rand"
|
||||
"reflect"
|
||||
@@ -11,6 +9,9 @@ import (
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/testutil"
|
||||
)
|
||||
|
||||
func TestSendRecvHarness(t *testing.T) {
|
||||
|
||||
@@ -2,13 +2,14 @@ package wal
|
||||
|
||||
import (
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/atomicheader"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
|
||||
@@ -3,6 +3,7 @@ package change
|
||||
import (
|
||||
"encoding/binary"
|
||||
"io"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ package change
|
||||
|
||||
import (
|
||||
"io"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"hash/crc64"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"unsafe"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
|
||||
"github.com/google/btree"
|
||||
)
|
||||
|
||||
@@ -20,12 +20,10 @@ type Collection[T any] struct {
|
||||
sanitize func(*T)
|
||||
validate func(*T) error
|
||||
|
||||
indices []Index[T]
|
||||
uniqueIndices []Index[T]
|
||||
indices []*Index[T]
|
||||
uniqueIndices []*Index[T]
|
||||
|
||||
ByID Index[T]
|
||||
|
||||
buf *bytes.Buffer
|
||||
ByID *Index[T]
|
||||
}
|
||||
|
||||
type CollectionConfig[T any] struct {
|
||||
@@ -64,9 +62,8 @@ func NewCollection[T any](db *Database, name string, conf *CollectionConfig[T])
|
||||
copy: conf.Copy,
|
||||
sanitize: conf.Sanitize,
|
||||
validate: conf.Validate,
|
||||
indices: []Index[T]{},
|
||||
uniqueIndices: []Index[T]{},
|
||||
buf: &bytes.Buffer{},
|
||||
indices: []*Index[T]{},
|
||||
uniqueIndices: []*Index[T]{},
|
||||
}
|
||||
|
||||
db.addCollection(c.collectionID, c, &collectionState[T]{
|
||||
@@ -91,7 +88,7 @@ func NewCollection[T any](db *Database, name string, conf *CollectionConfig[T])
|
||||
return c
|
||||
}
|
||||
|
||||
func (c Collection[T]) Name() string {
|
||||
func (c *Collection[T]) Name() string {
|
||||
return c.name
|
||||
}
|
||||
|
||||
@@ -107,35 +104,8 @@ type indexConfig[T any] struct {
|
||||
Include func(item *T) bool
|
||||
}
|
||||
|
||||
func (c Collection[T]) Get(tx *Snapshot, id uint64) (*T, bool) {
|
||||
x := new(T)
|
||||
c.setID(x, id)
|
||||
return c.ByID.Get(tx, x)
|
||||
}
|
||||
|
||||
func (c Collection[T]) List(tx *Snapshot, ids []uint64, out []*T) []*T {
|
||||
if len(ids) == 0 {
|
||||
return out[:0]
|
||||
}
|
||||
|
||||
if cap(out) < len(ids) {
|
||||
out = make([]*T, len(ids))
|
||||
}
|
||||
out = out[:0]
|
||||
|
||||
for _, id := range ids {
|
||||
item, ok := c.Get(tx, id)
|
||||
if ok {
|
||||
out = append(out, item)
|
||||
}
|
||||
}
|
||||
|
||||
return out
|
||||
}
|
||||
|
||||
// AddIndex: Add an index to the collection.
|
||||
func (c *Collection[T]) addIndex(conf indexConfig[T]) Index[T] {
|
||||
|
||||
func (c *Collection[T]) addIndex(conf indexConfig[T]) *Index[T] {
|
||||
var less func(*T, *T) bool
|
||||
|
||||
if conf.Unique {
|
||||
@@ -159,7 +129,8 @@ func (c *Collection[T]) addIndex(conf indexConfig[T]) Index[T] {
|
||||
BTree: btree.NewG(256, less),
|
||||
}
|
||||
|
||||
index := Index[T]{
|
||||
index := &Index[T]{
|
||||
db: c.db,
|
||||
collectionID: c.collectionID,
|
||||
name: conf.Name,
|
||||
indexID: c.getState(c.db.Snapshot()).addIndex(indexState),
|
||||
@@ -175,7 +146,35 @@ func (c *Collection[T]) addIndex(conf indexConfig[T]) Index[T] {
|
||||
return index
|
||||
}
|
||||
|
||||
func (c Collection[T]) Insert(tx *Snapshot, userItem *T) error {
|
||||
func (c *Collection[T]) Get(tx *Snapshot, id uint64) *T {
|
||||
if tx == nil {
|
||||
tx = c.db.Snapshot()
|
||||
}
|
||||
item := new(T)
|
||||
c.setID(item, id)
|
||||
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 {
|
||||
if tx == nil {
|
||||
return c.db.Update(func(tx *Snapshot) error {
|
||||
return c.insert(tx, userItem)
|
||||
})
|
||||
}
|
||||
|
||||
return c.insert(tx, userItem)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) insert(tx *Snapshot, userItem *T) error {
|
||||
if err := c.ensureMutable(tx); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -189,7 +188,7 @@ func (c Collection[T]) Insert(tx *Snapshot, userItem *T) error {
|
||||
|
||||
for i := range c.uniqueIndices {
|
||||
if c.uniqueIndices[i].insertConflict(tx, item) {
|
||||
return ErrDuplicate.WithCollection(c.name).WithIndex(c.uniqueIndices[i].name)
|
||||
return errs.Duplicate.WithCollection(c.name).WithIndex(c.uniqueIndices[i].name)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -202,7 +201,16 @@ func (c Collection[T]) Insert(tx *Snapshot, userItem *T) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c Collection[T]) Update(tx *Snapshot, userItem *T) error {
|
||||
func (c *Collection[T]) Update(tx *Snapshot, userItem *T) error {
|
||||
if tx == nil {
|
||||
return c.db.Update(func(tx *Snapshot) error {
|
||||
return c.update(tx, userItem)
|
||||
})
|
||||
}
|
||||
return c.update(tx, userItem)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) update(tx *Snapshot, userItem *T) error {
|
||||
if err := c.ensureMutable(tx); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -216,12 +224,12 @@ func (c Collection[T]) Update(tx *Snapshot, userItem *T) error {
|
||||
|
||||
old, ok := c.ByID.get(tx, item)
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
return errs.NotFound
|
||||
}
|
||||
|
||||
for i := range c.uniqueIndices {
|
||||
if c.uniqueIndices[i].updateConflict(tx, item) {
|
||||
return ErrDuplicate.WithCollection(c.name).WithIndex(c.uniqueIndices[i].name)
|
||||
return errs.Duplicate.WithCollection(c.name).WithIndex(c.uniqueIndices[i].name)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,18 +242,87 @@ func (c Collection[T]) Update(tx *Snapshot, userItem *T) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c Collection[T]) Upsert(tx *Snapshot, item *T) error {
|
||||
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 {
|
||||
if tx == nil {
|
||||
return c.db.Update(func(tx *Snapshot) error {
|
||||
return c.upsert(tx, item)
|
||||
})
|
||||
}
|
||||
return c.upsert(tx, item)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) upsert(tx *Snapshot, item *T) error {
|
||||
err := c.Insert(tx, item)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
if errors.Is(err, ErrDuplicate) {
|
||||
if errors.Is(err, errs.Duplicate) {
|
||||
return c.Update(tx, item)
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (c Collection[T]) Delete(tx *Snapshot, itemID uint64) error {
|
||||
func (c *Collection[T]) UpsertFunc(tx *Snapshot, id uint64, update func(item *T) error) error {
|
||||
if tx == nil {
|
||||
c.db.Update(func(tx *Snapshot) error {
|
||||
return c.upsertFunc(tx, id, update)
|
||||
})
|
||||
}
|
||||
return c.upsertFunc(tx, id, update)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) upsertFunc(tx *Snapshot, id uint64, update func(item *T) error) error {
|
||||
insert := false
|
||||
|
||||
item := c.Get(tx, id)
|
||||
if item == nil {
|
||||
item = new(T)
|
||||
insert = true
|
||||
}
|
||||
|
||||
if err := update(item); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
c.setID(item, id) // Don't allow the ID to change.
|
||||
|
||||
if insert {
|
||||
return c.insert(tx, item)
|
||||
}
|
||||
return c.update(tx, item)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) Delete(tx *Snapshot, itemID uint64) error {
|
||||
if tx == nil {
|
||||
return c.db.Update(func(tx *Snapshot) error {
|
||||
return c.delete(tx, itemID)
|
||||
})
|
||||
}
|
||||
return c.delete(tx, itemID)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) delete(tx *Snapshot, itemID uint64) error {
|
||||
if err := c.ensureMutable(tx); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -253,15 +330,22 @@ func (c Collection[T]) Delete(tx *Snapshot, itemID uint64) error {
|
||||
return c.deleteItem(tx, itemID)
|
||||
}
|
||||
|
||||
func (c Collection[T]) getByID(tx *Snapshot, itemID uint64) (*T, bool) {
|
||||
func (c *Collection[T]) Count(tx *Snapshot) int {
|
||||
if tx == nil {
|
||||
tx = c.db.Snapshot()
|
||||
}
|
||||
return c.ByID.Count(tx)
|
||||
}
|
||||
|
||||
func (c *Collection[T]) getByID(tx *Snapshot, itemID uint64) (*T, bool) {
|
||||
x := new(T)
|
||||
c.setID(x, itemID)
|
||||
return c.ByID.get(tx, x)
|
||||
}
|
||||
|
||||
func (c Collection[T]) ensureMutable(tx *Snapshot) error {
|
||||
func (c *Collection[T]) ensureMutable(tx *Snapshot) error {
|
||||
if !tx.writable() {
|
||||
return ErrReadOnly
|
||||
return errs.ReadOnly
|
||||
}
|
||||
|
||||
state := c.getState(tx)
|
||||
@@ -273,7 +357,7 @@ func (c Collection[T]) ensureMutable(tx *Snapshot) error {
|
||||
}
|
||||
|
||||
// For initial data loading.
|
||||
func (c Collection[T]) insertItem(tx *Snapshot, itemID uint64, data []byte) error {
|
||||
func (c *Collection[T]) insertItem(tx *Snapshot, itemID uint64, data []byte) error {
|
||||
item := new(T)
|
||||
if err := json.Unmarshal(data, item); err != nil {
|
||||
return errs.Encoding.WithErr(err).WithCollection(c.name)
|
||||
@@ -282,7 +366,7 @@ func (c Collection[T]) insertItem(tx *Snapshot, itemID uint64, data []byte) erro
|
||||
// Check for insert conflict.
|
||||
for _, index := range c.uniqueIndices {
|
||||
if index.insertConflict(tx, item) {
|
||||
return ErrDuplicate
|
||||
return errs.Duplicate
|
||||
}
|
||||
}
|
||||
|
||||
@@ -294,10 +378,10 @@ func (c Collection[T]) insertItem(tx *Snapshot, itemID uint64, data []byte) erro
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c Collection[T]) deleteItem(tx *Snapshot, itemID uint64) error {
|
||||
func (c *Collection[T]) deleteItem(tx *Snapshot, itemID uint64) error {
|
||||
item, ok := c.getByID(tx, itemID)
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
return errs.NotFound
|
||||
}
|
||||
|
||||
tx.delete(c.collectionID, itemID)
|
||||
@@ -311,7 +395,7 @@ func (c Collection[T]) deleteItem(tx *Snapshot, itemID uint64) error {
|
||||
|
||||
// upsertItem inserts or updates the item with itemID and the given serialized
|
||||
// form. It's called by
|
||||
func (c Collection[T]) upsertItem(tx *Snapshot, itemID uint64, data []byte) error {
|
||||
func (c *Collection[T]) upsertItem(tx *Snapshot, itemID uint64, data []byte) error {
|
||||
item, ok := c.getByID(tx, itemID)
|
||||
if ok {
|
||||
tx.delete(c.collectionID, itemID)
|
||||
@@ -334,14 +418,14 @@ func (c Collection[T]) upsertItem(tx *Snapshot, itemID uint64, data []byte) erro
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c Collection[T]) getID(t *T) uint64 {
|
||||
func (c *Collection[T]) getID(t *T) uint64 {
|
||||
return *((*uint64)(unsafe.Pointer(t)))
|
||||
}
|
||||
|
||||
func (c Collection[T]) setID(t *T, id uint64) {
|
||||
func (c *Collection[T]) setID(t *T, id uint64) {
|
||||
*((*uint64)(unsafe.Pointer(t))) = id
|
||||
}
|
||||
|
||||
func (c Collection[T]) getState(tx *Snapshot) *collectionState[T] {
|
||||
func (c *Collection[T]) getState(tx *Snapshot) *collectionState[T] {
|
||||
return tx.collections[c.collectionID].(*collectionState[T])
|
||||
}
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"log"
|
||||
"os"
|
||||
"os/exec"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func TestCrashConsistency(t *testing.T) {
|
||||
|
||||
@@ -1,67 +0,0 @@
|
||||
package mdb
|
||||
|
||||
/*
|
||||
func (db *Database) openPrimary() (err error) {
|
||||
wal, err := cwal.Open(db.walRootDir, cwal.Config{
|
||||
SegMinCount: db.conf.WALSegMinCount,
|
||||
SegMaxAgeSec: db.conf.WALSegMaxAgeSec,
|
||||
})
|
||||
|
||||
pFile, err := pfile.Open(db.pageFilePath,
|
||||
|
||||
pFile, err := openPageFileAndReplayWAL(db.rootDir)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer pFile.Close()
|
||||
|
||||
pfHeader, err := pFile.ReadHeader()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
tx := db.Snapshot()
|
||||
tx.seqNum = pfHeader.SeqNum
|
||||
tx.updatedAt = pfHeader.UpdatedAt
|
||||
|
||||
pIndex, err := pagefile.NewIndex(pFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = pFile.IterateAllocated(pIndex, func(cID, iID uint64, data []byte) error {
|
||||
return db.loadItem(tx, cID, iID, data)
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
w, err := cwal.OpenWriter(db.walRootDir, &cwal.WriterConfig{
|
||||
SegMinCount: db.conf.WALSegMinCount,
|
||||
SegMaxAgeSec: db.conf.WALSegMaxAgeSec,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
db.done.Add(1)
|
||||
go txAggregator{
|
||||
Stop: db.stop,
|
||||
Done: db.done,
|
||||
ModChan: db.modChan,
|
||||
W: w,
|
||||
Index: pIndex,
|
||||
Snapshot: db.snapshot,
|
||||
}.Run()
|
||||
|
||||
db.done.Add(1)
|
||||
go (&fileWriter{
|
||||
Stop: db.stop,
|
||||
Done: db.done,
|
||||
PageFilePath: db.pageFilePath,
|
||||
WALRootDir: db.walRootDir,
|
||||
}).Run()
|
||||
|
||||
return nil
|
||||
}
|
||||
*/
|
||||
@@ -1,13 +1,14 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"git.crumpington.com/public/jldb/mdb/pfile"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
)
|
||||
|
||||
func (db *Database) repSendState(conn net.Conn) error {
|
||||
@@ -98,6 +99,7 @@ func (db *Database) repApply(rec wal.Record) (err error) {
|
||||
}
|
||||
tx.seqNum = rec.SeqNum
|
||||
tx.timestampMS = rec.TimestampMS
|
||||
tx.setReadOnly()
|
||||
db.snapshot.Store(tx)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,129 +0,0 @@
|
||||
package mdb
|
||||
|
||||
/*
|
||||
func (db *Database) openSecondary() (err error) {
|
||||
if db.shouldLoadFromPrimary() {
|
||||
if err := db.loadFromPrimary(); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
log.Printf("Opening page-file...")
|
||||
|
||||
pFile, err := openPageFileAndReplayWAL(db.rootDir)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer pFile.Close()
|
||||
|
||||
pfHeader, err := pFile.ReadHeader()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
log.Printf("Building page-file index...")
|
||||
|
||||
pIndex, err := pagefile.NewIndex(pFile)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
tx := db.Snapshot()
|
||||
tx.seqNum = pfHeader.SeqNum
|
||||
tx.updatedAt = pfHeader.UpdatedAt
|
||||
|
||||
log.Printf("Loading data into memory...")
|
||||
|
||||
err = pFile.IterateAllocated(pIndex, func(cID, iID uint64, data []byte) error {
|
||||
return db.loadItem(tx, cID, iID, data)
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
log.Printf("Creating writer...")
|
||||
|
||||
w, err := cswal.OpenWriter(db.walRootDir, &cswal.WriterConfig{
|
||||
SegMinCount: db.conf.WALSegMinCount,
|
||||
SegMaxAgeSec: db.conf.WALSegMaxAgeSec,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
db.done.Add(1)
|
||||
go (&walFollower{
|
||||
Stop: db.stop,
|
||||
Done: db.done,
|
||||
W: w,
|
||||
Client: NewClient(db.conf.PrimaryURL, db.conf.ReplicationPSK, db.conf.NetTimeout),
|
||||
}).Run()
|
||||
|
||||
db.done.Add(1)
|
||||
go (&follower{
|
||||
Stop: db.stop,
|
||||
Done: db.done,
|
||||
WALRootDir: db.walRootDir,
|
||||
SeqNum: pfHeader.SeqNum,
|
||||
ApplyChanges: db.applyChanges,
|
||||
}).Run()
|
||||
|
||||
db.done.Add(1)
|
||||
go (&fileWriter{
|
||||
Stop: db.stop,
|
||||
Done: db.done,
|
||||
PageFilePath: db.pageFilePath,
|
||||
WALRootDir: db.walRootDir,
|
||||
}).Run()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (db *Database) shouldLoadFromPrimary() bool {
|
||||
if _, err := os.Stat(db.walRootDir); os.IsNotExist(err) {
|
||||
log.Printf("WAL doesn't exist.")
|
||||
return true
|
||||
}
|
||||
if _, err := os.Stat(db.pageFilePath); os.IsNotExist(err) {
|
||||
log.Printf("Page-file doesn't exist.")
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func (db *Database) loadFromPrimary() error {
|
||||
client := NewClient(db.conf.PrimaryURL, db.conf.ReplicationPSK, db.conf.NetTimeout)
|
||||
defer client.Disconnect()
|
||||
|
||||
log.Printf("Loading data from primary...")
|
||||
|
||||
if err := os.RemoveAll(db.pageFilePath); err != nil {
|
||||
log.Printf("Failed to remove page-file: %s", err)
|
||||
return errs.IO.WithErr(err) // Caller can retry.
|
||||
}
|
||||
|
||||
if err := os.RemoveAll(db.walRootDir); err != nil {
|
||||
log.Printf("Failed to remove WAL: %s", err)
|
||||
return errs.IO.WithErr(err) // Caller can retry.
|
||||
}
|
||||
|
||||
err := client.DownloadPageFile(db.pageFilePath+".tmp", db.pageFilePath)
|
||||
if err != nil {
|
||||
log.Printf("Failed to get page-file from primary: %s", err)
|
||||
return err // Caller can retry.
|
||||
}
|
||||
|
||||
pfHeader, err := pagefile.ReadHeader(db.pageFilePath)
|
||||
if err != nil {
|
||||
log.Printf("Failed to read page-file sequence number: %s", err)
|
||||
return err // Caller can retry.
|
||||
}
|
||||
|
||||
if err = cswal.CreateEx(db.walRootDir, pfHeader.SeqNum+1); err != nil {
|
||||
log.Printf("Failed to initialize WAL: %s", err)
|
||||
return err // Caller can retry.
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
*/
|
||||
@@ -6,6 +6,8 @@ import (
|
||||
"reflect"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
type DBTestCase struct {
|
||||
@@ -52,9 +54,9 @@ var testDBTestCases = []DBTestCase{{
|
||||
Name: "Update",
|
||||
|
||||
Update: func(t *testing.T, db TestDB, tx *Snapshot) error {
|
||||
user, ok := db.Users.ByID.Get(tx, &User{ID: 1})
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
user := db.Users.ByID.Get(tx, &User{ID: 1})
|
||||
if user == nil {
|
||||
return errs.NotFound
|
||||
}
|
||||
user.Name = "Bob"
|
||||
user.Email = "b@c.com"
|
||||
@@ -111,7 +113,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Insert(tx, user2)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrDuplicate,
|
||||
ExpectedUpdateError: errs.Duplicate,
|
||||
|
||||
State: DBState{},
|
||||
}},
|
||||
@@ -131,7 +133,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Insert(tx, user2)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrDuplicate,
|
||||
ExpectedUpdateError: errs.Duplicate,
|
||||
|
||||
State: DBState{},
|
||||
}},
|
||||
@@ -162,7 +164,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Insert(tx, user)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrDuplicate,
|
||||
ExpectedUpdateError: errs.Duplicate,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{{ID: 1, Name: "Alice", Email: "a@b.com"}},
|
||||
@@ -197,7 +199,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Insert(tx, user)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrDuplicate,
|
||||
ExpectedUpdateError: errs.Duplicate,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{{ID: 1, Name: "Alice", Email: "a@b.com"}},
|
||||
@@ -218,7 +220,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Insert(db.Snapshot(), user)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrReadOnly,
|
||||
ExpectedUpdateError: errs.ReadOnly,
|
||||
}},
|
||||
}, {
|
||||
|
||||
@@ -290,7 +292,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Update(tx, user)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrNotFound,
|
||||
ExpectedUpdateError: errs.NotFound,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{{ID: 5, Name: "Alice", Email: "a@b.com"}},
|
||||
@@ -321,16 +323,16 @@ var testDBTestCases = []DBTestCase{{
|
||||
Name: "Update",
|
||||
|
||||
Update: func(t *testing.T, db TestDB, tx *Snapshot) error {
|
||||
user, ok := db.Users.ByID.Get(tx, &User{ID: 1})
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
user := db.Users.ByID.Get(tx, &User{ID: 1})
|
||||
if user == nil {
|
||||
return errs.NotFound
|
||||
}
|
||||
user.Name = "Bob"
|
||||
user.Email = "b@c.com"
|
||||
return db.Users.Update(db.Snapshot(), user)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrReadOnly,
|
||||
ExpectedUpdateError: errs.ReadOnly,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{{ID: 1, Name: "Alice", Email: "a@b.com"}},
|
||||
@@ -451,7 +453,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Update(tx, user2)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrDuplicate,
|
||||
ExpectedUpdateError: errs.Duplicate,
|
||||
|
||||
State: DBState{},
|
||||
}},
|
||||
@@ -491,16 +493,16 @@ var testDBTestCases = []DBTestCase{{
|
||||
Name: "Update",
|
||||
|
||||
Update: func(t *testing.T, db TestDB, tx *Snapshot) error {
|
||||
u, ok := db.Users.ByID.Get(tx, &User{ID: 2})
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
u := db.Users.ByID.Get(tx, &User{ID: 2})
|
||||
if u == nil {
|
||||
return errs.NotFound
|
||||
}
|
||||
|
||||
u.Email = "a@b.com"
|
||||
return db.Users.Update(tx, u)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrDuplicate,
|
||||
ExpectedUpdateError: errs.Duplicate,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{
|
||||
@@ -542,7 +544,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Delete(db.Snapshot(), 1)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrReadOnly,
|
||||
ExpectedUpdateError: errs.ReadOnly,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{{ID: 1, Name: "Alice", Email: "a@b.com"}},
|
||||
@@ -575,7 +577,7 @@ var testDBTestCases = []DBTestCase{{
|
||||
return db.Users.Delete(tx, 2)
|
||||
},
|
||||
|
||||
ExpectedUpdateError: ErrNotFound,
|
||||
ExpectedUpdateError: errs.NotFound,
|
||||
|
||||
State: DBState{
|
||||
UsersByID: []User{{ID: 1, Name: "Alice", Email: "a@b.com"}},
|
||||
@@ -607,17 +609,17 @@ var testDBTestCases = []DBTestCase{{
|
||||
Update: func(t *testing.T, db TestDB, tx *Snapshot) error {
|
||||
expected := &User{ID: 1, Name: "Alice", Email: "a@b.com"}
|
||||
|
||||
u, ok := db.Users.ByID.Get(tx, &User{ID: 1})
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
u := db.Users.ByID.Get(tx, &User{ID: 1})
|
||||
if u == nil {
|
||||
return errs.NotFound
|
||||
}
|
||||
if !reflect.DeepEqual(u, expected) {
|
||||
return errors.New("Not equal (id)")
|
||||
}
|
||||
|
||||
u, ok = db.Users.ByEmail.Get(tx, &User{Email: "a@b.com"})
|
||||
if !ok {
|
||||
return ErrNotFound
|
||||
u = db.Users.ByEmail.Get(tx, &User{Email: "a@b.com"})
|
||||
if u == nil {
|
||||
return errs.NotFound
|
||||
}
|
||||
if !reflect.DeepEqual(u, expected) {
|
||||
return errors.New("Not equal (email)")
|
||||
@@ -635,11 +637,11 @@ var testDBTestCases = []DBTestCase{{
|
||||
Name: "Get not found",
|
||||
|
||||
Update: func(t *testing.T, db TestDB, tx *Snapshot) error {
|
||||
if _, ok := db.Users.ByID.Get(tx, &User{ID: 2}); ok {
|
||||
if u := db.Users.ByID.Get(tx, &User{ID: 2}); u != nil {
|
||||
return errors.New("Found (id)")
|
||||
}
|
||||
|
||||
if _, ok := db.Users.ByEmail.Get(tx, &User{Email: "x@b.com"}); ok {
|
||||
if u := db.Users.ByEmail.Get(tx, &User{Email: "x@b.com"}); u != nil {
|
||||
return errors.New("Found (email)")
|
||||
}
|
||||
|
||||
@@ -741,29 +743,25 @@ var testDBTestCases = []DBTestCase{{
|
||||
|
||||
first := true
|
||||
pivot := User{Name: "User1"}
|
||||
db.Users.ByName.AscendAfter(tx, &pivot, func(u *User) bool {
|
||||
for u := range db.Users.ByName.AscendAfter(tx, &pivot) {
|
||||
u.Name += "Mod"
|
||||
if err = db.Users.Update(tx, u); err != nil {
|
||||
return false
|
||||
return err
|
||||
}
|
||||
if first {
|
||||
first = false
|
||||
return true
|
||||
continue
|
||||
}
|
||||
|
||||
prev, ok := db.Users.ByID.Get(tx, &User{ID: u.ID - 1})
|
||||
if !ok {
|
||||
err = errors.New("Previous user not found")
|
||||
return false
|
||||
prev := db.Users.ByID.Get(tx, &User{ID: u.ID - 1})
|
||||
if prev == nil {
|
||||
return errors.New("Previous user not found")
|
||||
}
|
||||
|
||||
if !strings.HasSuffix(prev.Name, "Mod") {
|
||||
err = errors.New("Incorrect user name: " + prev.Name)
|
||||
return false
|
||||
return errors.New("Incorrect user name: " + prev.Name)
|
||||
}
|
||||
|
||||
return true
|
||||
})
|
||||
}
|
||||
return nil
|
||||
},
|
||||
|
||||
@@ -799,29 +797,26 @@ var testDBTestCases = []DBTestCase{{
|
||||
}
|
||||
|
||||
first := true
|
||||
db.Users.ByName.DescendAfter(tx, &User{Name: "User5Mod"}, func(u *User) bool {
|
||||
for u := range db.Users.ByName.DescendAfter(tx, &User{Name: "User5Mod"}) {
|
||||
u.Name = strings.TrimSuffix(u.Name, "Mod")
|
||||
if err = db.Users.Update(tx, u); err != nil {
|
||||
return false
|
||||
return err
|
||||
}
|
||||
if first {
|
||||
first = false
|
||||
return true
|
||||
continue
|
||||
}
|
||||
|
||||
prev, ok := db.Users.ByID.Get(tx, &User{ID: u.ID + 1})
|
||||
if !ok {
|
||||
err = errors.New("Previous user not found")
|
||||
return false
|
||||
prev := db.Users.ByID.Get(tx, &User{ID: u.ID + 1})
|
||||
if prev == nil {
|
||||
return errors.New("Previous user not found")
|
||||
}
|
||||
|
||||
if strings.HasSuffix(prev.Name, "Mod") {
|
||||
err = errors.New("Incorrect user name: " + prev.Name)
|
||||
return false
|
||||
return errors.New("Incorrect user name: " + prev.Name)
|
||||
}
|
||||
}
|
||||
|
||||
return true
|
||||
})
|
||||
return nil
|
||||
},
|
||||
|
||||
|
||||
@@ -1,138 +0,0 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestDBList(t *testing.T) {
|
||||
db := NewTestDBPrimary(t, t.TempDir())
|
||||
|
||||
var (
|
||||
user1 = User{
|
||||
ID: NewID(),
|
||||
Name: "User1",
|
||||
Email: "user1@gmail.com",
|
||||
}
|
||||
|
||||
user2 = User{
|
||||
ID: NewID(),
|
||||
Name: "User2",
|
||||
Email: "user2@gmail.com",
|
||||
}
|
||||
|
||||
user3 = User{
|
||||
ID: NewID(),
|
||||
Name: "User3",
|
||||
Email: "user3@gmail.com",
|
||||
}
|
||||
user1Data = make([]UserDataItem, 10)
|
||||
user2Data = make([]UserDataItem, 4)
|
||||
user3Data = make([]UserDataItem, 8)
|
||||
)
|
||||
|
||||
err := db.Update(func(tx *Snapshot) error {
|
||||
if err := db.Users.Insert(tx, &user1); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := db.Users.Insert(tx, &user2); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for i := range user1Data {
|
||||
user1Data[i] = UserDataItem{
|
||||
ID: NewID(),
|
||||
UserID: user1.ID,
|
||||
Name: fmt.Sprintf("Name1: %d", i),
|
||||
Data: fmt.Sprintf("Data: %d", i),
|
||||
}
|
||||
|
||||
if err := db.UserData.Insert(tx, &user1Data[i]); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
for i := range user2Data {
|
||||
user2Data[i] = UserDataItem{
|
||||
ID: NewID(),
|
||||
UserID: user2.ID,
|
||||
Name: fmt.Sprintf("Name2: %d", i),
|
||||
Data: fmt.Sprintf("Data: %d", i),
|
||||
}
|
||||
|
||||
if err := db.UserData.Insert(tx, &user2Data[i]); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
for i := range user3Data {
|
||||
user3Data[i] = UserDataItem{
|
||||
ID: NewID(),
|
||||
UserID: user3.ID,
|
||||
Name: fmt.Sprintf("Name3: %d", i),
|
||||
Data: fmt.Sprintf("Data: %d", i),
|
||||
}
|
||||
|
||||
if err := db.UserData.Insert(tx, &user3Data[i]); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
type TestCase struct {
|
||||
Name string
|
||||
Args ListArgs[UserDataItem]
|
||||
Expected []UserDataItem
|
||||
}
|
||||
|
||||
cases := []TestCase{
|
||||
{
|
||||
Name: "User1 all",
|
||||
Args: ListArgs[UserDataItem]{
|
||||
After: &UserDataItem{
|
||||
UserID: user1.ID,
|
||||
},
|
||||
While: func(item *UserDataItem) bool {
|
||||
return item.UserID == user1.ID
|
||||
},
|
||||
},
|
||||
Expected: user1Data,
|
||||
}, {
|
||||
Name: "User1 limited",
|
||||
Args: ListArgs[UserDataItem]{
|
||||
After: &UserDataItem{
|
||||
UserID: user1.ID,
|
||||
},
|
||||
While: func(item *UserDataItem) bool {
|
||||
return item.UserID == user1.ID
|
||||
},
|
||||
Limit: 4,
|
||||
},
|
||||
Expected: user1Data[:4],
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.Name, func(t *testing.T) {
|
||||
tx := db.Snapshot()
|
||||
l := db.UserData.ByName.List(tx, tc.Args, nil)
|
||||
if len(l) != len(tc.Expected) {
|
||||
t.Fatal(tc.Name, l)
|
||||
}
|
||||
|
||||
for i := range l {
|
||||
if !reflect.DeepEqual(*l[i], tc.Expected[i]) {
|
||||
t.Fatal(tc.Name, l)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -72,7 +72,7 @@ func testRunner_testCase(t *testing.T, testCase DBTestCase) {
|
||||
}
|
||||
|
||||
// TODO: Why is this necessary?
|
||||
time.Sleep(time.Second)
|
||||
//time.Sleep(time.Second)
|
||||
finalStep := testCase.Steps[len(testCase.Steps)-1]
|
||||
|
||||
secondarySnapshot := db2.Snapshot()
|
||||
@@ -134,23 +134,23 @@ func checkSlicesEqual[T any](t *testing.T, name string, actual, expected []T) {
|
||||
}
|
||||
}
|
||||
|
||||
func checkMinMaxEqual[T any](t *testing.T, name string, tx *Snapshot, index Index[T], expected []T) {
|
||||
func checkMinMaxEqual[T any](t *testing.T, name string, tx *Snapshot, index *Index[T], expected []T) {
|
||||
if len(expected) == 0 {
|
||||
if min, ok := index.Min(tx); ok {
|
||||
if min := index.Min(tx); min != nil {
|
||||
t.Fatal(min)
|
||||
}
|
||||
if max, ok := index.Max(tx); ok {
|
||||
if max := index.Max(tx); max != nil {
|
||||
t.Fatal(max)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
min, ok := index.Min(tx)
|
||||
if !ok {
|
||||
min := index.Min(tx)
|
||||
if min == nil {
|
||||
t.Fatal("No min")
|
||||
}
|
||||
max, ok := index.Max(tx)
|
||||
if !ok {
|
||||
max := index.Max(tx)
|
||||
if max == nil {
|
||||
t.Fatal("No max")
|
||||
}
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ package mdb
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
)
|
||||
|
||||
|
||||
@@ -14,7 +14,7 @@ type UserDataItem struct {
|
||||
|
||||
type UserData struct {
|
||||
*Collection[UserDataItem]
|
||||
ByName Index[UserDataItem] // Unique index on (Token).
|
||||
ByName *Index[UserDataItem] // Unique index on (Token).
|
||||
}
|
||||
|
||||
func NewUserDataCollection(db *Database) UserData {
|
||||
|
||||
@@ -12,9 +12,9 @@ type User struct {
|
||||
|
||||
type Users struct {
|
||||
*Collection[User]
|
||||
ByEmail Index[User] // Unique index on (Email).
|
||||
ByName Index[User] // Index on (Name).
|
||||
ByBlocked Index[User] // Partial index on (Blocked,Email).
|
||||
ByEmail *Index[User] // Unique index on (Email).
|
||||
ByName *Index[User] // Index on (Name).
|
||||
ByBlocked *Index[User] // Partial index on (Blocked,Email).
|
||||
}
|
||||
|
||||
func NewUserCollection(db *Database) Users {
|
||||
|
||||
13
mdb/db.go
13
mdb/db.go
@@ -2,15 +2,16 @@ package mdb
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/rep"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"git.crumpington.com/public/jldb/mdb/pfile"
|
||||
"net/http"
|
||||
"os"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/lib/rep"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"git.crumpington.com/public/jldb/mdb/pfile"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
@@ -72,6 +73,10 @@ type Database struct {
|
||||
}
|
||||
|
||||
func New(conf Config) *Database {
|
||||
if conf.NetTimeout <= 0 {
|
||||
conf.NetTimeout = time.Minute
|
||||
}
|
||||
|
||||
if conf.MaxConcurrentUpdates <= 0 {
|
||||
conf.MaxConcurrentUpdates = 32
|
||||
}
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
@@ -20,18 +19,16 @@ func (i Index[T]) AssertEqual(t *testing.T, tx1, tx2 *Snapshot) {
|
||||
}
|
||||
|
||||
errStr := ""
|
||||
i.Ascend(tx1, func(item1 *T) bool {
|
||||
item2, ok := i.Get(tx2, item1)
|
||||
if !ok {
|
||||
errStr = fmt.Sprintf("Indices don't match. %v not found.", item1)
|
||||
return false
|
||||
iter := i.Ascend(tx1)
|
||||
for item1 := range iter {
|
||||
item2 := i.Get(tx2, item1)
|
||||
if item2 == nil {
|
||||
t.Fatalf("Indices don't match. %v not found.", item1)
|
||||
}
|
||||
if !reflect.DeepEqual(item1, item2) {
|
||||
errStr = fmt.Sprintf("%v != %v", item1, item2)
|
||||
return false
|
||||
t.Fatalf("%v != %v", item1, item2)
|
||||
}
|
||||
return true
|
||||
})
|
||||
}
|
||||
|
||||
if errStr != "" {
|
||||
t.Fatal(errStr)
|
||||
|
||||
@@ -1,11 +0,0 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrNotFound = errs.NotFound
|
||||
ErrReadOnly = errs.ReadOnly
|
||||
ErrDuplicate = errs.Duplicate
|
||||
)
|
||||
163
mdb/index.go
163
mdb/index.go
@@ -1,6 +1,7 @@
|
||||
package mdb
|
||||
|
||||
import (
|
||||
"iter"
|
||||
"unsafe"
|
||||
|
||||
"github.com/google/btree"
|
||||
@@ -10,7 +11,7 @@ func NewIndex[T any](
|
||||
c *Collection[T],
|
||||
name string,
|
||||
compare func(lhs, rhs *T) int,
|
||||
) Index[T] {
|
||||
) *Index[T] {
|
||||
return c.addIndex(indexConfig[T]{
|
||||
Name: name,
|
||||
Unique: false,
|
||||
@@ -24,7 +25,7 @@ func NewPartialIndex[T any](
|
||||
name string,
|
||||
compare func(lhs, rhs *T) int,
|
||||
include func(*T) bool,
|
||||
) Index[T] {
|
||||
) *Index[T] {
|
||||
return c.addIndex(indexConfig[T]{
|
||||
Name: name,
|
||||
Unique: false,
|
||||
@@ -37,7 +38,7 @@ func NewUniqueIndex[T any](
|
||||
c *Collection[T],
|
||||
name string,
|
||||
compare func(lhs, rhs *T) int,
|
||||
) Index[T] {
|
||||
) *Index[T] {
|
||||
return c.addIndex(indexConfig[T]{
|
||||
Name: name,
|
||||
Unique: true,
|
||||
@@ -51,7 +52,7 @@ func NewUniquePartialIndex[T any](
|
||||
name string,
|
||||
compare func(lhs, rhs *T) int,
|
||||
include func(*T) bool,
|
||||
) Index[T] {
|
||||
) *Index[T] {
|
||||
return c.addIndex(indexConfig[T]{
|
||||
Name: name,
|
||||
Unique: true,
|
||||
@@ -63,6 +64,7 @@ func NewUniquePartialIndex[T any](
|
||||
// ----------------------------------------------------------------------------
|
||||
|
||||
type Index[T any] struct {
|
||||
db *Database
|
||||
name string
|
||||
collectionID uint64
|
||||
indexID uint64
|
||||
@@ -70,117 +72,94 @@ type Index[T any] struct {
|
||||
copy func(*T) *T
|
||||
}
|
||||
|
||||
func (i Index[T]) Get(tx *Snapshot, in *T) (item *T, ok bool) {
|
||||
tPtr, ok := i.get(tx, in)
|
||||
if !ok {
|
||||
return item, false
|
||||
func (i *Index[T]) ensureSnapshot(tx *Snapshot) *Snapshot {
|
||||
if tx == nil {
|
||||
tx = i.db.Snapshot()
|
||||
}
|
||||
return i.copy(tPtr), true
|
||||
return tx
|
||||
}
|
||||
|
||||
func (i Index[T]) get(tx *Snapshot, in *T) (*T, bool) {
|
||||
func (i *Index[T]) Get(tx *Snapshot, in *T) *T {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
if tPtr, ok := i.get(tx, in); ok {
|
||||
return i.copy(tPtr)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (i *Index[T]) get(tx *Snapshot, in *T) (*T, bool) {
|
||||
return i.btree(tx).Get(in)
|
||||
}
|
||||
|
||||
func (i Index[T]) Has(tx *Snapshot, in *T) bool {
|
||||
func (i *Index[T]) Has(tx *Snapshot, in *T) bool {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
return i.btree(tx).Has(in)
|
||||
}
|
||||
|
||||
func (i Index[T]) Min(tx *Snapshot) (item *T, ok bool) {
|
||||
tPtr, ok := i.btree(tx).Min()
|
||||
if !ok {
|
||||
return item, false
|
||||
func (i *Index[T]) Min(tx *Snapshot) *T {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
if tPtr, ok := i.btree(tx).Min(); ok {
|
||||
return i.copy(tPtr)
|
||||
}
|
||||
return i.copy(tPtr), true
|
||||
return nil
|
||||
}
|
||||
|
||||
func (i Index[T]) Max(tx *Snapshot) (item *T, ok bool) {
|
||||
tPtr, ok := i.btree(tx).Max()
|
||||
if !ok {
|
||||
return item, false
|
||||
func (i *Index[T]) Max(tx *Snapshot) *T {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
if tPtr, ok := i.btree(tx).Max(); ok {
|
||||
return i.copy(tPtr)
|
||||
}
|
||||
return i.copy(tPtr), true
|
||||
return nil
|
||||
}
|
||||
|
||||
func (i Index[T]) Ascend(tx *Snapshot, each func(*T) bool) {
|
||||
i.btreeForIter(tx).Ascend(func(t *T) bool {
|
||||
return each(i.copy(t))
|
||||
})
|
||||
func (i *Index[T]) Ascend(tx *Snapshot) iter.Seq[*T] {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
return func(yield func(*T) bool) {
|
||||
i.btreeForIter(tx).Ascend(func(t *T) bool {
|
||||
return yield(i.copy(t))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func (i Index[T]) AscendAfter(tx *Snapshot, after *T, each func(*T) bool) {
|
||||
i.btreeForIter(tx).AscendGreaterOrEqual(after, func(t *T) bool {
|
||||
return each(i.copy(t))
|
||||
})
|
||||
func (i *Index[T]) AscendAfter(tx *Snapshot, after *T) iter.Seq[*T] {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
return func(yield func(*T) bool) {
|
||||
i.btreeForIter(tx).AscendGreaterOrEqual(after, func(t *T) bool {
|
||||
return yield(i.copy(t))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func (i Index[T]) Descend(tx *Snapshot, each func(*T) bool) {
|
||||
i.btreeForIter(tx).Descend(func(t *T) bool {
|
||||
return each(i.copy(t))
|
||||
})
|
||||
func (i *Index[T]) Descend(tx *Snapshot) iter.Seq[*T] {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
return func(yield func(*T) bool) {
|
||||
i.btreeForIter(tx).Descend(func(t *T) bool {
|
||||
return yield(i.copy(t))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func (i Index[T]) DescendAfter(tx *Snapshot, after *T, each func(*T) bool) {
|
||||
i.btreeForIter(tx).DescendLessOrEqual(after, func(t *T) bool {
|
||||
return each(i.copy(t))
|
||||
})
|
||||
func (i *Index[T]) DescendAfter(tx *Snapshot, after *T) iter.Seq[*T] {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
return func(yield func(*T) bool) {
|
||||
i.btreeForIter(tx).DescendLessOrEqual(after, func(t *T) bool {
|
||||
return yield(i.copy(t))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
type ListArgs[T any] struct {
|
||||
Desc bool // True for descending order, otherwise ascending.
|
||||
After *T // If after is given, iterate after (and including) the value.
|
||||
While func(*T) bool // Continue iterating until While is false.
|
||||
Limit int // Maximum number of items to return. 0 => All.
|
||||
}
|
||||
|
||||
func (i Index[T]) List(tx *Snapshot, args ListArgs[T], out []*T) []*T {
|
||||
if args.Limit < 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
if args.While == nil {
|
||||
args.While = func(*T) bool { return true }
|
||||
}
|
||||
|
||||
size := args.Limit
|
||||
if size == 0 {
|
||||
size = 32 // Why not?
|
||||
}
|
||||
|
||||
items := out[:0]
|
||||
|
||||
each := func(item *T) bool {
|
||||
if !args.While(item) {
|
||||
return false
|
||||
}
|
||||
items = append(items, item)
|
||||
return args.Limit == 0 || len(items) < args.Limit
|
||||
}
|
||||
|
||||
if args.Desc {
|
||||
if args.After != nil {
|
||||
i.DescendAfter(tx, args.After, each)
|
||||
} else {
|
||||
i.Descend(tx, each)
|
||||
}
|
||||
} else {
|
||||
if args.After != nil {
|
||||
i.AscendAfter(tx, args.After, each)
|
||||
} else {
|
||||
i.Ascend(tx, each)
|
||||
}
|
||||
}
|
||||
|
||||
return items
|
||||
func (i *Index[T]) Count(tx *Snapshot) int {
|
||||
tx = i.ensureSnapshot(tx)
|
||||
return i.btree(tx).Len()
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
|
||||
func (i Index[T]) insertConflict(tx *Snapshot, item *T) bool {
|
||||
func (i *Index[T]) insertConflict(tx *Snapshot, item *T) bool {
|
||||
return i.btree(tx).Has(item)
|
||||
}
|
||||
|
||||
func (i Index[T]) updateConflict(tx *Snapshot, item *T) bool {
|
||||
func (i *Index[T]) updateConflict(tx *Snapshot, item *T) bool {
|
||||
current, ok := i.btree(tx).Get(item)
|
||||
return ok && i.getID(current) != i.getID(item)
|
||||
}
|
||||
@@ -188,7 +167,7 @@ func (i Index[T]) updateConflict(tx *Snapshot, item *T) bool {
|
||||
// This should only be called after insertConflict. Additionally, the caller
|
||||
// should ensure that the index has been properly cloned for write before
|
||||
// writing.
|
||||
func (i Index[T]) insert(tx *Snapshot, item *T) {
|
||||
func (i *Index[T]) insert(tx *Snapshot, item *T) {
|
||||
if i.include != nil && !i.include(item) {
|
||||
return
|
||||
}
|
||||
@@ -196,7 +175,7 @@ func (i Index[T]) insert(tx *Snapshot, item *T) {
|
||||
i.btree(tx).ReplaceOrInsert(item)
|
||||
}
|
||||
|
||||
func (i Index[T]) update(tx *Snapshot, old, new *T) {
|
||||
func (i *Index[T]) update(tx *Snapshot, old, new *T) {
|
||||
bt := i.btree(tx)
|
||||
bt.Delete(old)
|
||||
|
||||
@@ -204,22 +183,22 @@ func (i Index[T]) update(tx *Snapshot, old, new *T) {
|
||||
i.insert(tx, new)
|
||||
}
|
||||
|
||||
func (i Index[T]) delete(tx *Snapshot, item *T) {
|
||||
func (i *Index[T]) delete(tx *Snapshot, item *T) {
|
||||
i.btree(tx).Delete(item)
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
|
||||
func (i Index[T]) getState(tx *Snapshot) indexState[T] {
|
||||
func (i *Index[T]) getState(tx *Snapshot) indexState[T] {
|
||||
return tx.collections[i.collectionID].(*collectionState[T]).Indices[i.indexID]
|
||||
}
|
||||
|
||||
// Get the current btree for get/has/update/delete, etc.
|
||||
func (i Index[T]) btree(tx *Snapshot) *btree.BTreeG[*T] {
|
||||
func (i *Index[T]) btree(tx *Snapshot) *btree.BTreeG[*T] {
|
||||
return i.getState(tx).BTree
|
||||
}
|
||||
|
||||
func (i Index[T]) btreeForIter(tx *Snapshot) *btree.BTreeG[*T] {
|
||||
func (i *Index[T]) btreeForIter(tx *Snapshot) *btree.BTreeG[*T] {
|
||||
cState := tx.collections[i.collectionID].(*collectionState[T])
|
||||
bt := cState.Indices[i.indexID].BTree
|
||||
|
||||
@@ -231,6 +210,6 @@ func (i Index[T]) btreeForIter(tx *Snapshot) *btree.BTreeG[*T] {
|
||||
return bt
|
||||
}
|
||||
|
||||
func (i Index[T]) getID(t *T) uint64 {
|
||||
func (i *Index[T]) getID(t *T) uint64 {
|
||||
return *((*uint64)(unsafe.Pointer(t)))
|
||||
}
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
package mdb
|
||||
|
||||
func (i Index[T]) Dump(tx *Snapshot) (l []T) {
|
||||
i.Ascend(tx, func(t *T) bool {
|
||||
for t := range i.Ascend(tx) {
|
||||
l = append(l, *t)
|
||||
return true
|
||||
})
|
||||
}
|
||||
return l
|
||||
}
|
||||
|
||||
@@ -2,8 +2,9 @@ package pfile
|
||||
|
||||
import (
|
||||
crand "crypto/rand"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"math/rand"
|
||||
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
)
|
||||
|
||||
func randomChangeList() (changes []change.Change) {
|
||||
|
||||
@@ -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 {
|
||||
out = out[:0]
|
||||
|
||||
|
||||
@@ -13,14 +13,19 @@ type Index struct {
|
||||
}
|
||||
|
||||
func NewIndex(f *File) (*Index, error) {
|
||||
firstPage, err := f.pageCount()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
idx := &Index{
|
||||
fList: newFreeList(0),
|
||||
fList: newFreeList(firstPage),
|
||||
aList: *newAllocList(),
|
||||
seen: map[[2]uint64]struct{}{},
|
||||
mask: []bool{},
|
||||
}
|
||||
|
||||
err := f.iterate(func(pageID uint64, page dataPage) error {
|
||||
err = f.iterate(func(pageID uint64, page dataPage) error {
|
||||
header := page.Header()
|
||||
switch header.PageType {
|
||||
case pageTypeHead:
|
||||
|
||||
@@ -3,10 +3,11 @@ package pfile
|
||||
import (
|
||||
"bytes"
|
||||
crand "crypto/rand"
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
)
|
||||
|
||||
func newForTesting(t *testing.T) (*File, *Index) {
|
||||
|
||||
@@ -2,8 +2,9 @@ package pfile
|
||||
|
||||
import (
|
||||
"hash/crc32"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"unsafe"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
// ----------------------------------------------------------------------------
|
||||
@@ -30,7 +31,7 @@ var emptyPage = func() dataPage {
|
||||
type pageHeader struct {
|
||||
CRC uint32 // IEEE CRC-32 checksum.
|
||||
PageType uint32 // One of the PageType* constants.
|
||||
CollectionID uint64 //
|
||||
CollectionID uint64 //
|
||||
ItemID uint64
|
||||
DataSize uint64
|
||||
NextPage uint64
|
||||
|
||||
@@ -3,9 +3,10 @@ package pfile
|
||||
import (
|
||||
"bytes"
|
||||
crand "crypto/rand"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"math/rand"
|
||||
"testing"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
)
|
||||
|
||||
func randomPage(t *testing.T) dataPage {
|
||||
|
||||
@@ -6,12 +6,13 @@ import (
|
||||
"compress/gzip"
|
||||
"encoding/binary"
|
||||
"io"
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"net"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/errs"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
)
|
||||
|
||||
type File struct {
|
||||
@@ -133,6 +134,21 @@ func (pf *File) writePage(page dataPage, id uint64) error {
|
||||
// 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 {
|
||||
pf.lock.RLock()
|
||||
defer pf.lock.RUnlock()
|
||||
|
||||
@@ -2,6 +2,7 @@ package pfile
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
|
||||
"git.crumpington.com/public/jldb/lib/wal"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
)
|
||||
|
||||
@@ -3,8 +3,9 @@ package mdb
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
"sync/atomic"
|
||||
|
||||
"git.crumpington.com/public/jldb/mdb/change"
|
||||
)
|
||||
|
||||
type Snapshot struct {
|
||||
|
||||
@@ -133,8 +133,8 @@ func (db DataDB) modifyOnce() {
|
||||
func (db DataDB) ComputeCRC(tx *Snapshot) uint32 {
|
||||
h := crc32.NewIEEE()
|
||||
for dataID := uint64(1); dataID < 10; dataID++ {
|
||||
d, ok := db.Datas.ByID.Get(tx, &DataItem{ID: dataID})
|
||||
if !ok {
|
||||
d := db.Datas.ByID.Get(tx, &DataItem{ID: dataID})
|
||||
if d == nil {
|
||||
continue
|
||||
}
|
||||
h.Write(d.Data)
|
||||
@@ -143,8 +143,8 @@ func (db DataDB) ComputeCRC(tx *Snapshot) uint32 {
|
||||
}
|
||||
|
||||
func (db DataDB) ReadCRC(tx *Snapshot) uint32 {
|
||||
r, ok := db.CRCs.ByID.Get(tx, &CRCItem{ID: 1})
|
||||
if !ok {
|
||||
r := db.CRCs.ByID.Get(tx, &CRCItem{ID: 1})
|
||||
if r == nil {
|
||||
return 0
|
||||
}
|
||||
return r.CRC32
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"crypto/rand"
|
||||
"errors"
|
||||
"hash/crc32"
|
||||
"git.crumpington.com/public/jldb/mdb"
|
||||
"log"
|
||||
mrand "math/rand"
|
||||
"os"
|
||||
@@ -13,6 +12,8 @@ import (
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"git.crumpington.com/public/jldb/mdb"
|
||||
)
|
||||
|
||||
type DataItem struct {
|
||||
@@ -135,8 +136,8 @@ func (db DataDB) modifyOnce() {
|
||||
func (db DataDB) ComputeCRC(tx *mdb.Snapshot) uint32 {
|
||||
h := crc32.NewIEEE()
|
||||
for dataID := uint64(1); dataID < 10; dataID++ {
|
||||
d, ok := db.Datas.ByID.Get(tx, &DataItem{ID: dataID})
|
||||
if !ok {
|
||||
d := db.Datas.ByID.Get(tx, &DataItem{ID: dataID})
|
||||
if d == nil {
|
||||
continue
|
||||
}
|
||||
h.Write(d.Data)
|
||||
@@ -145,8 +146,8 @@ func (db DataDB) ComputeCRC(tx *mdb.Snapshot) uint32 {
|
||||
}
|
||||
|
||||
func (db DataDB) ReadCRC(tx *mdb.Snapshot) uint32 {
|
||||
r, ok := db.CRCs.ByID.Get(tx, &CRCItem{ID: 1})
|
||||
if !ok {
|
||||
r := db.CRCs.ByID.Get(tx, &CRCItem{ID: 1})
|
||||
if r == nil {
|
||||
return 0
|
||||
}
|
||||
return r.CRC32
|
||||
|
||||
@@ -1,92 +0,0 @@
|
||||
package mdb
|
||||
|
||||
/*
|
||||
type txAggregator struct {
|
||||
Stop chan struct{}
|
||||
Done *sync.WaitGroup
|
||||
ModChan chan txMod
|
||||
W *cswal.Writer
|
||||
Index *pagefile.Index
|
||||
Snapshot *atomic.Pointer[Snapshot]
|
||||
}
|
||||
|
||||
func (p txAggregator) Run() {
|
||||
defer p.Done.Done()
|
||||
defer p.W.Close()
|
||||
|
||||
var (
|
||||
tx *Snapshot
|
||||
mod txMod
|
||||
rec cswal.Record
|
||||
err error
|
||||
toNotify = make([]chan error, 0, 1024)
|
||||
)
|
||||
|
||||
READ_FIRST:
|
||||
|
||||
toNotify = toNotify[:0]
|
||||
|
||||
select {
|
||||
case mod = <-p.ModChan:
|
||||
goto BEGIN
|
||||
case <-p.Stop:
|
||||
goto END
|
||||
}
|
||||
|
||||
BEGIN:
|
||||
|
||||
tx = p.Snapshot.Load().begin()
|
||||
goto APPLY_MOD
|
||||
|
||||
CLONE:
|
||||
|
||||
tx = tx.clone()
|
||||
goto APPLY_MOD
|
||||
|
||||
APPLY_MOD:
|
||||
|
||||
if err = mod.Update(tx); err != nil {
|
||||
mod.Resp <- err
|
||||
goto ROLLBACK
|
||||
}
|
||||
|
||||
toNotify = append(toNotify, mod.Resp)
|
||||
goto NEXT
|
||||
|
||||
ROLLBACK:
|
||||
|
||||
if len(toNotify) == 0 {
|
||||
goto READ_FIRST
|
||||
}
|
||||
|
||||
tx = tx.rollback()
|
||||
goto NEXT
|
||||
|
||||
NEXT:
|
||||
|
||||
select {
|
||||
case mod = <-p.ModChan:
|
||||
goto CLONE
|
||||
default:
|
||||
goto WRITE
|
||||
}
|
||||
|
||||
WRITE:
|
||||
|
||||
rec, err = writeChangesToWAL(tx.changes, p.Index, p.W)
|
||||
if err == nil {
|
||||
tx.seqNum = rec.SeqNum
|
||||
tx.updatedAt = rec.CreatedAt
|
||||
tx.setReadOnly()
|
||||
p.Snapshot.Store(tx)
|
||||
}
|
||||
|
||||
for i := range toNotify {
|
||||
toNotify[i] <- err
|
||||
}
|
||||
|
||||
goto READ_FIRST
|
||||
|
||||
END:
|
||||
}
|
||||
*/
|
||||
Reference in New Issue
Block a user