Merge branch 'develop' into remove-cluster-refs

This commit is contained in:
Cody Soyland 2018-06-27 14:22:49 -05:00 committed by GitHub
commit 27de0e1682
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
2 changed files with 59 additions and 43 deletions

56
api.go
View file

@ -35,11 +35,10 @@ import (
// API provides the top level programmatic interface to Pilosa. It is usually
// wrapped by a handler which provides an external interface (e.g. HTTP).
type API struct {
Holder *Holder
Broadcaster Broadcaster
Cluster *Cluster
TranslateStore TranslateStore
server *Server
Holder *Holder
Broadcaster Broadcaster
Cluster *Cluster
server *Server
}
// APIOption is a functional option type for pilosa.API
@ -48,7 +47,6 @@ type APIOption func(*API) error
func OptAPIServer(s *Server) APIOption {
return func(a *API) error {
a.server = s
a.TranslateStore = s.primaryTranslateStore
a.Holder = s.holder
a.Broadcaster = s
a.Cluster = s.cluster
@ -142,9 +140,9 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er
}
// Translate column attributes, if necessary.
if api.TranslateStore != nil {
if api.server.primaryTranslateStore != nil {
for _, col := range resp.ColumnAttrSets {
v, err := api.TranslateStore.TranslateColumnToString(req.Index, col.ID)
v, err := api.server.primaryTranslateStore.TranslateColumnToString(req.Index, col.ID)
if err != nil {
return resp, err
}
@ -787,6 +785,48 @@ func (api *API) ResizeAbort() error {
return errors.Wrap(err, "complete current job")
}
// TranslateStoreBufferSize is the buffer size used for streaming data.
const TranslateStoreBufferSize = 65536
func (api *API) GetTranslateData(ctx context.Context, w io.WriteCloser, offset int64) error {
rc, err := api.server.primaryTranslateStore.Reader(ctx, offset)
if err != nil {
return errors.Wrap(err, "read from translate store")
}
// Ensure reader is closed when the client disconnects.
go func() { <-ctx.Done(); rc.Close() }()
go func() {
defer rc.Close()
defer w.Close()
buf := make([]byte, TranslateStoreBufferSize)
// Copy from reader to client until store or client disconnect.
for {
// Read from store.
n, err := rc.Read(buf)
if err == io.EOF {
return
} else if err != nil {
api.server.logger.Printf("api: translate store read error: %s", err)
return
} else if n == 0 {
continue
}
// Write to response & flush.
if _, err := w.Write(buf[:n]); err != nil {
api.server.logger.Printf("api: translate store response write error: %s", err)
return
}
}
}()
return nil
}
// State returns the cluster state which is usually "NORMAL", but could be
// "STARTING", "RESIZING", or potentially others. See cluster.go for more
// details.

View file

@ -1295,25 +1295,22 @@ func (h *Handler) GetAPI() *pilosa.API {
type defaultClusterMessageResponse struct{}
// TranslateStoreBufferSize is the buffer size used for streaming data.
const TranslateStoreBufferSize = 65536
func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) {
q := r.URL.Query()
offset, _ := strconv.ParseInt(q.Get("offset"), 10, 64)
rc, err := h.API.TranslateStore.Reader(r.Context(), offset)
if err == pilosa.ErrNotImplemented {
http.Error(w, err.Error(), http.StatusNotImplemented)
return
} else if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
pipeR, pipeW := io.Pipe()
err := h.API.GetTranslateData(r.Context(), pipeW, offset)
if err != nil {
if errors.Cause(err) == pilosa.ErrNotImplemented {
http.Error(w, err.Error(), http.StatusNotImplemented)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
defer rc.Close()
// Ensure reader is closed when the client disconnects.
go func() { <-r.Context().Done(); rc.Close() }()
// Flush header so client can continue.
w.WriteHeader(http.StatusOK)
@ -1321,28 +1318,7 @@ func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request)
w.Flush()
}
// Copy from reader to client until store or client disconnect.
buf := make([]byte, TranslateStoreBufferSize)
for {
// Read from store.
n, err := rc.Read(buf)
if err == io.EOF {
return
} else if err != nil {
h.Logger.Printf("http: translate store read error: %s", err)
return
} else if n == 0 {
continue
}
// Write to response & flush.
if _, err := w.Write(buf[:n]); err != nil {
h.Logger.Printf("http: translate store response write error: %s", err)
return
} else if w, ok := w.(http.Flusher); ok {
w.Flush()
}
}
io.Copy(w, pipeR)
}
type queryValidationSpec struct {