diff --git a/api.go b/api.go index b29b3c1fc..649a99cd2 100644 --- a/api.go +++ b/api.go @@ -756,46 +756,17 @@ 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 { +// GetTranslateData provides a reader for key translation logs starting at offset. +func (api *API) GetTranslateData(ctx context.Context, offset int64) (io.ReadCloser, error) { rc, err := api.server.translateFile.Reader(ctx, offset) if err != nil { - return errors.Wrap(err, "read from translate store") + return nil, 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 + return rc, nil } // State returns the cluster state which is usually "NORMAL", but could be diff --git a/executor.go b/executor.go index 6f0e0ce8b..10026ae2c 100644 --- a/executor.go +++ b/executor.go @@ -101,9 +101,12 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar } // Translate query keys to ids, if necessary. - for i := range q.Calls { - if err := e.translateCall(index, idx, q.Calls[i]); err != nil { - return nil, err + // No need to translate a remote call. + if !opt.Remote { + for i := range q.Calls { + if err := e.translateCall(index, idx, q.Calls[i]); err != nil { + return nil, err + } } } @@ -113,10 +116,13 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar } // Translate response objects from ids to keys, if necessary. - for i := range results { - results[i], err = e.translateResult(index, idx, q.Calls[i], results[i]) - if err != nil { - return nil, err + // No need to translate a remote call. + if !opt.Remote { + for i := range results { + results[i], err = e.translateResult(index, idx, q.Calls[i], results[i]) + if err != nil { + return nil, err + } } } return results, nil diff --git a/http/handler.go b/http/handler.go index 8e0fcb709..2ac314d34 100644 --- a/http/handler.go +++ b/http/handler.go @@ -1308,14 +1308,14 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques 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) - pipeR, pipeW := io.Pipe() - - err := h.api.GetTranslateData(r.Context(), pipeW, offset) - + rdr, err := h.api.GetTranslateData(r.Context(), offset) if err != nil { if errors.Cause(err) == pilosa.ErrNotImplemented { http.Error(w, err.Error(), http.StatusNotImplemented) @@ -1331,7 +1331,28 @@ func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) w.Flush() } - io.Copy(w, pipeR) + // Copy from reader to client until store or client disconnect. + buf := make([]byte, translateStoreBufferSize) + for { + // Read from store. + n, err := rdr.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() + } + } } type queryValidationSpec struct { diff --git a/http/translator.go b/http/translator.go index d094a8c28..80c73bf5e 100644 --- a/http/translator.go +++ b/http/translator.go @@ -16,7 +16,7 @@ import ( // Ensure implementation implements inteface. var _ pilosa.TranslateStore = (*translateStore)(nil) -// translateStore represents an implementation of translateStore that +// translateStore represents an implementation of pilosa.TranslateStore that // communicates over HTTP. This is used with the TranslateHandler. type translateStore struct { URL string diff --git a/mock/translator.go b/mock/translator.go index 186c81894..7a63d8cf0 100644 --- a/mock/translator.go +++ b/mock/translator.go @@ -12,8 +12,8 @@ var _ pilosa.TranslateStore = (*TranslateStore)(nil) type TranslateStore struct { TranslateColumnsToUint64Func func(index string, values []string) ([]uint64, error) TranslateColumnToStringFunc func(index string, values uint64) (string, error) - TranslateRowsToUint64Func func(index, frame string, values []string) ([]uint64, error) - TranslateRowToStringFunc func(index, frame string, values uint64) (string, error) + TranslateRowsToUint64Func func(index, field string, values []string) ([]uint64, error) + TranslateRowToStringFunc func(index, field string, values uint64) (string, error) ReaderFunc func(ctx context.Context, off int64) (io.ReadCloser, error) } @@ -25,12 +25,12 @@ func (s TranslateStore) TranslateColumnToString(index string, values uint64) (st return s.TranslateColumnToStringFunc(index, values) } -func (s TranslateStore) TranslateRowsToUint64(index, frame string, values []string) ([]uint64, error) { - return s.TranslateRowsToUint64Func(index, frame, values) +func (s TranslateStore) TranslateRowsToUint64(index, field string, values []string) ([]uint64, error) { + return s.TranslateRowsToUint64Func(index, field, values) } -func (s TranslateStore) TranslateRowToString(index, frame string, value uint64) (string, error) { - return s.TranslateRowToStringFunc(index, frame, value) +func (s TranslateStore) TranslateRowToString(index, field string, value uint64) (string, error) { + return s.TranslateRowToStringFunc(index, field, value) } func (s TranslateStore) Reader(ctx context.Context, off int64) (io.ReadCloser, error) { diff --git a/server/config.go b/server/config.go index 73f0b289e..39a769ce2 100644 --- a/server/config.go +++ b/server/config.go @@ -74,7 +74,7 @@ type Config struct { // Translation config supports translation store replication. Translation struct { PrimaryURL string `toml:"primary-url"` - } + } `toml:"translation"` AntiEntropy struct { Interval toml.Duration `toml:"interval"` diff --git a/translate.go b/translate.go index 6b3c1856b..09eed00b2 100644 --- a/translate.go +++ b/translate.go @@ -39,8 +39,8 @@ type TranslateStore interface { TranslateColumnsToUint64(index string, values []string) ([]uint64, error) TranslateColumnToString(index string, values uint64) (string, error) - TranslateRowsToUint64(index, frame string, values []string) ([]uint64, error) - TranslateRowToString(index, frame string, values uint64) (string, error) + TranslateRowsToUint64(index, field string, values []string) ([]uint64, error) + TranslateRowToString(index, field string, values uint64) (string, error) // Returns a reader from the given offset of the raw data file. // The returned reader must be closed by the caller when done. @@ -64,7 +64,7 @@ type TranslateFile struct { closing chan struct{} cols map[string]*index - rows map[frameKey]*index + rows map[fieldKey]*index Path string mapSize int @@ -82,7 +82,7 @@ func NewTranslateFile() *TranslateFile { writeNotify: make(chan struct{}), closing: make(chan struct{}), cols: make(map[string]*index), - rows: make(map[frameKey]*index), + rows: make(map[fieldKey]*index), mapSize: defaultMapSize, @@ -201,7 +201,7 @@ func (s *TranslateFile) applyEntry(entry *LogEntry, offset int64) error { idx = s.col(string(entry.Index)) case LogEntryTypeInsertRow: - idx = s.row(string(entry.Index), string(entry.Frame)) + idx = s.row(string(entry.Index), string(entry.Field)) default: return fmt.Errorf("enterprise.TranslateFile.applyEntry(): unknown log entry type: 0x%20x", entry.Type) @@ -318,11 +318,11 @@ func (s *TranslateFile) col(index string) *index { return idx } -func (s *TranslateFile) row(index, frame string) *index { - idx := s.rows[frameKey{index, frame}] +func (s *TranslateFile) row(index, field string) *index { + idx := s.rows[fieldKey{index, field}] if idx == nil { idx = newIndex(s.data) - s.rows[frameKey{index, frame}] = idx + s.rows[fieldKey{index, field}] = idx } return idx } @@ -433,8 +433,8 @@ func (s *TranslateFile) TranslateColumnToString(index string, value uint64) (str return "", nil } -func (s *TranslateFile) TranslateRowsToUint64(index, frame string, values []string) ([]uint64, error) { - key := frameKey{index, frame} +func (s *TranslateFile) TranslateRowsToUint64(index, field string, values []string) ([]uint64, error) { + key := fieldKey{index, field} ret := make([]uint64, len(values)) @@ -495,7 +495,7 @@ func (s *TranslateFile) TranslateRowsToUint64(index, frame string, values []stri entry := &LogEntry{ Type: LogEntryTypeInsertRow, Index: []byte(index), - Frame: []byte(frame), + Field: []byte(field), IDs: make([]uint64, 0, len(values)), Keys: make([][]byte, 0, len(values)), } @@ -524,9 +524,9 @@ func (s *TranslateFile) TranslateRowsToUint64(index, frame string, values []stri return ret, nil } -func (s *TranslateFile) TranslateRowToString(index, frame string, id uint64) (string, error) { +func (s *TranslateFile) TranslateRowToString(index, field string, id uint64) (string, error) { s.mu.RLock() - if idx := s.rows[frameKey{index, frame}]; idx != nil { + if idx := s.rows[fieldKey{index, field}]; idx != nil { if ret, ok := idx.keyByID(id); ok { s.mu.RUnlock() return string(ret), nil @@ -548,7 +548,7 @@ func (s *TranslateFile) Reader(ctx context.Context, offset int64) (io.ReadCloser type LogEntry struct { Type uint8 Index []byte - Frame []byte + Field []byte IDs []uint64 Keys [][]byte @@ -558,12 +558,12 @@ type LogEntry struct { Length uint64 } -// headerSize returns the number of bytes required for size, type, index, frame, & pair count. +// headerSize returns the number of bytes required for size, type, index, field, & pair count. func (e *LogEntry) headerSize() int64 { sz := uVarintSize(e.Length) + // total entry length 1 + // type uVarintSize(uint64(len(e.Index))) + len(e.Index) + // Index length and data - uVarintSize(uint64(len(e.Frame))) + len(e.Frame) + // Frame length and data + uVarintSize(uint64(len(e.Field))) + len(e.Field) + // Field length and data uVarintSize(uint64(len(e.IDs))) // ID/Key pair count return int64(sz) } @@ -604,14 +604,14 @@ func (e *LogEntry) ReadFrom(r io.Reader) (_ int64, err error) { } } - // Read frame name. + // Read field name. if sz, err := binary.ReadUvarint(br); err != nil { return n64, err } else if sz == 0 { - e.Frame = nil + e.Field = nil } else { - e.Frame = make([]byte, sz) - if _, err := io.ReadFull(r, e.Frame); err != nil { + e.Field = make([]byte, sz) + if _, err := io.ReadFull(r, e.Field); err != nil { return n64, err } } @@ -663,11 +663,11 @@ func (e *LogEntry) WriteTo(w io.Writer) (_ int64, err error) { return 0, err } - // Write frame name. - sz = binary.PutUvarint(b, uint64(len(e.Frame))) + // Write field name. + sz = binary.PutUvarint(b, uint64(len(e.Field))) if _, err := buf.Write(b[:sz]); err != nil { return 0, err - } else if _, err := buf.Write(e.Frame); err != nil { + } else if _, err := buf.Write(e.Field); err != nil { return 0, err } @@ -722,9 +722,9 @@ func validLogEntriesLen(p []byte) (n int) { } } -type frameKey struct { +type fieldKey struct { index string - frame string + field string } const defaultLoadFactor = 90 @@ -910,7 +910,7 @@ type translateFileReader struct { closing chan struct{} } -// newTranslateFileReader returns a new instance of TranslateFileReader. +// newTranslateFileReader returns a new instance of translateFileReader. func newTranslateFileReader(ctx context.Context, store *TranslateFile, offset int64) *translateFileReader { return &translateFileReader{ ctx: ctx, diff --git a/translate_test.go b/translate_test.go index 693cca759..e41a4dd79 100644 --- a/translate_test.go +++ b/translate_test.go @@ -136,42 +136,42 @@ func TestTranslateFile_TranslateRow(t *testing.T) { defer s.MustClose() // First translation should start id at zero. - if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"foo"}); err != nil { + if ids, err := s.TranslateRowsToUint64("IDX0", "FIELD0", []string{"foo"}); err != nil { t.Fatal(err) } else if !reflect.DeepEqual(ids, []uint64{1}) { t.Fatalf("unexpected id: %#v", ids) } // Next translation on the same index should move to one. - if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"bar"}); err != nil { + if ids, err := s.TranslateRowsToUint64("IDX0", "FIELD0", []string{"bar"}); err != nil { t.Fatal(err) } else if !reflect.DeepEqual(ids, []uint64{2}) { t.Fatalf("unexpected id: %#v", ids) } // Translation on a different index restarts at 0. - if ids, err := s.TranslateRowsToUint64("IDX1", "FRAME0", []string{"bar"}); err != nil { + if ids, err := s.TranslateRowsToUint64("IDX1", "FIELD0", []string{"bar"}); err != nil { t.Fatal(err) } else if !reflect.DeepEqual(ids, []uint64{1}) { t.Fatalf("unexpected id: %#v", ids) } - // Translation on a different frame restarts at 0. - if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME1", []string{"bar"}); err != nil { + // Translation on a different field restarts at 0. + if ids, err := s.TranslateRowsToUint64("IDX0", "FIELD1", []string{"bar"}); err != nil { t.Fatal(err) } else if !reflect.DeepEqual(ids, []uint64{1}) { t.Fatalf("unexpected id: %#v", ids) } // Ensure that string values can be looked up by ID. - if value, err := s.TranslateRowToString("IDX0", "FRAME0", 2); err != nil { + if value, err := s.TranslateRowToString("IDX0", "FIELD0", 2); err != nil { t.Fatal(err) } else if value != "bar" { t.Fatalf("unexpected value: %s", value) } // Ensure that non-existent values return blank. - if value, err := s.TranslateRowToString("IDX0", "FRAME0", 1000); err != nil { + if value, err := s.TranslateRowToString("IDX0", "FIELD0", 1000); err != nil { t.Fatal(err) } else if value != "" { t.Fatalf("unexpected value: %s", value) @@ -182,22 +182,22 @@ func TestTranslateFile_TranslateRow(t *testing.T) { t.Fatal(err) } - // Translation on a different frame restarts at 0. - if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME1", []string{"bar"}); err != nil { + // Translation on a different field restarts at 0. + if ids, err := s.TranslateRowsToUint64("IDX0", "FIELD1", []string{"bar"}); err != nil { t.Fatal(err) } else if !reflect.DeepEqual(ids, []uint64{1}) { t.Fatalf("unexpected id: %#v", ids) } // Ensure that string values can be looked up by ID. - if value, err := s.TranslateRowToString("IDX0", "FRAME0", 2); err != nil { + if value, err := s.TranslateRowToString("IDX0", "FIELD0", 2); err != nil { t.Fatal(err) } else if value != "bar" { t.Fatalf("unexpected value: %s", value) } // Translate new row and increment sequence. - if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"baz"}); err != nil { + if ids, err := s.TranslateRowsToUint64("IDX0", "FIELD0", []string{"baz"}); err != nil { t.Fatal(err) } else if !reflect.DeepEqual(ids, []uint64{3}) { t.Fatalf("unexpected id: %#v", ids) @@ -215,7 +215,7 @@ func TestTranslateFile_TranslateRow_Large(t *testing.T) { keys[j] = strconv.Itoa(i + j + 1) } - ids, err := s.TranslateRowsToUint64("IDX0", "FRAME0", keys) + ids, err := s.TranslateRowsToUint64("IDX0", "FIELD0", keys) if err != nil { t.Fatal(err) } @@ -230,7 +230,7 @@ func TestTranslateFile_TranslateRow_Large(t *testing.T) { // Verify values can be returned. for i := 0; i < 1000000; i++ { exp := strconv.Itoa(i + 1) - if key, err := s.TranslateRowToString("IDX0", "FRAME0", uint64(i+1)); err != nil { + if key, err := s.TranslateRowToString("IDX0", "FIELD0", uint64(i+1)); err != nil { t.Fatal(err) } else if key != exp { t.Fatalf("unexpected key: got=%q, exp=%q", key, exp) @@ -243,7 +243,7 @@ func TestTranslateFile_TranslateRow_Large(t *testing.T) { } for i := 0; i < 1000000; i++ { exp := strconv.Itoa(i + 1) - if key, err := s.TranslateRowToString("IDX0", "FRAME0", uint64(i+1)); err != nil { + if key, err := s.TranslateRowToString("IDX0", "FIELD0", uint64(i+1)); err != nil { t.Fatal(err) } else if key != exp { t.Fatalf("unexpected key: got=%q, exp=%q", key, exp) @@ -257,7 +257,7 @@ func TestTranslateFile_Reader(t *testing.T) { defer s.MustClose() if _, err := s.TranslateColumnsToUint64("IDX0", []string{"foo"}); err != nil { t.Fatal(err) - } else if _, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"bar", "baz"}); err != nil { + } else if _, err := s.TranslateRowsToUint64("IDX0", "FIELD0", []string{"bar", "baz"}); err != nil { t.Fatal(err) } @@ -290,7 +290,7 @@ func TestTranslateFile_Reader(t *testing.T) { } else if diff := cmp.Diff(entry, pilosa.LogEntry{ Type: pilosa.LogEntryTypeInsertRow, Index: []byte("IDX0"), - Frame: []byte("FRAME0"), + Field: []byte("FIELD0"), IDs: []uint64{1, 2}, Keys: [][]byte{[]byte("bar"), []byte("baz")}, Length: 24, @@ -329,7 +329,7 @@ func TestTranslateFile_Reader(t *testing.T) { defer s.MustClose() if _, err := s.TranslateColumnsToUint64("IDX0", []string{"foo"}); err != nil { t.Fatal(err) - } else if _, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"bar", "baz"}); err != nil { + } else if _, err := s.TranslateRowsToUint64("IDX0", "FIELD0", []string{"bar", "baz"}); err != nil { t.Fatal(err) } @@ -348,7 +348,7 @@ func TestTranslateFile_Reader(t *testing.T) { } else if diff := cmp.Diff(entry, pilosa.LogEntry{ Type: pilosa.LogEntryTypeInsertRow, Index: []byte("IDX0"), - Frame: []byte("FRAME0"), + Field: []byte("FIELD0"), IDs: []uint64{1, 2}, Keys: [][]byte{[]byte("bar"), []byte("baz")}, Length: 24, @@ -392,7 +392,7 @@ func TestTranslateFile_PrimaryTranslateStore(t *testing.T) { // Write to the primary. if _, err := primary.TranslateColumnsToUint64("IDX0", []string{"foo"}); err != nil { t.Fatal(err) - } else if _, err := primary.TranslateRowsToUint64("IDX0", "FRAME0", []string{"bar", "baz"}); err != nil { + } else if _, err := primary.TranslateRowsToUint64("IDX0", "FIELD0", []string{"bar", "baz"}); err != nil { t.Fatal(err) } @@ -405,13 +405,13 @@ func TestTranslateFile_PrimaryTranslateStore(t *testing.T) { return fmt.Errorf("unexpected column 1 value: %s", value) } - if value, err := replica.TranslateRowToString("IDX0", "FRAME0", 1); err != nil { + if value, err := replica.TranslateRowToString("IDX0", "FIELD0", 1); err != nil { return err } else if value != "bar" { return fmt.Errorf("unexpected row 1 value: %s", value) } - if value, err := replica.TranslateRowToString("IDX0", "FRAME0", 2); err != nil { + if value, err := replica.TranslateRowToString("IDX0", "FIELD0", 2); err != nil { return err } else if value != "baz" { return fmt.Errorf("unexpected row 2 value: %s", value)