diff --git a/Makefile b/Makefile index 0ffeb96bd..2c2ac0456 100644 --- a/Makefile +++ b/Makefile @@ -19,7 +19,7 @@ DOCKER_BUILD= # set to 1 to use `docker-build` instead of `build` when creating BUILD_TAGS += shardwidth$(SHARD_WIDTH) TEST_TAGS = roaringparanoia define LICENSE_HASH_CODE - head -13 $1 | sed -e 's/Copyright 20[0-9][0-9]/Copyright 20XX/g' | shasum | cut -f 1 -d " " + head -13 $1 | sed -e 's/Copyright 20[0-9][0-9]/Copyright 20XX/' -e 's/Pilosa Corp\./Molecula Corp./' | shasum | cut -f 1 -d " " endef LICENSE_HASH=$(shell $(call LICENSE_HASH_CODE, pilosa.go)) UNAME := $(shell uname -s) diff --git a/api.go b/api.go index c9c998848..bbdb8c6b4 100644 --- a/api.go +++ b/api.go @@ -1430,16 +1430,6 @@ var ErrAborted = fmt.Errorf("error: update was aborted") func (api *API) ImportAtomicRecord(ctx context.Context, qcx *Qcx, req *AtomicRecord, opts ...ImportOption) error { - // this is because some of the tests pass nil qcx for convenience. - isLocalQcx := false - if qcx == nil { - isLocalQcx = true - qcx = api.Txf().NewQcx() - defer func() { - qcx.Abort() - }() - } - simPowerLoss := false lossAfter := -1 var opt ImportOptions @@ -1491,11 +1481,6 @@ func (api *API) ImportAtomicRecord(ctx context.Context, qcx *Qcx, req *AtomicRec return errors.Wrap(err, "ImportAtomicRecord ImportWithTx") } } - - // got to the end succesfully, so commit if we made the qcx - if isLocalQcx { - return qcx.Finish() - } return nil } @@ -1521,21 +1506,10 @@ func (api *API) Import(ctx context.Context, qcx *Qcx, req *ImportRequest, opts . if req.Clear { opts = addClearToImportOptions(opts) } - isLocalQcx := false - if qcx == nil { - isLocalQcx = true - qcx = api.Txf().NewQcx() - defer func() { - qcx.Abort() - }() - } err = api.ImportWithTx(ctx, qcx, req, opts...) if err != nil { return err } - if isLocalQcx { - return qcx.Finish() - } return nil } @@ -1633,21 +1607,19 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest, return errors.Wrap(err, "validating shard ownership") } - // Convert timestamps to time.Time. - timestamps := make([]*time.Time, len(req.Timestamps)) - for i, ts := range req.Timestamps { - if ts == 0 { - continue + var timestamps []int64 + for _, v := range req.Timestamps { + if v != 0 { + timestamps = req.Timestamps + break } - t := time.Unix(0, ts).UTC() - timestamps[i] = &t } // Import columnIDs into existence field. // Note: req.Shard may not be the only shard imported into here, // so don't expect it to be invariant. if !options.Clear { - if err := importExistenceColumns(qcx, idx, req.ColumnIDs); err != nil { + if err := importExistenceColumns(qcx, idx, req.ColumnIDs, req.Shard); err != nil { api.server.logger.Errorf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) return err } @@ -1657,7 +1629,7 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest, } // Import into fragment. - err = field.Import(qcx, req.RowIDs, req.ColumnIDs, timestamps, opts...) + err = field.Import(qcx, req.RowIDs, req.ColumnIDs, timestamps, req.Shard, opts...) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) return errors.Wrap(err, "importing") @@ -1752,20 +1724,13 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu // don't keep that list around since we don't need it anymore req.scratch = nil } - isLocalQcx := false - if qcx == nil { - isLocalQcx = true - qcx = api.Txf().NewQcx() - defer func() { - qcx.Abort() - }() - } // if we're importing into a specific shard if req.Shard != math.MaxUint64 { // Check that column IDs match the stated shard. - if s1, s2 := req.ColumnIDs[0]/ShardWidth, req.ColumnIDs[len(req.ColumnIDs)-1]/ShardWidth; s1 != s2 && s2 != req.Shard { - return errors.Errorf("shard %d specified, but import spans shards %d to %d", req.Shard, s1, s2) + shard := req.ColumnIDs[0] / ShardWidth + if s2 := req.ColumnIDs[len(req.ColumnIDs)-1] / ShardWidth; (shard != s2) || (shard != req.Shard) { + return errors.Errorf("shard %d specified, but import spans shards %d to %d", req.Shard, shard, s2) } // Validate shard ownership. TODO - we should forward to the // correct node rather than barfing here. @@ -1774,7 +1739,7 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu } // Import columnIDs into existence field. if !options.Clear { - if err := importExistenceColumns(qcx, idx, req.ColumnIDs); err != nil { + if err := importExistenceColumns(qcx, idx, req.ColumnIDs, shard); err != nil { api.server.logger.Errorf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) return errors.Wrap(err, "importing existence columns") } @@ -1782,17 +1747,17 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu // Import into fragment. if len(req.Values) > 0 { - err = field.importValue(qcx, req.ColumnIDs, req.Values, options) + err = field.importValue(qcx, req.ColumnIDs, req.Values, shard, options) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } } else if len(req.TimestampValues) > 0 { - err = field.importTimestampValue(qcx, req.ColumnIDs, req.TimestampValues, options) + err = field.importTimestampValue(qcx, req.ColumnIDs, req.TimestampValues, shard, options) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } } else if len(req.FloatValues) > 0 { - err = field.importFloatValue(qcx, req.ColumnIDs, req.FloatValues, options) + err = field.importFloatValue(qcx, req.ColumnIDs, req.FloatValues, shard, options) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } @@ -1847,20 +1812,23 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu if err != nil { return err } - if isLocalQcx { - return qcx.Finish() - } return nil } -func importExistenceColumns(qcx *Qcx, index *Index, columnIDs []uint64) error { +func importExistenceColumns(qcx *Qcx, index *Index, columnIDs []uint64, shard uint64) error { ef := index.existenceField() if ef == nil { return nil } existenceRowIDs := make([]uint64, len(columnIDs)) - return ef.Import(qcx, existenceRowIDs, columnIDs, nil) + // If we don't gratuitously hand-duplicate things in field.Import, + // the fact that fragment.bulkImport rewrites its row and column + // lists can burn us if we don't make a copy before doing the + // existence field write. + columnCopy := make([]uint64, len(columnIDs)) + copy(columnCopy, columnIDs) + return ef.Import(qcx, existenceRowIDs, columnCopy, nil, shard) } // ShardDistribution returns an object representing the distribution of shards diff --git a/api_test.go b/api_test.go index 4a4418ed9..201a7c89a 100644 --- a/api_test.go +++ b/api_test.go @@ -482,10 +482,10 @@ func TestAPI_ClearFlagForImportAndImportValues(t *testing.T) { } qcx := m0api.Txf().NewQcx() - if err := m0api.Import(ctx, qcx, ir0); err != nil { + if err := m0api.Import(ctx, qcx, ir0.Clone()); err != nil { t.Fatal(err) } - if err := m0api.ImportValue(ctx, qcx, ivr0); err != nil { + if err := m0api.ImportValue(ctx, qcx, ivr0.Clone()); err != nil { t.Fatal(err) } PanicOn(qcx.Finish()) diff --git a/bolt.go b/bolt.go index 366b77391..e679b1400 100644 --- a/bolt.go +++ b/bolt.go @@ -281,7 +281,7 @@ func (w *BoltWrapper) DeleteIndex(indexName string) error { // index name in the key prefix, so we cannot allow indexNames // themselves to contain apostrophies. if strings.Contains(indexName, "/") { - return fmt.Errorf("error: bad indexName `%v` in BoltWrapper.DeleteIndex() call: indexName cannot contain '/'.", indexName) + return fmt.Errorf("error: bad indexName `%v` in BoltWrapper.DeleteIndex() call: indexName cannot contain '/'", indexName) } prefix := txkey.IndexOnlyPrefix(indexName) return w.DeletePrefix(prefix) diff --git a/executor.go b/executor.go index ec071f34f..8cd24f33e 100644 --- a/executor.go +++ b/executor.go @@ -5939,6 +5939,7 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu ch := make(chan mapResponse, len(shards)) expected := 0 + shardLoop: for _, shard := range shards { j := job{ shard: shard, @@ -5949,7 +5950,7 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu } select { case <-done: - break + break shardLoop case e.work <- j: expected++ } diff --git a/executor_test.go b/executor_test.go index 834906864..8440adb9d 100644 --- a/executor_test.go +++ b/executor_test.go @@ -492,7 +492,7 @@ func TestExecutor(t *testing.T) { Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) Set(8, f=1, %s) - + Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) Set(2, f=10, 2001-01-01T00:00)`, nextDayExclusive.Format("2006-01-02T15:04")) @@ -539,7 +539,7 @@ func TestExecutor(t *testing.T) { Set("five", f=1, 2000-02-01T00:00) Set("six", f=1, 2001-01-01T00:00) Set("seven", f=1, 2002-01-01T02:00) - + Set("two", f=1, 1999-12-30T00:00) Set("two", f=1, 2002-02-01T00:00) Set("two", f=10, 2001-01-01T00:00)` @@ -573,7 +573,7 @@ func TestExecutor(t *testing.T) { Set(5, f="foo", 2000-02-01T00:00) Set(6, f="foo", 2001-01-01T00:00) Set(7, f="foo", 2002-01-01T02:00) - + Set(2, f="foo", 1999-12-30T00:00) Set(2, f="foo", 2002-02-01T00:00) Set(2, f="bar", 2001-01-01T00:00)` @@ -608,7 +608,7 @@ func TestExecutor(t *testing.T) { Set("five", f="foo", 2000-02-01T00:00) Set("six", f="foo", 2001-01-01T00:00) Set("seven", f="foo", 2002-01-01T02:00) - + Set("two", f="foo", 1999-12-30T00:00) Set("two", f="foo", 2002-02-01T00:00) Set("two", f="bar", 2001-01-01T00:00)` @@ -643,7 +643,7 @@ func TestExecutor(t *testing.T) { Set(5, f=1, 2000-02-01T00:00) Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) - + Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) Set(2, f=10, 2001-01-01T00:00)` @@ -678,7 +678,7 @@ func TestExecutor(t *testing.T) { Set(5, f=1, 2000-02-01T00:00) Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) - + Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) Set(2, f=10, 2001-01-01T00:00)` @@ -724,7 +724,7 @@ func TestExecutor(t *testing.T) { Set("five", f=1, 2000-02-01T00:00) Set("six", f=1, 2001-01-01T00:00) Set("seven", f=1, 2002-01-01T02:00) - + Set("two", f=1, 1999-12-30T00:00) Set("two", f=1, 2002-02-01T00:00) Set("two", f=10, 2001-01-01T00:00)` @@ -758,7 +758,7 @@ func TestExecutor(t *testing.T) { Set(5, f="foo", 2000-02-01T00:00) Set(6, f="foo", 2001-01-01T00:00) Set(7, f="foo", 2002-01-01T02:00) - + Set(2, f="foo", 1999-12-30T00:00) Set(2, f="foo", 2002-02-01T00:00) Set(2, f="bar", 2001-01-01T00:00)` @@ -793,7 +793,7 @@ func TestExecutor(t *testing.T) { Set("five", f="foo", 2000-02-01T00:00) Set("six", f="foo", 2001-01-01T00:00) Set("seven", f="foo", 2002-01-01T02:00) - + Set("two", f="foo", 1999-12-30T00:00) Set("two", f="foo", 2002-02-01T00:00) Set("two", f="bar", 2001-01-01T00:00)` @@ -4131,10 +4131,17 @@ func TestExecutor_Execute_All(t *testing.T) { req.ColumnIDs[bitCount-1] = uint64((3 * ShardWidth) + 2) m0 := c.GetNode(0) + // the request gets altered by the Import operation now... + reqs, err := req.Clone().ShardSplit() + if err != nil { + t.Fatalf("splitting request into shards: %v", err) + } qcx := m0.API.Txf().NewQcx() - if err := m0.API.Import(context.Background(), qcx, req); err != nil { - t.Fatal(err) + for _, r := range reqs { + if err := m0.API.Import(context.Background(), qcx, r); err != nil { + t.Fatal(err) + } } PanicOn(qcx.Finish()) @@ -5154,10 +5161,13 @@ func TestExecutor_GroupByStrings(t *testing.T) { }); err != nil { t.Fatalf("importing: %v", err) } + m0 := c.GetNode(0) + qcx := m0.API.Txf().NewQcx() + defer qcx.Abort() var v1, v2, v3, v4, v5, v6, v7, v8, v9, v10 int64 = 1, 2, 3, 4, 5, 6, 7, 8, 9, 10 var nv1, nv2, nv3, nv4 int64 = -1, -2, -3, -4 - if err := c.GetNode(0).API.ImportValue(context.Background(), nil, &pilosa.ImportValueRequest{ + if err := m0.API.ImportValue(context.Background(), qcx, &pilosa.ImportValueRequest{ Index: "istring", Field: "v", Shard: 0, @@ -5167,7 +5177,7 @@ func TestExecutor_GroupByStrings(t *testing.T) { t.Fatalf("importing: %v", err) } - if err := c.GetNode(0).API.ImportValue(context.Background(), nil, &pilosa.ImportValueRequest{ + if err := m0.API.ImportValue(context.Background(), qcx, &pilosa.ImportValueRequest{ Index: "istring", Field: "vv", Shard: 0, @@ -5177,7 +5187,7 @@ func TestExecutor_GroupByStrings(t *testing.T) { t.Fatalf("importing: %v", err) } - if err := c.GetNode(0).API.ImportValue(context.Background(), nil, &pilosa.ImportValueRequest{ + if err := m0.API.ImportValue(context.Background(), qcx, &pilosa.ImportValueRequest{ Index: "istring", Field: "nv", Shard: 0, diff --git a/field.go b/field.go index 0bb88d6ce..5f147d735 100644 --- a/field.go +++ b/field.go @@ -935,7 +935,7 @@ func (f *Field) RowTime(tx Tx, rowID uint64, time time.Time, quantum string) (*R if !TimeQuantum(quantum).Valid() { return nil, ErrInvalidTimeQuantum } - viewname := viewsByTime(viewStandard, time, TimeQuantum(quantum[len(quantum)-1:]))[0] + viewname := viewByTimeUnit(viewStandard, time, rune(quantum[len(quantum)-1])) view := f.view(viewname) if view == nil { return nil, errors.Errorf("view with quantum %v not found.", quantum) @@ -1437,7 +1437,7 @@ func (f *Field) Range(qcx *Qcx, name string, op pql.Token, predicate int64) (*Ro } // Import bulk imports data. -func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time.Time, opts ...ImportOption) (err0 error) { +func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, shard uint64, opts ...ImportOption) (err0 error) { // Set up import options. options := &ImportOptions{} @@ -1450,18 +1450,81 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time. // Determine quantum if timestamps are set. q := f.TimeQuantum() - if hasTime(timestamps) { + if len(timestamps) > 0 { if q == "" { return errors.New("time quantum not set in field") } else if options.Clear { return errors.New("import clear is not supported with timestamps") } + } else { + // short path: if we don't have any timestamps, we only need + // to write to exactly one view, which is always viewStandard, + // and *every* bit goes into that view, and we already verified that + // everything is in the same shard, so we can skip most of this. + fieldType := f.Type() + if fieldType == FieldTypeBool { + for _, rowID := range rowIDs { + if rowID > 1 { + return errors.New("bool field imports only support values 0 and 1") + } + } + } + tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: shard}) + if err != nil { + return errors.Wrap(err, "qcx.GetTx") + } + var err1 error + defer finisher(&err1) + view, err := f.createViewIfNotExists(viewStandard) + if err != nil { + return errors.Wrapf(err, "creating view %s", viewStandard) + } + + frag, err := view.CreateFragmentIfNotExists(shard) + if err != nil { + return errors.Wrap(err, "creating fragment") + } + + err1 = frag.bulkImport(tx, rowIDs, columnIDs, options) + return err1 } fieldType := f.Type() // Split import data by fragment. - dataByFragment := make(map[importKey]importData) + views := make(map[string]*importData) + var timeStringBuf []byte + var timeViews [][]byte + if len(q) > 0 { + // We're supporting time quantums, so we need to store bits in a + // number of views for every entry with a timestamp. We want to compute + // time quantum view names for whatever combination of YMDH views + // we have. But we don't want to allocate four strings per entry, or + // recompute and recreate the entire string. We know that only the + // YYYYMMDDHH part of the string changes over time. + timeStringBuf = make([]byte, len(viewStandard) + 11) + copy(timeStringBuf, []byte(viewStandard)) + copy(timeStringBuf[len(viewStandard):], []byte("_YYYYMMDDHH")) + // Now we have a buffer that contains + // `standard_YYYYMMDDHH`. We also need storage space to hold several + // slice headers, one per entry in q. These will hold the view names + // corresponding to each letter in q. + timeViews = make([][]byte, len(q)) + } + // This helper function records that a given column/row pair is relevant + // to a specific view. We use a map lookup for the strings, but do the + // actual operations using a slice so we're only writing each map entry + // once, not once on every update. + see := func(name []byte, columnID uint64, rowID uint64) { + var ok bool + var data *importData + if data, ok = views[string(name)]; !ok { + data = &importData{} + views[string(name)] = data + } + data.RowIDs = append(data.RowIDs, rowID) + data.ColumnIDs = append(data.ColumnIDs, columnID) + } for i := range rowIDs { rowID, columnID := rowIDs[i], columnIDs[i] @@ -1470,56 +1533,44 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time. return errors.New("bool field imports only support values 0 and 1") } - var timestamp *time.Time - if len(timestamps) > i { - timestamp = timestamps[i] - } + hasTime := len(timestamps) > i && timestamps[i] != 0 - var standard []string - if timestamp == nil { - standard = []string{viewStandard} - } else { - standard = viewsByTime(viewStandard, *timestamp, q) - if !f.options.NoStandardView { - // In order to match the logic of `SetBit()`, we want bits - // with timestamps to write to both time and standard views. - standard = append(standard, viewStandard) + // attach bit to standard view unless we have a timestamp and + // have the NoStandardView option set + if !hasTime || !f.options.NoStandardView { + see([]byte(viewStandard), columnID, rowID) + } + if hasTime { + // attach bit to all the views for this timestamp. note that the + // `timeViews` slice gets resliced and reused by this process, so + // we don't have to allocate millions of tiny slices of slice headers. + timeViews = viewsByTimeInto(timeStringBuf, timeViews, time.Unix(0, timestamps[i]).UTC(), q) + for _, v := range timeViews { + see(v, columnID, rowID) } } - - // Attach bit to each standard view. - for _, name := range standard { - key := importKey{View: name, Shard: columnID / ShardWidth} - data := dataByFragment[key] - data.RowIDs = append(data.RowIDs, rowID) - data.ColumnIDs = append(data.ColumnIDs, columnID) - dataByFragment[key] = data - } } - - // Import into each fragment. - for key, data := range dataByFragment { - view, err := f.createViewIfNotExists(key.View) + tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: shard}) + if err != nil { + return errors.Wrap(err, "qcx.GetTx") + } + var err1 error + defer finisher(&err1) + for viewName, data := range views { + view, err := f.createViewIfNotExists(viewName) if err != nil { - return errors.Wrap(err, "creating view") + return errors.Wrapf(err, "creating view %s", viewName) } - frag, err := view.CreateFragmentIfNotExists(key.Shard) + frag, err := view.CreateFragmentIfNotExists(shard) if err != nil { return errors.Wrap(err, "creating fragment") } - tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: frag.idx, Fragment: frag, Shard: frag.shard}) - if err != nil { - return errors.Wrap(err, "qcx.GetTx") - } - - err1 := frag.bulkImport(tx, data.RowIDs, data.ColumnIDs, options) + err1 = frag.bulkImport(tx, data.RowIDs, data.ColumnIDs, options) if err1 != nil { - finisher(&err1) return err1 } - finisher(nil) } return nil } @@ -1527,7 +1578,7 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time. // importFloatValue imports floating point values. In current usage, this // should only ever be called with data for a single shard; the API calls // around this are splitting it up per shard. -func (f *Field) importFloatValue(qcx *Qcx, columnIDs []uint64, values []float64, options *ImportOptions) error { +func (f *Field) importFloatValue(qcx *Qcx, columnIDs []uint64, values []float64, shard uint64, options *ImportOptions) error { // convert values to int64 values based on scale ivalues := make([]int64, len(values)) bsig := f.bsiGroup(f.name) @@ -1539,13 +1590,13 @@ func (f *Field) importFloatValue(qcx *Qcx, columnIDs []uint64, values []float64, ivalues[i] = int64(fval * mult) } // then call importValue - return f.importValue(qcx, columnIDs, ivalues, options) + return f.importValue(qcx, columnIDs, ivalues, shard, options) } // importFloatValue imports timestamp values. In current usage, this // should only ever be called with data for a single shard; the API calls // around this are splitting it up per shard. -func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time.Time, options *ImportOptions) error { +func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time.Time, shard uint64, options *ImportOptions) error { ivalues := make([]int64, len(values)) bsig := f.bsiGroup(f.name) if bsig == nil { @@ -1555,13 +1606,13 @@ func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time for i, t := range values { ivalues[i] = t.UnixNano() / TimeUnitNanos(f.options.TimeUnit) } - return f.importValue(qcx, columnIDs, ivalues, options) + return f.importValue(qcx, columnIDs, ivalues, shard, options) } // importValue bulk imports range-encoded value data. This function should // only be called with data for a single shard; the API calls that wrap // this handle splitting the data up per-shard. -func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, options *ImportOptions) (err0 error) { +func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, shard uint64, options *ImportOptions) (err0 error) { // no data to import if len(columnIDs) == 0 { return nil @@ -1614,9 +1665,9 @@ func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, option } f.mu.Unlock() - // Since all data should be for the same shard, we can just compute - // this from the first value. - shard := columnIDs[0] / ShardWidth + if columnIDs[0]/ShardWidth != shard { + return fmt.Errorf("requested import for shard %d, got record ID for shard %d", shard, columnIDs[0]/ShardWidth) + } view, err := f.createViewIfNotExists(viewName) if err != nil { @@ -2092,13 +2143,14 @@ const ( TimeUnitSeconds = "s" TimeUnitMilliseconds = "ms" TimeUnitMicroseconds = "µs" + TimeUnitUSeconds = "us" TimeUnitNanoseconds = "ns" ) // IsValidTimeUnit returns true if unit is valid. func IsValidTimeUnit(unit string) bool { switch unit { - case TimeUnitSeconds, TimeUnitMilliseconds, TimeUnitMicroseconds, TimeUnitNanoseconds: + case TimeUnitSeconds, TimeUnitMilliseconds, TimeUnitMicroseconds, TimeUnitUSeconds, TimeUnitNanoseconds: return true default: return false @@ -2112,7 +2164,7 @@ func TimeUnitNanos(unit string) int64 { return int64(time.Second) case TimeUnitMilliseconds: return int64(time.Millisecond) - case TimeUnitMicroseconds: + case TimeUnitMicroseconds, TimeUnitUSeconds: return int64(time.Microsecond) default: return int64(time.Nanosecond) diff --git a/field_internal_test.go b/field_internal_test.go index 41f288016..fc2296845 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -494,7 +494,7 @@ func TestBSIGroup_importValue(t *testing.T) { []uint64{100}, }, } { - if err := f.importValue(qcx, tt.columnIDs, tt.values, options); err != nil { + if err := f.importValue(qcx, tt.columnIDs, tt.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } PanicOn(qcx.Finish()) @@ -514,7 +514,8 @@ func TestBSIGroup_importValue(t *testing.T) { func benchmarkFieldImportValues(b *testing.B, qcx *Qcx, bitDepth uint64, f *TestField, cfunc func(uint64) uint64) { batches := makeBenchmarkImportValueData(b, bitDepth, cfunc) for _, req := range batches { - err := f.importValue(qcx, req.ColumnIDs, req.Values, &ImportOptions{}) + // NOTE: We assume everything's in Shard 0 for now. + err := f.importValue(qcx, req.ColumnIDs, req.Values, 0, &ImportOptions{}) if err != nil { b.Fatalf("error importing values: %s", err) } @@ -591,7 +592,7 @@ func TestIntField_MinMaxForShard(t *testing.T) { }, } { t.Run(test.name+strconv.Itoa(i), func(t *testing.T) { - if err := f.importValue(qcx, test.columnIDs, test.values, options); err != nil { + if err := f.importValue(qcx, test.columnIDs, test.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } PanicOn(qcx.Finish()) @@ -751,7 +752,7 @@ func TestDecimalField_MinMaxForShard(t *testing.T) { }, } { t.Run(test.name+strconv.Itoa(i), func(t *testing.T) { - if err := f.importFloatValue(qcx, test.columnIDs, test.values, options); err != nil { + if err := f.importFloatValue(qcx, test.columnIDs, test.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } @@ -811,7 +812,7 @@ func TestBSIGroup_TxReopenDB(t *testing.T) { []uint64{100}, }, } { - if err := f.importValue(qcx, tt.columnIDs, tt.values, options); err != nil { + if err := f.importValue(qcx, tt.columnIDs, tt.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } PanicOn(qcx.Finish()) diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 7de7b5c1f..3820d25fa 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1020,6 +1020,10 @@ func BenchmarkFragment_SetValue(b *testing.B) { } } +// makeBenchmarkImportValueData produces data that's supposed to be all within +// the same shard; for fragment purposes, implicitly shard 0. This also gets +// used by the field tests, but import requests are supposed to be per-shard, +// so it's important that we generate values only within a given shard. func makeBenchmarkImportValueData(b *testing.B, bitDepth uint64, cfunc func(uint64) uint64) []ImportValueRequest { b.StopTimer() column := uint64(0) diff --git a/handler.go b/handler.go index 49c4b1ff5..a07c3d6c9 100644 --- a/handler.go +++ b/handler.go @@ -18,6 +18,7 @@ import ( "encoding/json" "time" + "github.com/molecula/featurebase/v2/shardwidth" "github.com/molecula/featurebase/v2/tracing" "github.com/pkg/errors" ) @@ -129,6 +130,41 @@ type ImportValueRequest struct { scratch []int // scratch space to allow us to get a stable sort in reasonable time } +func (ivr *ImportValueRequest) Clone() *ImportValueRequest { + newIVR := &ImportValueRequest{} + if ivr == nil { + return newIVR + } + *newIVR = *ivr + // don't copy the internal scratch buffer + newIVR.scratch = nil + if len(ivr.ColumnIDs) > 0 { + newIVR.ColumnIDs = make([]uint64, len(ivr.ColumnIDs)) + copy(newIVR.ColumnIDs, ivr.ColumnIDs) + } + if len(ivr.ColumnKeys) > 0 { + newIVR.ColumnKeys = make([]string, len(ivr.ColumnKeys)) + copy(newIVR.ColumnKeys, ivr.ColumnKeys) + } + if len(ivr.Values) > 0 { + newIVR.Values = make([]int64, len(ivr.Values)) + copy(newIVR.Values, ivr.Values) + } + if len(ivr.FloatValues) > 0 { + newIVR.FloatValues = make([]float64, len(ivr.FloatValues)) + copy(newIVR.FloatValues, ivr.FloatValues) + } + if len(ivr.TimestampValues) > 0 { + newIVR.TimestampValues = make([]time.Time, len(ivr.TimestampValues)) + copy(newIVR.TimestampValues, ivr.TimestampValues) + } + if len(ivr.StringValues) > 0 { + newIVR.StringValues = make([]string, len(ivr.StringValues)) + copy(newIVR.StringValues, ivr.StringValues) + } + return newIVR +} + // AtomicRecord applies all its Ivr and Ivr atomically, in a Tx. // The top level Shard has to agree with Ivr[i].Shard and the Iv[i].Shard // for all i included (in Ivr and Ir). The same goes for the top level Index: all records @@ -142,6 +178,19 @@ type AtomicRecord struct { Ir []*ImportRequest // other field types, e.g. single bit } +func (ar *AtomicRecord) Clone() *AtomicRecord { + newAR := &AtomicRecord{Index: ar.Index, Shard: ar.Shard} + newAR.Ivr = make([]*ImportValueRequest, len(ar.Ivr)) + for i, vr := range ar.Ivr { + newAR.Ivr[i] = vr.Clone() + } + newAR.Ir = make([]*ImportRequest, len(ar.Ir)) + for i, vr := range ar.Ir { + newAR.Ir[i] = vr.Clone() + } + return newAR +} + func (ivr *ImportValueRequest) Len() int { return len(ivr.ColumnIDs) } func (ivr *ImportValueRequest) Less(i, j int) bool { if ivr.ColumnIDs[i] < ivr.ColumnIDs[j] { @@ -225,6 +274,75 @@ type ImportRequest struct { Clear bool } +// Clone allows copying an import request. Normally you wouldn't, but +// some import functions are destructive on their inputs, and if you +// want to *re-use* an import request, you might need this. If you're +// using this outside tx_test, something is probably wrong. +func (ir *ImportRequest) Clone() *ImportRequest { + newIR := &ImportRequest{} + if ir == nil { + return newIR + } + *newIR = *ir + if ir.RowIDs != nil { + newIR.RowIDs = make([]uint64, len(ir.RowIDs)) + copy(newIR.RowIDs, ir.RowIDs) + } + if ir.ColumnIDs != nil { + newIR.ColumnIDs = make([]uint64, len(ir.ColumnIDs)) + copy(newIR.ColumnIDs, ir.ColumnIDs) + } + if ir.RowKeys != nil { + newIR.RowKeys = make([]string, len(ir.RowKeys)) + copy(newIR.RowKeys, ir.RowKeys) + } + if ir.ColumnKeys != nil { + newIR.ColumnKeys = make([]string, len(ir.ColumnKeys)) + copy(newIR.ColumnKeys, ir.ColumnKeys) + } + if ir.Timestamps != nil { + newIR.Timestamps = make([]int64, len(ir.Timestamps)) + copy(newIR.Timestamps, ir.Timestamps) + } + return newIR +} + +// ShardSplit splits the request into a slice of import requests. It requires +// that the original request have all elements sorted, and already have +// column IDs, not column keys. +func (ir *ImportRequest) ShardSplit() ([]*ImportRequest, error) { + if ir == nil { + return nil, nil + } + // fix shard + if len(ir.ColumnIDs) < 2 { + ir.Shard = ir.ColumnIDs[0] >> shardwidth.Exponent + return []*ImportRequest{ir}, nil + } + shards, ends := shardwidth.FindShards(ir.ColumnIDs) + out := make([]*ImportRequest, len(shards)) + prev := 0 + for i, shard := range shards { + next := ends[i] + newIR := &ImportRequest{} + *newIR = *ir + newIR.ColumnIDs = ir.ColumnIDs[prev:next:next] + if ir.RowIDs != nil { + newIR.RowIDs = ir.RowIDs[prev:next:next] + } + if ir.RowKeys != nil { + newIR.RowKeys = ir.RowKeys[prev:next:next] + } + if ir.Timestamps != nil { + newIR.Timestamps = ir.Timestamps[prev:next:next] + } + newIR.Shard = shard + out[i] = newIR + prev = next + } + return out, nil +} + // ValidateWithTimestamp ensures that the payload of the request is valid. func (ir *ImportRequest) ValidateWithTimestamp(indexCreatedAt, fieldCreatedAt int64) error { if (ir.IndexCreatedAt != 0 && ir.IndexCreatedAt != indexCreatedAt) || diff --git a/http/client.go b/http/client.go index a51f54233..ce31012d7 100644 --- a/http/client.go +++ b/http/client.go @@ -30,7 +30,7 @@ import ( "strings" "time" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/encoding/proto" pnet "github.com/molecula/featurebase/v2/net" "github.com/molecula/featurebase/v2/topology" diff --git a/index.go b/index.go index b7204007f..25b7bd5ee 100644 --- a/index.go +++ b/index.go @@ -22,7 +22,6 @@ import ( "sort" "strconv" "sync" - "time" "github.com/molecula/featurebase/v2/disco" "github.com/molecula/featurebase/v2/roaring" @@ -841,21 +840,6 @@ type IndexOptions struct { TrackExistence bool `json:"trackExistence"` } -// hasTime returns true if a contains a non-nil time. -func hasTime(a []*time.Time) bool { - for _, t := range a { - if t != nil { - return true - } - } - return false -} - -type importKey struct { - View string - Shard uint64 -} - type importData struct { RowIDs []uint64 ColumnIDs []uint64 diff --git a/rbf.go b/rbf.go index e75e7768e..64a268281 100644 --- a/rbf.go +++ b/rbf.go @@ -542,7 +542,7 @@ func (w *RbfDBWrapper) DeleteField(index, field, fieldPath string) error { func (w *RbfDBWrapper) DeleteIndex(indexName string) error { if strings.Contains(indexName, "'") { - return fmt.Errorf("error: bad indexName `%v` in RbfDBWrapper.DeleteIndex() call: indexName cannot contain apostrophes/single quotes.", indexName) + return fmt.Errorf("error: bad indexName `%v` in RbfDBWrapper.DeleteIndex() call: indexName cannot contain apostrophes/single quotes", indexName) } prefix := txkey.IndexOnlyPrefix(indexName) diff --git a/shardwidth/helper.go b/shardwidth/helper.go new file mode 100644 index 000000000..c04e2e72b --- /dev/null +++ b/shardwidth/helper.go @@ -0,0 +1,70 @@ +// Copyright 2021 Molecula Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package shardwidth + +import ( + "math/bits" +) + +// FindNextShard returns the index of the first item which is not in the +// same shard as i. The index it returns may be equal to the length of the +// haystack, indicatincg that the rest of the list is in the same shard. +func FindNextShard(i int, haystack []uint64) int { + // compute the last thing that's in the same shard as haystack[i]. + if i >= len(haystack) { + return i + } + // current shard: + shard := (haystack[i] >> Exponent) + // last value in shard: + shardEnd := ((shard + 1) << Exponent) - 1 + j := i + // We want to do a binary search of the haystack. For any length of + // haystack, its topmost bit gives us a reasonable halfway point; it may + // not actually be halfway, but the number of steps it'll take to search + // it will be the same as if it were. sort.Search has interface overhead + // and makes us sad. + for incr := 1 << (bits.Len64(uint64(len(haystack) - i))); incr > 0; incr >>= 1 { + if j+incr < len(haystack) { + if haystack[j+incr] <= shardEnd { + j += incr + } + } + } + // we've found the last item that is in the same shard as i, so... + return j + 1 +} + +// FindShards finds the shards in a given haystack +func FindShards(haystack []uint64) (shards []uint64, endIndexes []int) { + if len(haystack) == 0 { + return nil, nil + } + index := 0 + // the steady state of this loop is that shards contains the current + // shard, but not its ending index; each time we find a new ending + // index, we record that index as the end for the current shard, and + // the new shard, until we reach the end and append len(haystack) + // as the last index. + shards = []uint64{haystack[index] >> Exponent} + index = FindNextShard(index, haystack) + for index < len(haystack) { + shards = append(shards, haystack[index]>>Exponent) + endIndexes = append(endIndexes, index) + index = FindNextShard(index, haystack) + } + endIndexes = append(endIndexes, index) + return shards, endIndexes +} diff --git a/shardwidth/helper_test.go b/shardwidth/helper_test.go new file mode 100644 index 000000000..cd1806cd8 --- /dev/null +++ b/shardwidth/helper_test.go @@ -0,0 +1,103 @@ +// Copyright 2021 Molecula Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package shardwidth_test + +import ( + "math/rand" + "testing" + + "github.com/molecula/featurebase/v2/shardwidth" +) + +type nextShardTestCase struct { + name string + haystack [][2]uint64 // stored as shard, offset pairs + shardIndexes []int +} + +var nextShardTestCases = []nextShardTestCase{ + { + name: "all-in-one", + haystack: [][2]uint64{ + {0, 0}, + {0, 1}, + }, + shardIndexes: []int{2}, + }, + { + name: "split", + haystack: [][2]uint64{ + {0, 0}, + {1, 1}, + }, + shardIndexes: []int{1, 2}, + }, + { + name: "two-and-one", + haystack: [][2]uint64{ + {0, 0}, + {0, 1}, + {1, 1}, + }, + shardIndexes: []int{2, 3}, + }, +} + +func TestFindShards(t *testing.T) { + for _, c := range nextShardTestCases { + haystack := make([]uint64, len(c.haystack)) + for i, h := range c.haystack { + haystack[i] = (h[0] << shardwidth.Exponent) + h[1] + } + _, indexes := shardwidth.FindShards(haystack) + if len(indexes) != len(c.shardIndexes) { + t.Fatalf("%s: expected %d, got %d", c.name, c.shardIndexes, indexes) + } + for i, expected := range c.shardIndexes { + if indexes[i] != expected { + t.Fatalf("%s: expected index %d to be %d, got %d", c.name, i, expected, indexes[i]) + } + } + } + // fake up some more test cases + for i := 0; i < 100; i++ { + haystack := make([]uint64, 100) + shard := uint64(0) + bit := uint64(0) + shardIndexes := []int{} + for j := range haystack { + if rand.Intn(30) == 0 { + if j > 0 { + shardIndexes = append(shardIndexes, j) + } + shard++ + bit = 0 + } else { + bit += uint64(rand.Intn(30)) + } + haystack[j] = (shard << shardwidth.Exponent) + bit + } + shardIndexes = append(shardIndexes, len(haystack)) + _, indexes := shardwidth.FindShards(haystack) + if len(indexes) != len(shardIndexes) { + t.Fatalf("trial %d: expected %d, got %d", i, shardIndexes, indexes) + } + for idx, expected := range shardIndexes { + if indexes[idx] != expected { + t.Fatalf("trial %d: expected index %d to be %d, got %d", i, idx, expected, indexes[idx]) + } + } + } +} diff --git a/test/cluster.go b/test/cluster.go index 7ff157094..126245a46 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -23,7 +23,7 @@ import ( "testing" "time" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/api/client" "github.com/molecula/featurebase/v2/disco" "github.com/molecula/featurebase/v2/logger" @@ -213,32 +213,37 @@ func (c *Cluster) ImportBitsWithTimestamp(t testing.TB, index, field string, row if com.API.Node().ID != node.ID { continue } - if len(timestamps) == 0 { - err := com.API.Import(context.Background(), nil, &pilosa.ImportRequest{ - Index: index, - Field: field, - Shard: shard, - RowIDs: rowIDs, - ColumnIDs: colIDs, - }) - if err != nil { - t.Fatalf("importing data: %v", err) - } - } else { - ts := byShardTs[shard] - err := com.API.Import(context.Background(), nil, &pilosa.ImportRequest{ - Index: index, - Field: field, - Shard: shard, - RowIDs: rowIDs, - ColumnIDs: colIDs, - Timestamps: ts, - }) - if err != nil { - t.Fatalf("importing data: %v", err) - } + func() { + qcx := com.API.Txf().NewQcx() + defer qcx.Abort() + if len(timestamps) == 0 { + err := com.API.Import(context.Background(), qcx, &pilosa.ImportRequest{ + Index: index, + Field: field, + Shard: shard, + RowIDs: rowIDs, + ColumnIDs: colIDs, + }) + if err != nil { + t.Fatalf("importing data: %v", err) + } + } else { + ts := byShardTs[shard] + err := com.API.Import(context.Background(), qcx, &pilosa.ImportRequest{ + Index: index, + Field: field, + Shard: shard, + RowIDs: rowIDs, + ColumnIDs: colIDs, + Timestamps: ts, + }) + if err != nil { + t.Fatalf("importing data: %v", err) + } + + } + }() - } } } } @@ -262,7 +267,9 @@ func (c *Cluster) ImportKeyKey(t testing.TB, index, field string, valAndRecKeys importRequest.RowKeys[i] = vk[0] importRequest.ColumnKeys[i] = vk[1] } - err := c.GetPrimary().API.Import(context.Background(), nil, importRequest) + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + err := c.GetPrimary().API.Import(context.Background(), qcx, importRequest) if err != nil { t.Fatalf("importing keykey data: %v", err) } @@ -292,7 +299,9 @@ func (c *Cluster) ImportTimeQuantumKey(t testing.TB, index, field string, entrie importRequest.Timestamps[i] = entry.Ts } - err := c.GetPrimary().API.Import(context.Background(), nil, importRequest) + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + err := c.GetPrimary().API.Import(context.Background(), qcx, importRequest) if err != nil { t.Fatalf("importing keykey data: %v", err) } @@ -318,7 +327,9 @@ func (c *Cluster) ImportIntKey(t testing.TB, index, field string, pairs []IntKey importRequest.Values[i] = pair.Val importRequest.ColumnKeys[i] = pair.Key } - if err := c.GetPrimary().API.ImportValue(context.Background(), nil, importRequest); err != nil { + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + if err := c.GetPrimary().API.ImportValue(context.Background(), qcx, importRequest); err != nil { t.Fatalf("importing IntKey data: %v", err) } } @@ -342,7 +353,9 @@ func (c *Cluster) ImportIntID(t testing.TB, index, field string, pairs []IntID) importRequest.Values[i] = pair.Val importRequest.ColumnIDs[i] = pair.ID } - if err := c.GetPrimary().API.ImportValue(context.Background(), nil, importRequest); err != nil { + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + if err := c.GetPrimary().API.ImportValue(context.Background(), qcx, importRequest); err != nil { t.Fatalf("importing IntID data: %v", err) } } @@ -367,7 +380,9 @@ func (c *Cluster) ImportIDKey(t testing.TB, index, field string, pairs []KeyID) importRequest.RowIDs[i] = pair.ID importRequest.ColumnKeys[i] = pair.Key } - err := c.GetPrimary().API.Import(context.Background(), nil, importRequest) + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + err := c.GetPrimary().API.Import(context.Background(), qcx, importRequest) if err != nil { t.Fatalf("importing IDKey data: %v", err) } diff --git a/time.go b/time.go index b87774f1a..a17d37726 100644 --- a/time.go +++ b/time.go @@ -89,15 +89,69 @@ func viewByTimeUnit(name string, t time.Time, unit rune) string { } } +// YYYYMMDDHH lengths. Note that this is a []int, not a map[byte]int, so +// the lookups can be cheaper. +var lengthsByQuantum = []int{ + 'Y': 4, + 'M': 6, + 'D': 8, + 'H': 10, +} + +// viewsByTimeInto computes the list of views for a given time. It expects +// to be given an initial buffer of the form `name_YYYYMMDDHH`, and a slice +// of []bytes. This allows us to reuse the buffer for all the sub-buffers, +// and also to reuse the slice of slices, to eliminate all those allocations. +// This might seem crazy, but even including the JSON parsing and all the +// disk activity, the straightforward viewsByTime implementation was 25% +// of runtime in an ingest test. +func viewsByTimeInto(fullBuf []byte, into [][]byte, t time.Time, q TimeQuantum) [][]byte { + l := len(fullBuf) - 10 + date := fullBuf[l : l+10] + y, m, d := t.Date() + h := t.Hour() + // Did you know that Sprintf, Printf, and other things like that all + // do allocations, and that doing allocations in a tight loop like this + // is stunningly expensive? viewsByTime was 25% of an ingest test's + // total CPU, not counting the garbage collector overhead. This is about + // 3%. No, I'm not totally sure that justifies it. + if y < 1000 { + ys := fmt.Sprintf("%04d", y) + copy(date[0:4], []byte(ys)) + } else if y >= 10000 { + // This is probably a bad answer but there isn't really a + // good answer. + ys := fmt.Sprintf("%04d", y%1000) + copy(date[0:4], []byte(ys)) + } else { + strconv.AppendInt(date[:0], int64(y), 10) + } + date[4] = '0' + byte(m/10) + date[5] = '0' + byte(m%10) + date[6] = '0' + byte(d/10) + date[7] = '0' + byte(d%10) + date[8] = '0' + byte(h/10) + date[9] = '0' + byte(h%10) + into = into[:0] + for _, unit := range q { + if int(unit) < len(lengthsByQuantum) && lengthsByQuantum[unit] != 0 { + into = append(into, fullBuf[:l+lengthsByQuantum[unit]]) + } + } + return into +} + // viewsByTime returns a list of views for a given timestamp. func viewsByTime(name string, t time.Time, q TimeQuantum) []string { // nolint: unparam + y, m, d := t.Date() + h := t.Hour() + full := fmt.Sprintf("%s_%04d%02d%02d%02d", name, y, m, d, h) + l := len(name) + 1 a := make([]string, 0, len(q)) for _, unit := range q { - view := viewByTimeUnit(name, t, unit) - if view == "" { - continue + if int(unit) < len(lengthsByQuantum) && lengthsByQuantum[unit] != 0 { + a = append(a, full[:l+lengthsByQuantum[unit]]) } - a = append(a, view) } return a } @@ -257,16 +311,16 @@ func parsePartialTime(t string) (time.Time, error) { // has minutes minute, err = strconv.Atoi(subStrings[1]) if err != nil { - return -1, -1, errors.New("Invalid Time") + return -1, -1, errors.New("invalid time") } fallthrough case 1: hour, err = strconv.Atoi(subStrings[0]) if err != nil { - return -1, -1, errors.New("Invalid Time") + return -1, -1, errors.New("invalid time") } default: - return -1, -1, errors.New("Invalid Time") + return -1, -1, errors.New("invalid time") } return @@ -282,7 +336,7 @@ func parsePartialTime(t string) (time.Time, error) { return true } if len(subMatches) <= 1 { - return nil, errors.New("Invalid time") + return nil, errors.New("invalid time") } // ignore full match which is at index 0 subMatches = subMatches[1:] // ignore full match which is at index 0 @@ -295,7 +349,7 @@ func parsePartialTime(t string) (time.Time, error) { } else { // rest must be empty for date-time to be valid if !restAreEmpty(subMatches[i:]) { - return nil, errors.New("Invalid date-time") + return nil, errors.New("invalid date-time") } break } diff --git a/time_internal_test.go b/time_internal_test.go index 8356d5dbe..2627ec65d 100644 --- a/time_internal_test.go +++ b/time_internal_test.go @@ -83,6 +83,38 @@ func TestViewsByTime(t *testing.T) { }) } +func TestViewsByTimeInto(t *testing.T) { + ts := time.Date(2000, time.January, 2, 3, 4, 5, 6, time.UTC) + s := []byte("F_YYYYMMDDHH") + var timeViews [][]byte + + t.Run("YMDH", func(t *testing.T) { + a := viewsByTime("F", ts, mustParseTimeQuantum("YMDH")) + b := viewsByTimeInto(s, timeViews, ts, mustParseTimeQuantum("YMDH")) + if len(a) != len(b) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + for i := range a { + if a[i] != string(b[i]) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + } + }) + + t.Run("D", func(t *testing.T) { + a := viewsByTime("F", ts, mustParseTimeQuantum("D")) + b := viewsByTimeInto(s, timeViews, ts, mustParseTimeQuantum("D")) + if len(a) != len(b) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + for i := range a { + if a[i] != string(b[i]) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + } + }) +} + // Ensure sets of fields can be returned for a given time range. func TestViewsByTimeRange(t *testing.T) { t.Run("Y", func(t *testing.T) { diff --git a/tx_test.go b/tx_test.go index e53a91dbe..e383a86b1 100644 --- a/tx_test.go +++ b/tx_test.go @@ -20,7 +20,7 @@ import ( "strings" "testing" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/http" "github.com/molecula/featurebase/v2/server" "github.com/molecula/featurebase/v2/storage" @@ -164,7 +164,12 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { //vv("BEFORE the first ImportAtomicRecord!") - if err := m0api.ImportAtomicRecord(ctx, nil, air); err != nil { + qcx := m0api.Txf().NewQcx() + if err := m0api.ImportAtomicRecord(ctx, qcx, air); err != nil { + qcx.Abort() + t.Fatal(err) + } + if err := qcx.Finish(); err != nil { t.Fatal(err) } @@ -196,9 +201,9 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { air = createAIRUpdate(expectedBalEndingAcct0, expectedBalEndingAcct1) - qcx := m0api.Txf().NewQcx() + qcx = m0api.Txf().NewQcx() //vv("just before the SECOND ImportAtomicRecord, qcx is %p, should NOT BE NIL", qcx) - err = m0api.ImportAtomicRecord(ctx, qcx, air, opt) + err = m0api.ImportAtomicRecord(ctx, qcx, air.Clone(), opt) //err = m0api.ImportAtomicRecord(ctx, nil, air, opt) if err != pilosa.ErrAborted { PanicOn(fmt.Sprintf("expected ErrTxnAborted but got err='%#v'", err)) @@ -223,9 +228,12 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { // happy path with no power failure half-way through. - err = m0api.ImportAtomicRecord(ctx, nil, air) + qcx = m0api.Txf().NewQcx() + err = m0api.ImportAtomicRecord(ctx, qcx, air.Clone()) PanicOn(err) - + if err := qcx.Finish(); err != nil { + t.Fatal(err) + } eb0, eb1 := queryBalances(m0api, acctOwnerID, fieldAcct0, fieldAcct1, index) // should have been applied this time. @@ -240,7 +248,11 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { air.Ivr[1].Clear = true air.Ir[0].Clear = true - err = m0api.ImportAtomicRecord(ctx, nil, air) + qcx = m0api.Txf().NewQcx() + err = m0api.ImportAtomicRecord(ctx, qcx, air) + if err := qcx.Finish(); err != nil { + t.Fatal(err) + } PanicOn(err) eb0, eb1 = queryBalances(m0api, acctOwnerID, fieldAcct0, fieldAcct1, index) diff --git a/txfactory.go b/txfactory.go index 84742cd2e..93e9449df 100644 --- a/txfactory.go +++ b/txfactory.go @@ -1383,7 +1383,7 @@ func (f *TxFactory) greenHasData() (hasData bool, err error) { case 2: return f.dbPerShard.HasData(1) } - err = fmt.Errorf("unsupported len(f.types): %v. Must be 1 or 2.", n) + err = fmt.Errorf("unsupported len(f.types): %v; must be 1 or 2", n) PanicOn(err) return } @@ -1410,7 +1410,7 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) { greenSrc := f.types[1] if blueDest == roaringTxn { - return fmt.Errorf("error: cannot migrate to 'roaring': not implemented.") + return fmt.Errorf("error: cannot migrate to 'roaring': not implemented") } idxs := holder.Indexes() @@ -1427,13 +1427,13 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) { return errors.Wrap(err, "TxFactory.green2blue f.greenHasData()") } if !blueHasData && !greenHasData { - holder.Logger.Infof("no data in blue or green. No migration or verification to do.") + holder.Logger.Infof("no data in blue or green. No migration or verification to do") return nil } // INVAR: blue has data. if !greenHasData { - holder.Logger.Errorf("cannot migrate from green '%v' because it has no data in it.", greenSrc) - return fmt.Errorf("error: cannot migrate from green '%v' because it has no data in it.", greenSrc) + holder.Logger.Errorf("cannot migrate from green '%v' because it has no data in it", greenSrc) + return fmt.Errorf("error: cannot migrate from green '%v' because it has no data in it", greenSrc) } nGoro := runtime.NumCPU()