diff --git a/api.go b/api.go index 38483a4ea..b455f4301 100644 --- a/api.go +++ b/api.go @@ -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. diff --git a/http/handler.go b/http/handler.go index a079ae479..42236f2f2 100644 --- a/http/handler.go +++ b/http/handler.go @@ -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 {