Merge pull request #279 from tgruben/fix-closers

Closed all post request bodies and optimized available shard with new view capabilities
This commit is contained in:
tgruben 2020-04-10 20:29:04 -05:00 • committed by GitHub
commit 550fcec9ee
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
4 changed files with 13 additions and 29 deletions

1
api.go
View file

@ -414,7 +414,6 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string,
if field == nil {
return newNotFoundError(ErrFieldNotFound)
}
errCh := make(chan error, len(nodes))
for _, node := range nodes {

View file

@ -413,13 +413,6 @@ func (f *Field) AvailableShards() *roaring.Bitmap {
return b
}
// constainsShard is used for limiting unnecessary CreateShard broadcast
func (f *Field) containsShard(shard uint64) bool {
f.mu.RLock()
defer f.mu.RUnlock()
return f.remoteAvailableShards.Contains(shard)
}
// AddRemoteAvailableShards merges the set of available shards into the current known set
// and saves the set to a file.
func (f *Field) AddRemoteAvailableShards(b *roaring.Bitmap) error {
@ -1186,7 +1179,6 @@ func (f *Field) newView(path, name string) *view {
view.rowAttrStore = f.rowAttrStore
view.stats = f.Stats
view.broadcaster = f.broadcaster
view.remoteShardPresent = f.containsShard
if f.snapshotQueue != nil {
view.snapshotQueue = f.snapshotQueue
}
@ -1914,7 +1906,6 @@ func (f *Field) importRoaring(ctx context.Context, data []byte, shard uint64, vi
if err != nil {
return errors.Wrap(err, "creating fragment")
}
if err := frag.importRoaring(ctx, data, clear); err != nil {
return err
}
@ -1939,7 +1930,6 @@ func (f *Field) importRoaringOverwrite(ctx context.Context, data []byte, shard u
if err != nil {
return errors.Wrap(err, "creating fragment")
}
if err := frag.importRoaringOverwrite(ctx, data, block); err != nil {
return err
}

View file

@ -1626,7 +1626,6 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques
http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType)
return
}
defer r.Body.Close()
err := h.api.ClusterMessage(r.Context(), r.Body)
if err != nil {
// TODO this was the previous behavior, but perhaps not everything is a bad request
@ -1647,7 +1646,6 @@ func (h *Handler) handlePostTranslateData(w http.ResponseWriter, r *http.Request
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Stream all translation data.
rd, err := h.api.GetTranslateEntryReader(r.Context(), offsets)
if errors.Cause(err) == pilosa.ErrNotImplemented {
@ -1795,7 +1793,6 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request
http.Error(w, "Not acceptable", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
fieldName := mux.Vars(r)["field"]
@ -1833,7 +1830,6 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request
http.Error(w, "shard should be an unsigned integer", http.StatusBadRequest)
return
}
resp := &pilosa.ImportResponse{}
// TODO give meaningful stats for import
err = h.api.ImportRoaring(ctx, indexName, fieldName, shard, remote, req)
@ -1893,7 +1889,6 @@ func (h *Handler) handlePostTranslateIDs(w http.ResponseWriter, r *http.Request)
http.Error(w, "Not acceptable", http.StatusNotAcceptable)
return
}
buf, err := h.api.TranslateIDs(r.Context(), r.Body)
if err != nil {
http.Error(w, fmt.Sprintf("translate ids: %v", err), http.StatusInternalServerError)

26
view.go
View file

@ -56,12 +56,11 @@ type view struct {
// Fragments by shard.
fragments map[uint64]*fragment
broadcaster broadcaster
stats stats.StatsClient
rowAttrStore AttrStore
logger logger.Logger
snapshotQueue snapshotQueue
remoteShardPresent func(uint64) bool
broadcaster broadcaster
stats stats.StatsClient
rowAttrStore AttrStore
logger logger.Logger
snapshotQueue snapshotQueue
knownShards *roaring.Bitmap
knownShardsCopied uint32
@ -81,11 +80,10 @@ func newView(path, index, field, name string, fieldOptions FieldOptions) *view {
fragments: make(map[uint64]*fragment),
broadcaster: NopBroadcaster,
stats: stats.NopStatsClient,
logger: logger.NopLogger,
remoteShardPresent: func(uint64) bool { return false },
knownShards: roaring.NewSliceBitmap(),
broadcaster: NopBroadcaster,
stats: stats.NopStatsClient,
logger: logger.NopLogger,
knownShards: roaring.NewSliceBitmap(),
}
}
@ -322,15 +320,17 @@ func (v *view) CreateFragmentIfNotExists(shard uint64) (*fragment, error) {
frag.RowAttrStore = v.rowAttrStore
v.fragments[shard] = frag
v.addKnownShard(shard)
v.notifyIfNewShard(shard)
v.addKnownShard(shard)
return frag, nil
}
func (v *view) notifyIfNewShard(shard uint64) {
if v.remoteShardPresent(shard) { //checks the fields remoteShards bitmap to see if broadcast needed
if v.knownShards.Contains(shard) { //checks the fields remoteShards bitmap to see if broadcast needed
return
}
broadcastChan := make(chan struct{})
go func() {