add ID auto-generation on the coordinator

This commit is contained in:
Nia Weiss 2020-12-03 11:11:31 -05:00
parent df13c72351
commit 7ea2e4b7e3
No known key found for this signature in database
GPG key ID: 895E83409BFDA1BB
7 changed files with 902 additions and 0 deletions

42
api.go
View file

@ -2082,6 +2082,42 @@ func (api *API) PastQueries(ctx context.Context, remote bool) ([]PastQueryStatus
return clusterQueries, nil
}
func (api *API) ReserveIDs(key IDAllocKey, session [32]byte, offset uint64, count uint64) ([]IDRange, error) {
if err := api.validate(apiIDReserve); err != nil {
return nil, errors.Wrap(err, "validating api method")
}
if api.holder.isCoordinator() {
return api.holder.ida.reserve(key, session, offset, count)
}
return nil, errors.New("cannot reserve IDs on a non-coordinator node")
}
func (api *API) CommitIDs(key IDAllocKey, session [32]byte, count uint64) error {
if err := api.validate(apiIDCommit); err != nil {
return errors.Wrap(err, "validating api method")
}
if api.holder.isCoordinator() {
return api.holder.ida.commit(key, session, count)
}
return errors.New("cannot commit IDs on a non-coordinator node")
}
func (api *API) ResetIDAlloc(index string) error {
if err := api.validate(apiIDReset); err != nil {
return errors.Wrap(err, "validating api method")
}
if api.holder.isCoordinator() {
return api.holder.ida.reset(index)
}
return errors.New("cannot reset IDs on a non-coordinator node")
}
// TranslateIndexDB is an internal function to load the index keys database
// rd is a boltdb file.
func (api *API) TranslateIndexDB(ctx context.Context, indexName string, partitionID int, rd io.Reader) error {
@ -2157,6 +2193,9 @@ const (
apiGetTransaction
apiActiveQueries
apiPastQueries
apiIDReserve
apiIDCommit
apiIDReset
)
var methodsCommon = map[apiMethod]struct{}{
@ -2198,4 +2237,7 @@ var methodsNormal = map[apiMethod]struct{}{
apiGetTransaction: {},
apiActiveQueries: {},
apiPastQueries: {},
apiIDReserve: {},
apiIDCommit: {},
apiIDReset: {},
}

View file

@ -98,11 +98,16 @@ type Holder struct {
// Func to open whatever implementation of transaction store we're using.
OpenTransactionStore OpenTransactionStoreFunc
// Func to open the ID allocator.
OpenIDAllocator func(string) (*idAllocator, error)
// transactionManager
transactionManager *TransactionManager
translationSyncer TranslationSyncer
ida *idAllocator
// Queue of fields (having a foreign index) which have
// opened before their foreign index has opened.
foreignIndexFields []*Field
@ -195,6 +200,7 @@ type HolderConfig struct {
OpenTranslateStore OpenTranslateStoreFunc
OpenTranslateReader OpenTranslateReaderFunc
OpenTransactionStore OpenTransactionStoreFunc
OpenIDAllocator OpenIDAllocatorFunc
TranslationSyncer TranslationSyncer
CacheFlushInterval time.Duration
StatsClient stats.StatsClient
@ -213,6 +219,7 @@ func DefaultHolderConfig() *HolderConfig {
OpenTranslateStore: OpenInMemTranslateStore,
OpenTranslateReader: nil,
OpenTransactionStore: OpenInMemTransactionStore,
OpenIDAllocator: func(string) (*idAllocator, error) { return &idAllocator{}, nil },
TranslationSyncer: NopTranslationSyncer,
CacheFlushInterval: defaultCacheFlushInterval,
StatsClient: stats.NopStatsClient,
@ -253,6 +260,7 @@ func NewHolder(path string, cfg *HolderConfig) *Holder {
OpenTranslateStore: cfg.OpenTranslateStore,
OpenTranslateReader: cfg.OpenTranslateReader,
OpenTransactionStore: cfg.OpenTransactionStore,
OpenIDAllocator: cfg.OpenIDAllocator,
translationSyncer: cfg.TranslationSyncer,
Logger: cfg.Logger,
Opts: HolderOpts{Txsrc: cfg.Txsrc, RowcacheOff: cfg.RowcacheOff},
@ -596,6 +604,12 @@ func (h *Holder) Open() error {
h.transactionManager = NewTransactionManager(tstore)
h.transactionManager.Log = h.Logger
// Open ID allocator.
h.ida, err = h.OpenIDAllocator(filepath.Join(h.path, "idalloc.db"))
if err != nil {
return errors.Wrap(err, "opening ID allocator")
}
// Open path to read all index directories.
f, err := os.Open(h.path)
if err != nil {
@ -742,6 +756,9 @@ func (h *Holder) Close() error {
if err := h.txf.Close(); err != nil {
return errors.Wrap(err, "holder.Txf.Close()")
}
if err := h.ida.Close(); err != nil {
return errors.Wrap(err, "closing ID allocator")
}
// Reset opened in case Holder needs to be reopened.
h.txf = nil

View file

@ -430,6 +430,10 @@ func newRouter(handler *Handler) http.Handler {
router.HandleFunc("/internal/translate/field/{index}/{field}/keys/find", handler.handleFindFieldKeys).Methods("POST").Name("FindFieldKeys")
router.HandleFunc("/internal/translate/field/{index}/{field}/keys/create", handler.handleCreateFieldKeys).Methods("POST").Name("CreateFieldKeys")
router.HandleFunc("/internal/idalloc/reserve", handler.handleReserveIDs).Methods("POST").Name("ReserveIDs")
router.HandleFunc("/internal/idalloc/commit", handler.handleCommitIDs).Methods("POST").Name("CommitIDs")
router.HandleFunc("/internal/idalloc/reset/{index}", handler.handleResetIDAlloc).Methods("POST").Name("ResetIDAlloc")
// endpoints for collecting cpu profiles from a chosen begin point to
// when the client wants to stop. Used for profiling imports that
// could be long or short.
@ -2845,3 +2849,91 @@ func (h *Handler) handleCreateFieldKeys(w http.ResponseWriter, r *http.Request)
return
}
}
func (h *Handler) handleReserveIDs(w http.ResponseWriter, r *http.Request) {
// Verify input and output types
if r.Header.Get("Content-Type") != "application/json" {
http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType)
return
} else if !validHeaderAcceptJSON(r.Header) {
http.Error(w, "Not acceptable", http.StatusNotAcceptable)
return
}
bd, err := readBody(r)
if err != nil {
http.Error(w, "failed to read body", http.StatusBadRequest)
return
}
var req pilosa.IDAllocReserveRequest
req.Offset = ^uint64(0)
err = json.Unmarshal(bd, &req)
if err != nil {
http.Error(w, "failed to decode request", http.StatusBadRequest)
return
}
ids, err := h.api.ReserveIDs(req.Key, req.Session, req.Offset, req.Count)
if err != nil {
http.Error(w, fmt.Sprintf("reserving IDs: %v", err.Error()), http.StatusBadRequest)
return
}
err = json.NewEncoder(w).Encode(ids)
if err != nil {
http.Error(w, "encoding result", http.StatusBadRequest)
return
}
}
func (h *Handler) handleCommitIDs(w http.ResponseWriter, r *http.Request) {
// Verify input and output types
if r.Header.Get("Content-Type") != "application/json" {
http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType)
return
} else if !validHeaderAcceptJSON(r.Header) {
http.Error(w, "Not acceptable", http.StatusNotAcceptable)
return
}
bd, err := readBody(r)
if err != nil {
http.Error(w, "failed to read body", http.StatusBadRequest)
return
}
var req pilosa.IDAllocCommitRequest
err = json.Unmarshal(bd, &req)
if err != nil {
http.Error(w, "failed to decode request", http.StatusBadRequest)
return
}
err = h.api.CommitIDs(req.Key, req.Session, req.Count)
if err != nil {
http.Error(w, fmt.Sprintf("committing IDs: %v", err.Error()), http.StatusBadRequest)
return
}
w.WriteHeader(http.StatusNoContent)
}
func (h *Handler) handleResetIDAlloc(w http.ResponseWriter, r *http.Request) {
if !validHeaderAcceptType(r.Header, "text", "plain") {
http.Error(w, "text/plain is not an acceptable response type", http.StatusNotAcceptable)
}
indexName, ok := mux.Vars(r)["index"]
if !ok {
http.Error(w, "index name is required", http.StatusBadRequest)
return
}
err := h.api.ResetIDAlloc(indexName)
if err != nil {
http.Error(w, fmt.Sprintf("resetting ID allocation: %v", err.Error()), http.StatusBadRequest)
return
}
w.Header().Add("Content-Type", "text/plain")
w.WriteHeader(http.StatusOK)
w.Write([]byte("OK")) //nolint:errcheck
}

521
idalloc.go Normal file
View file

@ -0,0 +1,521 @@
// Copyright 2020 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa
import (
"encoding/binary"
"math/bits"
"sort"
"time"
"github.com/pkg/errors"
bolt "go.etcd.io/bbolt"
)
// IDAllocKey is an ID allocation key.
type IDAllocKey struct {
Index string `json:"index"`
Key string `json:"key,omitempty"`
}
type IDAllocReserveRequest struct {
Key IDAllocKey `json:"key"`
Session [32]byte `json:"session"`
Offset uint64 `json:"offset"`
Count uint64 `json:"count"`
}
type IDAllocCommitRequest struct {
Key IDAllocKey `json:"key"`
Session [32]byte `json:"session"`
Count uint64 `json:"count"`
}
func (k IDAllocKey) String() string {
// This is a pretty printing format.
// **DO NOT** attempt to parse this.
return k.Index + ":" + k.Key
}
type idAllocator struct {
db *bolt.DB
}
type OpenIDAllocatorFunc func(path string) (*idAllocator, error) // whyyyyyyyyy
func OpenIDAllocator(path string) (*idAllocator, error) {
db, err := bolt.Open(path, 0666, &bolt.Options{Timeout: 1 * time.Second})
if err != nil {
return nil, err
}
return &idAllocator{db}, nil
}
func (ida *idAllocator) Close() error {
if ida == nil || ida.db == nil {
return nil
}
defer func() { ida.db = nil }()
return ida.db.Close()
}
func (ida *idAllocator) reserve(key IDAllocKey, session [32]byte, offset, count uint64) ([]IDRange, error) {
if session == [32]byte{} {
return nil, errors.New("detected a broken session key")
}
var ranges []IDRange
err := ida.db.Update(func(tx *bolt.Tx) error {
// Find the bucket associated with the key.
bkt, err := key.findBucket(tx)
if err != nil {
return err
}
// Fetch the old reservation.
var res idReservation
err = res.decode(bkt.Get([]byte(key.Key + "\x00")))
if err != nil {
return errors.Wrap(err, "decoding old reservation")
}
res.lock = session
if offset != ^uint64(0) && offset != res.offset {
// Offset control is in use and the offset differs.
if offset < res.offset {
// The client is confused.
// Abort before we break something.
return errors.New("unable to rewind ID generation")
}
// Roll the IDs forward.
diff := offset - res.offset
for diff > 0 && len(res.ranges) > 0 {
if res.ranges[0].width() > diff-1 {
res.ranges[0].First += diff
break
}
diff -= res.ranges[0].width() + 1
res.ranges = res.ranges[1:]
}
}
res.offset = offset
// Count already-reserved ranges.
var preReserved uint64
for _, r := range res.ranges {
sum, overflow := bits.Add64(preReserved, r.width(), 1)
if overflow != 0 {
return errors.New("impossibly large reservation")
}
preReserved = sum
}
switch {
case preReserved == count:
// The client probbably failed while ingesting this.
// Spit out the same ID ranges a second time.
ranges = res.ranges
case preReserved > count:
// This could happen????
ranges = make([]IDRange, 0, len(res.ranges))
for _, r := range res.ranges {
if count == 0 {
break
}
if count-1 < r.width() {
ranges = append(ranges, IDRange{
First: r.First,
Last: r.First + count - 1,
})
break
}
ranges = append(ranges, r)
count -= r.width() + 1
}
default:
// Reserve additional IDs.
reserved, err := doReserveIDs(bkt, count-preReserved)
if err != nil {
return err
}
reserved = append(res.ranges, reserved...)
res.ranges = reserved
ranges = reserved
}
// Save the updated reservation.
encoded, err := res.encode()
if err != nil {
return errors.Wrap(err, "encoding updated ID reservation entry")
}
err = bkt.Put([]byte(key.Key+"\x00"), encoded)
if err != nil {
return errors.Wrap(err, "saving updated ID reservation entry")
}
return nil
})
if err != nil {
return nil, errors.Wrapf(err, "acquiring IDs from key %s in session %v", key, session)
}
return ranges, nil
}
func (ida *idAllocator) commit(key IDAllocKey, session [32]byte, count uint64) error {
if session == [32]byte{} {
return errors.New("detected a broken session key")
}
err := ida.db.Update(func(tx *bolt.Tx) error {
// Find the bucket associated with the key.
bkt, err := key.findBucket(tx)
if err != nil {
return err
}
// Fetch the old reservation.
prev := bkt.Get([]byte(key.Key + "\x00"))
if prev == nil {
// There is nothing to commit.
return errors.New("nothing to commit")
}
var res idReservation
err = res.decode(prev)
if err != nil {
return errors.Wrap(err, "decoding old reservation")
}
if res.lock != session {
// This probbably means that there are multiple clients attempting to manipulate the same ID allocator key.
return errors.New("lock collision")
}
if res.offset != ^uint64(0) {
// Update the offset.
if count > ^uint64(1)-res.offset {
return errors.New("offset overflow")
}
res.offset += count
}
// Discard processed ID ranges.
n := count
for n > 0 && len(res.ranges) > 0 {
r := &res.ranges[0]
w := r.width()
if n-1 < w {
r.First += n
n = 0
break
}
n -= w + 1
res.ranges = res.ranges[1:]
}
if n > 0 {
// The client attempted to commit more IDs than they requested.
return errors.New("overcommitted processed IDs")
}
if res.offset == ^uint64(0) && len(res.ranges) > 0 {
// Return unused ranges to the central list.
err = doReleaseIDs(bkt, res.ranges...)
if err != nil {
return errors.Wrap(err, "releasing unused IDs")
}
res.ranges = nil
}
if res.offset == ^uint64(0) && len(res.ranges) == 0 {
// Delete the reservation entry, as it no longer contains any useful information.
err = bkt.Delete([]byte(key.Key + "\x00"))
if err != nil {
return errors.Wrap(err, "deleting used ID reservation metadata")
}
return nil
}
// Save the updated reservation.
encoded, err := res.encode()
if err != nil {
return errors.Wrap(err, "encoding updated ID reservation entry")
}
err = bkt.Put([]byte(key.Key+"\x00"), encoded)
if err != nil {
return errors.Wrap(err, "saving updated ID reservation entry")
}
return nil
})
if err != nil {
return errors.Wrapf(err, "committing %d IDs in allocation session %x with key %s", count, session, key)
}
return nil
}
func (ida *idAllocator) reset(index string) error {
err := ida.db.Update(func(tx *bolt.Tx) error {
if tx.Bucket([]byte(index)) == nil {
return nil
}
err := tx.DeleteBucket([]byte(index))
if err != nil {
return errors.Wrap(err, "deleting bucket")
}
return nil
})
if err != nil {
return errors.Wrapf(err, "resetting ID allocation for index %q", index)
}
return nil
}
func (k IDAllocKey) findBucket(tx *bolt.Tx) (*bolt.Bucket, error) {
ibkt, err := tx.CreateBucketIfNotExists([]byte(k.Index))
if err != nil {
return nil, errors.Wrap(err, "creating index bucket")
}
return ibkt, nil
}
// IDRange is a reserved ID range.
type IDRange struct {
First uint64 `json:"first"`
Last uint64 `json:"last"`
}
func (r *IDRange) decode(data []byte) error {
if len(data) != 16 {
return errors.New("wrong size for ID range")
}
first := binary.LittleEndian.Uint64(data[:8])
last := binary.LittleEndian.Uint64(data[8:16])
if last < first {
return errors.New("backwards ID range")
}
*r = IDRange{First: first, Last: last}
return nil
}
// width of the range (count-1).
func (r IDRange) width() uint64 {
return r.Last - r.First
}
func decodeRanges(data []byte) ([]IDRange, error) {
if len(data)%16 != 0 {
return nil, errors.New("incorrectly sized ranges list")
}
ranges := make([]IDRange, len(data)/16)
for i := range ranges {
err := ranges[i].decode(data[16*i:][:16])
if err != nil {
return nil, errors.Wrapf(err, "decoding range %d", i)
}
}
return ranges, nil
}
func encodeRanges(ranges []IDRange) ([]byte, error) {
data := make([]byte, 16*len(ranges))
for i, r := range ranges {
if r.Last < r.First {
return nil, errors.Errorf("backwards ID range at index %d", i)
}
rdat := data[16*i:][:16]
binary.LittleEndian.PutUint64(rdat[0:8], r.First)
binary.LittleEndian.PutUint64(rdat[8:16], r.Last)
}
return data, nil
}
type idReservation struct {
ranges []IDRange
offset uint64
lock [32]byte
}
func (r *idReservation) decode(data []byte) error {
if data == nil {
*r = idReservation{}
return nil
}
if len(data) < 8+len(r.lock) {
return errors.New("invalid ID reservation entry")
}
offset := binary.LittleEndian.Uint64(data[0:8])
lock := data[8 : 8+len(r.lock)]
ranges, err := decodeRanges(data[8+len(r.lock):])
if err != nil {
return errors.Wrap(err, "decoding ranges in reservation entry")
}
r.ranges = ranges
r.offset = offset
copy(r.lock[:], lock)
return nil
}
func (r idReservation) encode() ([]byte, error) {
var odat [8]byte
binary.LittleEndian.PutUint64(odat[:], r.offset)
rdat, err := encodeRanges(r.ranges)
if err != nil {
return nil, err
}
return append(append(odat[:], r.lock[:]...), rdat...), nil
}
func doReserveIDs(bkt *bolt.Bucket, count uint64) ([]IDRange, error) {
var avail []IDRange
if prev := bkt.Get([]byte{1}); prev != nil {
r, err := decodeRanges(prev)
if err != nil {
return nil, err
}
avail = r
} else {
avail = []IDRange{{1, ^uint64(0)}}
}
reserved, avail, err := reserveIDs(avail, count)
if err != nil {
return nil, errors.Wrap(err, "reserving IDs")
}
encoded, err := encodeRanges(avail)
if err != nil {
return nil, errors.Wrap(err, "encoding available ID ranges")
}
err = bkt.Put([]byte{1}, encoded)
if err != nil {
return nil, errors.Wrap(err, "saving available ID ranges")
}
return reserved, nil
}
func reserveIDs(avail []IDRange, count uint64) (reserved []IDRange, unused []IDRange, err error) {
// Min-heapify the available ranges by size.
for i, v := range avail {
width := v.width()
for width < avail[(i-1)/2].width() {
avail[(i-1)/2], avail[i] = v, avail[(i-1)/2]
i = (i - 1) / 2
}
}
// Reserve ranges in ascending size order.
for count > 0 && len(avail) > 0 {
if count-1 < avail[0].width() {
// The smallest available range is larger than what is needed.
reserved = append(reserved, IDRange{
First: avail[0].First,
Last: avail[0].First + count - 1,
})
avail[0].First += count
count = 0
break
}
// Completely reserve the smallest range.
count -= avail[0].width() + 1
reserved = append(reserved, avail[0])
avail[0] = avail[len(avail)-1]
avail = avail[:len(avail)-1]
i := 0
for {
min := i
if l := 2*i + 1; l < len(avail) && avail[l].width() < avail[min].width() {
min = l
}
if r := 2*i + 2; r < len(avail) && avail[r].width() < avail[min].width() {
min = r
}
if min == i {
break
}
avail[i], avail[min] = avail[min], avail[i]
i = min
}
}
if count > 0 {
return nil, nil, errors.Errorf("insufficient available ID space (missing %d IDs)", count)
}
return reserved, avail, nil
}
func doReleaseIDs(bkt *bolt.Bucket, ranges ...IDRange) error {
var avail []IDRange
if prev := bkt.Get([]byte{1}); prev != nil {
r, err := decodeRanges(prev)
if err != nil {
return err
}
avail = r
} else {
return errors.New("cannot return IDs to the void")
}
avail, err := mergeIDs(avail, ranges)
if err != nil {
return errors.Wrap(err, "returning IDs")
}
encoded, err := encodeRanges(avail)
if err != nil {
return errors.Wrap(err, "encoding available ID ranges")
}
err = bkt.Put([]byte{1}, encoded)
if err != nil {
return errors.Wrap(err, "saving available ID ranges")
}
return nil
}
func mergeIDs(main, returned []IDRange) ([]IDRange, error) {
if len(main) == 0 && len(returned) == 0 {
return nil, nil
}
main = append(main, returned...)
sort.Slice(main, func(i, j int) bool { return main[i].First < main[j].First })
i := 1
for _, v := range main[1:] {
prev := &main[i-1]
switch {
case v.First == prev.First || v.First <= prev.Last:
return nil, errors.Errorf("found unexpected overlap between ranges %v and %v", prev, v)
case v.First == prev.Last+1:
prev.Last = v.Last
default:
main[i] = v
i++
}
}
return main[:i], nil
}

219
idalloc_test.go Normal file
View file

@ -0,0 +1,219 @@
// Copyright 2020 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa
import (
"crypto/rand"
"io"
"io/ioutil"
"os"
"reflect"
"testing"
"time"
bolt "go.etcd.io/bbolt"
)
func TestIDAlloc(t *testing.T) {
// Acquire a temporary file.
f, err := ioutil.TempFile("", "")
if err != nil {
t.Errorf("acquiring temporary file: %v", err)
return
}
defer func() {
cerr := f.Close()
if cerr != nil {
t.Errorf("closing temporary file: %v", cerr)
}
}()
// Open bolt.
db, err := bolt.Open(f.Name(), 0666, &bolt.Options{Timeout: 1 * time.Second})
if rerr := os.Remove(f.Name()); rerr != nil {
t.Errorf("removing temporary file: %v", rerr)
return
}
if err != nil {
t.Errorf("opening bolt: %v", err)
return
}
defer func() {
cerr := db.Close()
if cerr != nil {
t.Errorf("closing bolt: %v", cerr)
}
}()
alloc := idAllocator{db: db}
a := IDAllocKey{
Index: "h",
Key: "a",
}
b := IDAllocKey{
Index: "h",
Key: "b",
}
var sessions [11][32]byte
for i := range sessions {
_, err = io.ReadFull(rand.Reader, sessions[i][:])
if err != nil {
t.Errorf("failed to obtain entropy: %v", err)
return
}
}
// Reserve 2 batches of IDs and then commit both.
if ranges, err := alloc.reserve(a, sessions[0], ^uint64(0), 3); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
expect := []IDRange{{1, 3}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if ranges, err := alloc.reserve(b, sessions[1], ^uint64(0), 3); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
expect := []IDRange{{4, 6}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if err := alloc.commit(a, sessions[0], 3); err != nil {
t.Errorf("failed to commit reservation of 3 IDs: %v", err)
}
if err := alloc.commit(b, sessions[1], 3); err != nil {
t.Errorf("failed to commit reservation of 3 IDs: %v", err)
}
// Test collision handling.
if ranges, err := alloc.reserve(a, sessions[2], ^uint64(0), 4); err != nil {
t.Errorf("failed to reserve 4 IDs: %v", err)
} else {
expect := []IDRange{{7, 10}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if ranges, err := alloc.reserve(a, sessions[3], ^uint64(0), 4); err != nil {
t.Errorf("failed to reserve 4 IDs: %v", err)
} else {
// Instead of reserving new ranges, we have taken over the previous reservation.
expect := []IDRange{{7, 10}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if err := alloc.commit(a, sessions[2], 4); err == nil {
// This was taken over by reservation 3, so this should have warned us that something went wrong.
t.Error("unexpected commit success for clobbered reservation")
}
if err := alloc.commit(a, sessions[3], 4); err != nil {
t.Errorf("failed to commit reservation of 4 IDs: %v", err)
}
// Test key return.
if ranges, err := alloc.reserve(a, sessions[4], ^uint64(0), 100); err != nil {
t.Errorf("failed to reserve 100 IDs: %v", err)
} else {
// Instead of reserving new ranges, we have taken over the previous reservation.
expect := []IDRange{{11, 110}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if err := alloc.commit(a, sessions[4], 1); err != nil {
t.Errorf("failed to commit reservation of 1 ID: %v", err)
}
if ranges, err := alloc.reserve(a, sessions[5], ^uint64(0), 100); err != nil {
t.Errorf("failed to reserve 100 IDs: %v", err)
} else {
// Instead of reserving new ranges, we have taken over the previous reservation.
expect := []IDRange{{12, 111}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if err := alloc.commit(a, sessions[5], 100); err != nil {
t.Errorf("failed to commit reservation of 100 IDs: %v", err)
}
// Test offset-based allocation.
p1 := IDAllocKey{
Index: "xyzzy",
Key: "kafka-partition-1",
}
p2 := IDAllocKey{
Index: "xyzzy",
Key: "kafka-partition-2",
}
// Start processsing on 2 partitions.
if ranges, err := alloc.reserve(p1, sessions[6], 0, 3); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
expect := []IDRange{{1, 3}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
if ranges, err := alloc.reserve(p2, sessions[7], 0, 3); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
expect := []IDRange{{4, 6}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
// Partition 2 finishes a batch of 2 IDs.
if err := alloc.commit(p2, sessions[7], 2); err != nil {
t.Errorf("failed to commit reservation of 2 IDs: %v", err)
}
// Partition 2 reserves more IDs.
if ranges, err := alloc.reserve(p2, sessions[8], 2, 3); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
// This could be coalesced here, but that is not true in the general case.
expect := []IDRange{{6, 6}, {7, 8}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
// Partition 1 ingester failed and was restarted.
if ranges, err := alloc.reserve(p1, sessions[9], 0, 3); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
expect := []IDRange{{1, 3}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
// This time, the partition 1 ingester finishes the batch.
// It commits the offset to Kafka but fails to commit to the ID allocator.
// It is restarted at Kafka's next offset.
if ranges, err := alloc.reserve(p1, sessions[10], 3, 4); err != nil {
t.Errorf("failed to reserve 3 IDs: %v", err)
} else {
expect := []IDRange{{9, 12}}
if !reflect.DeepEqual(expect, ranges) {
t.Errorf("expeced %v but got %v", expect, ranges)
}
}
// Partition 1 finishes a batch of 4 IDs.
if err := alloc.commit(p1, sessions[10], 4); err != nil {
t.Errorf("failed to commit reservation of 4 IDs: %v", err)
}
}

View file

@ -326,6 +326,16 @@ func OptServerOpenTranslateStore(fn OpenTranslateStoreFunc) ServerOption {
}
}
// OptServerOpenIDAllocator is a functional option on Server
// used to specify the ID allocator data store type.
// Except not really.
func OptServerOpenIDAllocator(fn OpenIDAllocatorFunc) ServerOption {
return func(s *Server) error {
s.holderConfig.OpenIDAllocator = fn
return nil
}
}
// OptServerOpenTranslateReader is a functional option on Server
// used to specify the remote translation data reader.
func OptServerOpenTranslateReader(fn OpenTranslateReaderFunc) ServerOption {

View file

@ -398,6 +398,7 @@ func (m *Command) SetupServer() error {
pilosa.OptServerExecutorPoolSize(m.Config.WorkerPoolSize),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderWithLockerFunc(c, &sync.Mutex{})),
pilosa.OptServerOpenIDAllocator(pilosa.OpenIDAllocator),
pilosa.OptServerLogger(m.logger),
pilosa.OptServerAttrStoreFunc(boltdb.NewAttrStore),
pilosa.OptServerSystemInfo(gopsutil.NewSystemInfo()),