diff --git a/api.go b/api.go index 84c2f5146..f369ba10a 100644 --- a/api.go +++ b/api.go @@ -25,6 +25,7 @@ import ( "reflect" "strconv" "strings" + "time" "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" @@ -221,6 +222,12 @@ func (api *API) DeleteFrame(ctx context.Context, indexName string, frameName str } func (api *API) ExportCSV(ctx context.Context, indexName string, frameName string, viewName string, slice uint64, w io.Writer) error { + // Validate that this handler owns the slice. + if !api.Cluster.OwnsFragment(api.URI.HostPort(), indexName, slice) { + api.logger.Printf("host does not own slice %s-%s slice:%d", api.URI, indexName, slice) + return ErrClusterDoesNotOwnSlice + } + // Find the fragment. f := api.Holder.Fragment(indexName, frameName, viewName, slice) if f == nil { @@ -600,7 +607,166 @@ func (api *API) DeleteView(ctx context.Context, indexName string, frameName stri return err } -// InputJSONDataParser validates input json file and executes SetBit. +func (api *API) IndexAttrDiff(ctx context.Context, indexName string, blocks []AttrBlock) (map[uint64]map[string]interface{}, error) { + // Retrieve index from holder. + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + // Retrieve local blocks. + localBlocks, err := index.ColumnAttrStore().Blocks() + if err != nil { + return nil, err + } + + // Read all attributes from all mismatched blocks. + attrs := make(map[uint64]map[string]interface{}) + for _, blockID := range AttrBlocks(localBlocks).Diff(blocks) { + // Retrieve block data. + m, err := index.ColumnAttrStore().BlockData(blockID) + if err != nil { + return nil, err + } + + // Copy to index-wide struct. + for k, v := range m { + attrs[k] = v + } + } + return attrs, nil +} + +func (api *API) FrameAttrDiff(ctx context.Context, indexName string, frameName string, blocks []AttrBlock) (map[uint64]map[string]interface{}, error) { + // Retrieve index from holder. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return nil, ErrFrameNotFound + } + + // Retrieve local blocks. + localBlocks, err := f.RowAttrStore().Blocks() + if err != nil { + return nil, err + } + + // Read all attributes from all mismatched blocks. + attrs := make(map[uint64]map[string]interface{}) + for _, blockID := range AttrBlocks(localBlocks).Diff(blocks) { + // Retrieve block data. + m, err := f.RowAttrStore().BlockData(blockID) + if err != nil { + return nil, err + } + + // Copy to index-wide struct. + for k, v := range m { + attrs[k] = v + } + } + return attrs, nil +} + +func (api *API) Import(ctx context.Context, req internal.ImportRequest) error { + _, frame, err := api.indexFrame(req.Index, req.Frame, req.Slice) + if err != nil { + return err + } + + // Convert timestamps to time.Time. + timestamps := make([]*time.Time, len(req.Timestamps)) + for i, ts := range req.Timestamps { + if ts == 0 { + continue + } + t := time.Unix(0, ts) + timestamps[i] = &t + } + + // Import into fragment. + err = frame.Import(req.RowIDs, req.ColumnIDs, timestamps) + if err != nil { + api.logger.Printf("import error: index=%s, frame=%s, slice=%d, bits=%d, err=%s", req.Index, req.Frame, req.Slice, len(req.ColumnIDs), err) + } + return err +} + +func (api *API) ImportValue(ctx context.Context, req internal.ImportValueRequest) error { + _, frame, err := api.indexFrame(req.Index, req.Frame, req.Slice) + if err != nil { + return err + } + + // Import into fragment. + err = frame.ImportValue(req.Field, req.ColumnIDs, req.Values) + if err != nil { + api.logger.Printf("import error: index=%s, frame=%s, slice=%d, field=%s, bits=%d, err=%s", req.Index, req.Frame, req.Slice, req.Field, len(req.ColumnIDs), err) + } + return err +} + +func (api *API) ModifyIndexTimeQuantum(ctx context.Context, indexName string, timeQuantum TimeQuantum) error { + // Retrieve index by name. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + // Set default time quantum on index. + return index.SetTimeQuantum(timeQuantum) +} + +func (api *API) ModifyFrameTimeQuantum(ctx context.Context, indexName string, frameName string, timeQuantum TimeQuantum) error { + // Retrieve index by name. + frame := api.Holder.Frame(indexName, frameName) + if frame == nil { + return ErrFrameNotFound + } + + // Set default time quantum on index. + return frame.SetTimeQuantum(timeQuantum) +} + +func (api *API) SliceMax(ctx context.Context, inverse bool) map[string]uint64 { + if inverse { + return api.Holder.MaxInverseSlices() + } + return api.Holder.MaxSlices() +} + +func (api *API) StatsWithTags(tags []string) StatsClient { + return api.Holder.Stats.WithTags(tags...) +} + +func (api *API) ClusterLongQueryTime() time.Duration { + return api.Cluster.LongQueryTime +} + +func (api *API) indexFrame(indexName string, frameName string, slice uint64) (*Index, *Frame, error) { + // Validate that this handler owns the slice. + if !api.Cluster.OwnsFragment(api.URI.HostPort(), indexName, slice) { + api.logger.Printf("host does not own slice %s-%s slice:%d", api.URI, indexName, slice) + return nil, nil, ErrClusterDoesNotOwnSlice + } + + // Find the Index. + api.logger.Println("importing:", indexName, frameName, slice) + index := api.Holder.Index(indexName) + if index == nil { + api.logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", indexName, frameName, slice, ErrIndexNotFound.Error()) + return nil, nil, ErrIndexNotFound + } + + // Retrieve frame. + frame := index.Frame(frameName) + if frame == nil { + api.logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", indexName, frameName, slice, ErrFrameNotFound.Error()) + return nil, nil, ErrFrameNotFound + } + return index, frame, nil +} + +// inputJSONDataParser validates input json file and executes SetBit. func (api *API) inputJSONDataParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) { inputDef, err := index.InputDefinition(name) if err != nil { diff --git a/handler.go b/handler.go index f9fc73c52..0f04028e0 100644 --- a/handler.go +++ b/handler.go @@ -240,27 +240,26 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { dif := time.Since(t) // Calculate per request StatsD metrics when the handler is fully configured. - if h.Holder != nil && h.Cluster != nil { - statsTags := make([]string, 0, 3) + statsTags := make([]string, 0, 3) - if h.Cluster.LongQueryTime > 0 && dif > h.Cluster.LongQueryTime { - h.Logger.Printf("%s %s %v", r.Method, r.URL.String(), dif) - statsTags = append(statsTags, "slow_query") - } - - pathParts := strings.Split(r.URL.Path, "/") - endpointName := strings.Join(pathParts, "_") - - if externalPrefixFlag[pathParts[1]] { - statsTags = append(statsTags, "external") - } - - // useragent tag identifies internal/external endpoints - statsTags = append(statsTags, "useragent:"+r.UserAgent()) - - stats := h.Holder.Stats.WithTags(statsTags...) - stats.Histogram("http."+endpointName, float64(dif), 0.1) + longQueryTime := h.API.ClusterLongQueryTime() + if longQueryTime > 0 && dif > longQueryTime { + h.Logger.Printf("%s %s %v", r.Method, r.URL.String(), dif) + statsTags = append(statsTags, "slow_query") } + + pathParts := strings.Split(r.URL.Path, "/") + endpointName := strings.Join(pathParts, "_") + + if externalPrefixFlag[pathParts[1]] { + statsTags = append(statsTags, "external") + } + + // useragent tag identifies internal/external endpoints + statsTags = append(statsTags, "useragent:"+r.UserAgent()) + + stats := h.API.StatsWithTags(statsTags) + stats.Histogram("http."+endpointName, float64(dif), 0.1) } func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) { @@ -348,13 +347,23 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } } -// handleGetSlicesMax handles GET /schema requests. -func (h *Handler) handleGetSlicesMax(w http.ResponseWriter, r *http.Request) { - if err := json.NewEncoder(w).Encode(getSlicesMaxResponse{ - Standard: h.Holder.MaxSlices(), - Inverse: h.Holder.MaxInverseSlices(), - }); err != nil { - h.Logger.Printf("write slices-max response error: %s", err) +func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) { + inverse, err := strconv.ParseBool(r.URL.Query().Get("inverse")) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + ms := h.API.SliceMax(r.Context(), inverse) + if strings.Contains(r.Header.Get("Accept"), "application/x-protobuf") { + pb := &internal.MaxSlicesResponse{ + MaxSlices: ms, + } + if buf, err := proto.Marshal(pb); err != nil { + h.Logger.Printf("protobuf marshal error: %s", err) + } else if _, err := w.Write(buf); err != nil { + h.Logger.Printf("stream write error: %s", err) + } + return } } @@ -518,16 +527,12 @@ func (h *Handler) handlePatchIndexTimeQuantum(w http.ResponseWriter, r *http.Req return } - // Retrieve index by name. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Set default time quantum on index. - if err := index.SetTimeQuantum(tq); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + if err = h.API.ModifyIndexTimeQuantum(r.Context(), indexName, tq); err != nil { + if err == ErrIndexNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -554,34 +559,14 @@ func (h *Handler) handlePostIndexAttrDiff(w http.ResponseWriter, r *http.Request return } - // Retrieve index from holder. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve local blocks. - blks, err := index.ColumnAttrStore().Blocks() + attrs, err := h.API.IndexAttrDiff(r.Context(), indexName, req.Blocks) if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Read all attributes from all mismatched blocks. - attrs := make(map[uint64]map[string]interface{}) - for _, blockID := range AttrBlocks(blks).Diff(req.Blocks) { - // Retrieve block data. - m, err := index.ColumnAttrStore().BlockData(blockID) - if err != nil { + if err == ErrIndexNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Copy to index-wide struct. - for k, v := range m { - attrs[k] = v } + return } // Encode response. @@ -722,16 +707,12 @@ func (h *Handler) handlePatchFrameTimeQuantum(w http.ResponseWriter, r *http.Req return } - // Retrieve index by name. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Set default time quantum on index. - if err := f.SetTimeQuantum(tq); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + if err := h.API.ModifyFrameTimeQuantum(r.Context(), indexName, frameName, tq); err != nil { + if err == ErrFragmentNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -909,34 +890,15 @@ func (h *Handler) handlePostFrameAttrDiff(w http.ResponseWriter, r *http.Request return } - // Retrieve index from holder. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve local blocks. - blks, err := f.RowAttrStore().Blocks() + attrs, err := h.API.FrameAttrDiff(r.Context(), indexName, frameName, req.Blocks) if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Read all attributes from all mismatched blocks. - attrs := make(map[uint64]map[string]interface{}) - for _, blockID := range AttrBlocks(blks).Diff(req.Blocks) { - // Retrieve block data. - m, err := f.RowAttrStore().BlockData(blockID) - if err != nil { + switch err { + case ErrFragmentNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + default: http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Copy to index-wide struct. - for k, v := range m { - attrs[k] = v } + return } // Encode response. @@ -1094,44 +1056,17 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { return } - // Convert timestamps to time.Time. - timestamps := make([]*time.Time, len(req.Timestamps)) - for i, ts := range req.Timestamps { - if ts == 0 { - continue + if err := h.API.Import(r.Context(), req); err != nil { + switch err { + case ErrIndexNotFound: + fallthrough + case ErrFrameNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrClusterDoesNotOwnSlice: + http.Error(w, err.Error(), http.StatusPreconditionFailed) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) } - t := time.Unix(0, ts) - timestamps[i] = &t - } - - // Validate that this handler owns the slice. - if !h.Cluster.OwnsFragment(h.Node.ID, req.Index, req.Slice) { - msg := fmt.Sprintf("host does not own slice %s-%s slice:%d", h.Node.ID, req.Index, req.Slice) - http.Error(w, msg, http.StatusPreconditionFailed) - return - } - - // Find the Index. - h.Logger.Printf("importing: %s %s %d", req.Index, req.Frame, req.Slice) - index := h.Holder.Index(req.Index) - if index == nil { - h.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error()) - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve frame. - f := index.Frame(req.Frame) - if f == nil { - h.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error()) - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Import into fragment. - err = f.Import(req.RowIDs, req.ColumnIDs, timestamps) - if err != nil { - h.Logger.Printf("import error: index=%s, frame=%s, slice=%d, bits=%d, err=%s", req.Index, req.Frame, req.Slice, len(req.ColumnIDs), err) return } @@ -1174,34 +1109,17 @@ func (h *Handler) handlePostImportValue(w http.ResponseWriter, r *http.Request) return } - // Validate that this handler owns the slice. - if !h.Cluster.OwnsFragment(h.Node.ID, req.Index, req.Slice) { - msg := fmt.Sprintf("host does not own slice %s-%s slice:%d", h.Node.ID, req.Index, req.Slice) - http.Error(w, msg, http.StatusPreconditionFailed) - return - } - - // Find the Index. - h.Logger.Printf("importing: %s %s %d", req.Index, req.Frame, req.Slice) - index := h.Holder.Index(req.Index) - if index == nil { - h.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error()) - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve frame. - f := index.Frame(req.Frame) - if f == nil { - h.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error()) - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Import into fragment. - err = f.ImportValue(req.Field, req.ColumnIDs, req.Values) - if err != nil { - h.Logger.Printf("import error: index=%s, frame=%s, slice=%d, field=%s, bits=%d, err=%s", req.Index, req.Frame, req.Slice, req.Field, len(req.ColumnIDs), err) + if err = h.API.ImportValue(r.Context(), req); err != nil { + switch err { + case ErrIndexNotFound: + fallthrough + case ErrFrameNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrClusterDoesNotOwnSlice: + http.Error(w, err.Error(), http.StatusPreconditionFailed) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -1240,19 +1158,16 @@ func (h *Handler) handleGetExportCSV(w http.ResponseWriter, r *http.Request) { return } - // Validate that this handler owns the slice. - if !h.Cluster.OwnsFragment(h.Node.ID, index, slice) { - msg := fmt.Sprintf("host does not own slice %s-%s slice:%d", h.Node.ID, index, slice) - http.Error(w, msg, http.StatusPreconditionFailed) - return - } - if err = h.API.ExportCSV(r.Context(), index, frame, view, slice, w); err != nil { - if err == ErrFragmentNotFound { + switch err { + case ErrFragmentNotFound: http.Error(w, err.Error(), http.StatusNotFound) - } else { + case ErrClusterDoesNotOwnSlice: + http.Error(w, err.Error(), http.StatusPreconditionFailed) + default: http.Error(w, err.Error(), http.StatusInternalServerError) } + return } } diff --git a/pilosa.go b/pilosa.go index 98c67203f..04c76a133 100644 --- a/pilosa.go +++ b/pilosa.go @@ -74,6 +74,10 @@ var ( ErrTooManyWrites = errors.New("too many write commands") ErrConfigClusterEnabledHosts = errors.New("providing hosts to a non-disabled cluster is not allowed") + ErrConfigClusterTypeInvalid = errors.New("invalid cluster type") + ErrConfigHostsMissing = errors.New("missing bind address in cluster hosts") + + ErrClusterDoesNotOwnSlice = errors.New("cluster does not own slice") ) // Regular expression to validate index and frame names. diff --git a/server.go b/server.go index 1973f746e..338183cf7 100644 --- a/server.go +++ b/server.go @@ -115,9 +115,8 @@ func NewServer() *Server { Logger: NopLogger, } - s.Handler.Holder = s.Holder - s.diagnostics.server = s s.Handler.API = NewAPI(s.logger) + s.Handler.API.Holder = s.Holder return s } @@ -165,12 +164,11 @@ func (s *Server) Open() error { s.Cluster.MaxWritesPerRequest = s.MaxWritesPerRequest // Initialize HTTP handler. - s.Handler.Broadcaster = s.Broadcaster s.Handler.API.Broadcaster = s.Broadcaster - s.Handler.BroadcastHandler = s - s.Handler.StatusHandler = s - s.Handler.Node = node - s.Handler.Cluster = s.Cluster + s.Handler.API.BroadcastHandler = s + s.Handler.API.StatusHandler = s + s.Handler.API.URI = s.URI + s.Handler.API.Cluster = s.Cluster s.Handler.Executor = e s.Cluster.prefect = s.Handler diff --git a/server/server.go b/server/server.go index 66725f4ee..a50ad3e71 100644 --- a/server/server.go +++ b/server/server.go @@ -218,8 +218,7 @@ func (m *Command) SetupServer() error { } c := pilosa.GetHTTPClient(TLSConfig) m.Server.RemoteClient = c - m.Server.Handler.RemoteClient = c - m.Server.Cluster.RemoteClient = c + m.Server.Handler.API.RemoteClient = c // Statik file system. m.Server.Handler.FileSystem = &statik.FileSystem{}