Merge pull request #1676 from seebs/prep

ingest API prep work
This commit is contained in:
Matthew Jaffee 2021-08-19 09:08:28 -05:00 • committed by GitHub
commit 61f3fa23b0
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
20 changed files with 622 additions and 198 deletions

View file

@ -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)

76
api.go
View file

@ -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

View file

@ -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())

View file

@ -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)

View file

@ -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++
}

View file

@ -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,

152
field.go
View file

@ -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)

View file

@ -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())

View file

@ -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)

View file

@ -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) ||

View file

@ -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"

View file

@ -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

2
rbf.go
View file

@ -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)

70
shardwidth/helper.go Normal file
View file

@ -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
}

103
shardwidth/helper_test.go Normal file
View file

@ -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])
}
}
}
}

View file

@ -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)
}

72
time.go
View file

@ -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
}

View file

@ -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) {

View file

@ -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)

View file

@ -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()