Compare commits
	
		
			17 Commits
		
	
	
		
			ff8b87d6ea
			...
			v0.8.1
		
	
	| Author | SHA1 | Date | |
|---|---|---|---|
|  | 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,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,70 @@ | ||||
| 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" | ||||
| ) | ||||
|  | ||||
| // --------------------------------------------------------------------------- | ||||
|  | ||||
| // TODO: Remove this! | ||||
| 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) | ||||
| 	} | ||||
| } | ||||
|   | ||||
							
								
								
									
										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 { | ||||
| @@ -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 | ||||
|   | ||||
| @@ -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" | ||||
| ) | ||||
|  | ||||
|   | ||||
| @@ -5,9 +5,10 @@ import ( | ||||
| 	"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,10 +21,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] | ||||
| 	ByID *Index[T] | ||||
|  | ||||
| 	buf *bytes.Buffer | ||||
| } | ||||
| @@ -64,8 +65,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]{}, | ||||
| 		indices:       []*Index[T]{}, | ||||
| 		uniqueIndices: []*Index[T]{}, | ||||
| 		buf:           &bytes.Buffer{}, | ||||
| 	} | ||||
|  | ||||
| @@ -91,7 +92,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 +108,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 +133,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 +150,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 +192,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 +205,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 +228,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 +246,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 +334,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 +361,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 +370,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 +382,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 +399,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 +422,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 { | ||||
|   | ||||
| @@ -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)") | ||||
| 			} | ||||
|  | ||||
| @@ -751,8 +753,8 @@ var testDBTestCases = []DBTestCase{{ | ||||
| 					return true | ||||
| 				} | ||||
|  | ||||
| 				prev, ok := db.Users.ByID.Get(tx, &User{ID: u.ID - 1}) | ||||
| 				if !ok { | ||||
| 				prev := db.Users.ByID.Get(tx, &User{ID: u.ID - 1}) | ||||
| 				if prev == nil { | ||||
| 					err = errors.New("Previous user not found") | ||||
| 					return false | ||||
| 				} | ||||
| @@ -809,8 +811,8 @@ var testDBTestCases = []DBTestCase{{ | ||||
| 					return true | ||||
| 				} | ||||
|  | ||||
| 				prev, ok := db.Users.ByID.Get(tx, &User{ID: u.ID + 1}) | ||||
| 				if !ok { | ||||
| 				prev := db.Users.ByID.Get(tx, &User{ID: u.ID + 1}) | ||||
| 				if prev == nil { | ||||
| 					err = errors.New("Previous user not found") | ||||
| 					return false | ||||
| 				} | ||||
|   | ||||
| @@ -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) | ||||
| 				} | ||||
| 			} | ||||
| 		}) | ||||
| 	} | ||||
| } | ||||
| @@ -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 | ||||
| 	} | ||||
|   | ||||
| @@ -21,8 +21,8 @@ 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 { | ||||
| 		item2 := i.Get(tx2, item1) | ||||
| 		if item2 == nil { | ||||
| 			errStr = fmt.Sprintf("Indices don't match. %v not found.", item1) | ||||
| 			return false | ||||
| 		} | ||||
|   | ||||
| @@ -1,11 +0,0 @@ | ||||
| package mdb | ||||
|  | ||||
| import ( | ||||
| 	"git.crumpington.com/public/jldb/lib/errs" | ||||
| ) | ||||
|  | ||||
| var ( | ||||
| 	ErrNotFound  = errs.NotFound | ||||
| 	ErrReadOnly  = errs.ReadOnly | ||||
| 	ErrDuplicate = errs.Duplicate | ||||
| ) | ||||
							
								
								
									
										130
									
								
								mdb/index.go
									
									
									
									
									
								
							
							
						
						
									
										130
									
								
								mdb/index.go
									
									
									
									
									
								
							| @@ -10,7 +10,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 +24,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 +37,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 +51,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 +63,7 @@ func NewUniquePartialIndex[T any]( | ||||
| // ---------------------------------------------------------------------------- | ||||
|  | ||||
| type Index[T any] struct { | ||||
| 	db           *Database | ||||
| 	name         string | ||||
| 	collectionID uint64 | ||||
| 	indexID      uint64 | ||||
| @@ -70,117 +71,86 @@ 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) { | ||||
| func (i *Index[T]) Ascend(tx *Snapshot, each func(*T) bool) { | ||||
| 	tx = i.ensureSnapshot(tx) | ||||
| 	i.btreeForIter(tx).Ascend(func(t *T) bool { | ||||
| 		return each(i.copy(t)) | ||||
| 	}) | ||||
| } | ||||
|  | ||||
| func (i Index[T]) AscendAfter(tx *Snapshot, after *T, each func(*T) bool) { | ||||
| func (i *Index[T]) AscendAfter(tx *Snapshot, after *T, each func(*T) bool) { | ||||
| 	tx = i.ensureSnapshot(tx) | ||||
| 	i.btreeForIter(tx).AscendGreaterOrEqual(after, func(t *T) bool { | ||||
| 		return each(i.copy(t)) | ||||
| 	}) | ||||
| } | ||||
|  | ||||
| func (i Index[T]) Descend(tx *Snapshot, each func(*T) bool) { | ||||
| func (i *Index[T]) Descend(tx *Snapshot, each func(*T) bool) { | ||||
| 	tx = i.ensureSnapshot(tx) | ||||
| 	i.btreeForIter(tx).Descend(func(t *T) bool { | ||||
| 		return each(i.copy(t)) | ||||
| 	}) | ||||
| } | ||||
|  | ||||
| func (i Index[T]) DescendAfter(tx *Snapshot, after *T, each func(*T) bool) { | ||||
| func (i *Index[T]) DescendAfter(tx *Snapshot, after *T, each func(*T) bool) { | ||||
| 	tx = i.ensureSnapshot(tx) | ||||
| 	i.btreeForIter(tx).DescendLessOrEqual(after, func(t *T) bool { | ||||
| 		return each(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 +158,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 +166,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 +174,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 +201,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))) | ||||
| } | ||||
|   | ||||
| @@ -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 | ||||
|   | ||||
		Reference in New Issue
	
	Block a user