diff --git a/api.go b/api.go index 9be6427c8..57a7b4032 100644 --- a/api.go +++ b/api.go @@ -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: {}, } diff --git a/holder.go b/holder.go index f93be5ba7..4a22cfa43 100644 --- a/holder.go +++ b/holder.go @@ -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 diff --git a/http/handler.go b/http/handler.go index be621177b..a51466dfe 100644 --- a/http/handler.go +++ b/http/handler.go @@ -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 +} diff --git a/idalloc.go b/idalloc.go new file mode 100644 index 000000000..d13a54b13 --- /dev/null +++ b/idalloc.go @@ -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 +} diff --git a/idalloc_test.go b/idalloc_test.go new file mode 100644 index 000000000..5c6d4103c --- /dev/null +++ b/idalloc_test.go @@ -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) + } +} diff --git a/server.go b/server.go index 1b876bce3..1dda70b74 100644 --- a/server.go +++ b/server.go @@ -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 { diff --git a/server/server.go b/server/server.go index b510e0fb3..9fbf5c9bd 100644 --- a/server/server.go +++ b/server/server.go @@ -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()),