diff --git a/client.go b/client.go index 8fe8ee7d8..ed6a1a2ab 100644 --- a/client.go +++ b/client.go @@ -45,8 +45,8 @@ func NewClient(host string) (*Client, error) { // Host returns the host the client was initialized with. func (c *Client) Host() string { return c.host } -// SliceNs returns the number of slices on a server by database. -func (c *Client) SliceNs(ctx context.Context) (MaxSlices, error) { +// MaxSliceByDatabase returns the number of slices on a server by database. +func (c *Client) MaxSliceByDatabase(ctx context.Context) (map[string]uint64, error) { // Execute request against the host. u := url.URL{ Scheme: "http", @@ -362,13 +362,13 @@ func (c *Client) BackupTo(ctx context.Context, w io.Writer, db, frame string) er tw := tar.NewWriter(w) // Find the maximum number of slices. - sliceNs, err := c.SliceNs(ctx) + maxSlices, err := c.MaxSliceByDatabase(ctx) if err != nil { return fmt.Errorf("slice n: %s", err) } // Backup every slice to the tar file. - for i := uint64(0); i <= sliceNs[db]; i++ { + for i := uint64(0); i <= maxSlices[db]; i++ { if err := c.backupSliceTo(ctx, tw, db, frame, i); err != nil { return err } diff --git a/cmd/pilosactl/main.go b/cmd/pilosactl/main.go index b7ab3247f..88b4c328f 100644 --- a/cmd/pilosactl/main.go +++ b/cmd/pilosactl/main.go @@ -480,13 +480,13 @@ func (cmd *ExportCommand) Run(ctx context.Context) error { } // Determine slice count. - sliceNs, err := client.SliceNs(ctx) + maxSlices, err := client.MaxSliceByDatabase(ctx) if err != nil { return err } // Export each slice. - for slice := uint64(0); slice <= sliceNs[cmd.Database]; slice++ { + for slice := uint64(0); slice <= maxSlices[cmd.Database]; slice++ { logger.Printf("exporting slice: %d", slice) if err := client.ExportCSV(ctx, cmd.Database, cmd.Frame, slice, w); err != nil { return err diff --git a/db.go b/db.go index fb7724388..056214f7d 100644 --- a/db.go +++ b/db.go @@ -116,8 +116,8 @@ func (db *DB) Close() error { return nil } -// SliceN returns the max slice in the database according to this node. -func (db *DB) SliceN() uint64 { +// MaxSlice returns the max slice in the database according to this node. +func (db *DB) MaxSlice() uint64 { if db == nil { return 0 } @@ -126,7 +126,7 @@ func (db *DB) SliceN() uint64 { max := db.remoteMaxSlice for _, f := range db.frames { - if slice := f.SliceN(); slice > max { + if slice := f.MaxSlice(); slice > max { max = slice } } diff --git a/executor.go b/executor.go index f966a78ec..3ef806940 100644 --- a/executor.go +++ b/executor.go @@ -53,10 +53,10 @@ func (e *Executor) Execute(ctx context.Context, db string, q *pql.Query, slices // If slices aren't specified, then include all of them. if len(slices) == 0 { // Round up the number of slices. - sliceN := e.Index.DB(db).SliceN() + maxSlice := e.Index.DB(db).MaxSlice() // Generate a slices of all slices. - slices = make([]uint64, sliceN+1) + slices = make([]uint64, maxSlice+1) for i := range slices { slices[i] = uint64(i) } @@ -713,7 +713,7 @@ func (e *Executor) mapReduce(ctx context.Context, db string, slices []uint64, c // Iterate over all map responses and reduce. var result interface{} - var sliceN int + var maxSlice int for { select { case <-ctx.Done(): @@ -738,8 +738,8 @@ func (e *Executor) mapReduce(ctx context.Context, db string, slices []uint64, c result = reduceFn(result, resp.result) // If all slices have been processed then return. - sliceN += len(resp.slices) - if sliceN >= len(slices) { + maxSlice += len(resp.slices) + if maxSlice >= len(slices) { return result, nil } } @@ -797,7 +797,7 @@ func (e *Executor) mapperLocal(ctx context.Context, slices []uint64, mapFn mapFu } // Reduce results - var sliceN int + var maxSlice int var result interface{} for { select { @@ -808,11 +808,11 @@ func (e *Executor) mapperLocal(ctx context.Context, slices []uint64, mapFn mapFu return nil, resp.err } result = reduceFn(result, resp.result) - sliceN++ + maxSlice++ } // Exit once all slices are processed. - if sliceN == len(slices) { + if maxSlice == len(slices) { return result, nil } } diff --git a/frame.go b/frame.go index 56e4b3f1a..c8f18f8a1 100644 --- a/frame.go +++ b/frame.go @@ -58,8 +58,8 @@ func (f *Frame) Path() string { return f.path } // BitmapAttrStore returns the attribute storage. func (f *Frame) BitmapAttrStore() *AttrStore { return f.bitmapAttrStore } -// SliceN returns the max slice in the frame. -func (f *Frame) SliceN() uint64 { +// MaxSlice returns the max slice in the frame. +func (f *Frame) MaxSlice() uint64 { f.mu.Lock() defer f.mu.Unlock() @@ -128,7 +128,7 @@ func (f *Frame) openFragments() error { frag.BitmapAttrStore = f.bitmapAttrStore f.fragments[frag.Slice()] = frag - f.stats.Count("sliceN", 1) + f.stats.Count("maxSlice", 1) } return nil @@ -202,7 +202,7 @@ func (f *Frame) createFragmentIfNotExists(slice uint64) (*Fragment, error) { // Save to lookup. f.fragments[slice] = frag - f.stats.Count("sliceN", 1) + f.stats.Count("maxSlice", 1) return frag, nil } diff --git a/handler.go b/handler.go index 153548414..c278ddbd7 100644 --- a/handler.go +++ b/handler.go @@ -103,7 +103,7 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { case "/slices/max": switch r.Method { case "GET": - h.handleGetMaxSlices(w, r) + h.handleGetSliceMax(w, r) default: http.Error(w, "method not allowed", http.StatusMethodNotAllowed) } @@ -252,8 +252,8 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } } -func (h *Handler) handleGetMaxSlices(w http.ResponseWriter, r *http.Request) error { - ms := h.Index.SliceNs() +func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) error { + ms := h.Index.MaxSlices() if strings.Contains(r.Header.Get("Accept"), "application/x-protobuf") { pb := &internal.MaxSlicesResponse{ MaxSlices: ms, @@ -271,7 +271,7 @@ func (h *Handler) handleGetMaxSlices(w http.ResponseWriter, r *http.Request) err } type sliceMaxResponse struct { - MaxSlices MaxSlices `json:"MaxSlices"` + MaxSlices map[string]uint64 `json:"MaxSlices"` } // handleDeleteDB handles DELETE /db request. @@ -818,7 +818,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request) } // Determine the maximum number of slices. - sliceNs, err := client.SliceNs(r.Context()) + maxSlices, err := client.MaxSliceByDatabase(r.Context()) if err != nil { http.Error(w, "cannot determine remote slice count: "+err.Error(), http.StatusInternalServerError) return @@ -826,7 +826,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request) // Loop over each slice and import it if this node owns it. //travis - for slice := uint64(0); slice <= sliceNs[db]; slice++ { + for slice := uint64(0); slice <= maxSlices[db]; slice++ { // Ignore this slice if we don't own it. if !h.Cluster.OwnsFragment(h.Host, db, slice) { continue diff --git a/index.go b/index.go index 189919434..6b5de11d5 100644 --- a/index.go +++ b/index.go @@ -106,14 +106,11 @@ func (i *Index) Close() error { return nil } -// MaxSlices contains the max known slice -by this node- for each DB -type MaxSlices map[string]uint64 - -// SliceNs returns MaxSlice map for all databases. -func (i *Index) SliceNs() MaxSlices { - a := make(MaxSlices) +// MaxSlices returns MaxSlice map for all databases. +func (i *Index) MaxSlices() map[string]uint64 { + a := make(map[string]uint64) for _, db := range i.DBs() { - a[db.Name()] = db.SliceN() + a[db.Name()] = db.MaxSlice() } return a } @@ -345,7 +342,7 @@ func (s *IndexSyncer) SyncIndex() error { return fmt.Errorf("frame sync error: db=%s, frame=%s, err=%s", di.Name, fi.Name, err) } - for slice := uint64(0); slice <= s.Index.DB(di.Name).SliceN(); slice++ { + for slice := uint64(0); slice <= s.Index.DB(di.Name).MaxSlice(); slice++ { // Ignore slices that this host doesn't own. if !s.Cluster.OwnsFragment(s.Host, di.Name, slice) { continue diff --git a/server.go b/server.go index b9f0a38b7..7323935fe 100644 --- a/server.go +++ b/server.go @@ -197,19 +197,26 @@ func (s *Server) monitorMaxSlices() { case <-ticker.C: } - oldmaxslices := s.Index.SliceNs() + oldmaxslices := s.Index.MaxSlices() for _, node := range s.Cluster.Nodes { if s.Host != node.Host { maxSlices, _ := checkMaxSlices(node.Host) for db, newmax := range maxSlices { - // we're not going to create a db locally if we don't know about it already - // TODO: consider changing this so we DO create a db locally - // do we want/need nodes to have empty files structures? + // if we don't know about a db locally, create it + // so that the /schema endpoint can report it if localdb := s.Index.DB(db); localdb != nil { if newmax > oldmaxslices[db] { oldmaxslices[db] = newmax localdb.SetRemoteMaxSlice(newmax) } + } else { + d, err := s.Index.CreateDBIfNotExists(db) + if err != nil { + s.logger().Printf("Failed to create DB locally: %s", db) + return + } + oldmaxslices[db] = newmax + d.SetRemoteMaxSlice(newmax) } } } @@ -217,7 +224,7 @@ func (s *Server) monitorMaxSlices() { } } -func checkMaxSlices(hostport string) (MaxSlices, error) { +func checkMaxSlices(hostport string) (map[string]uint64, error) { // Create HTTP request. req, err := http.NewRequest("GET", (&url.URL{ Scheme: "http",