From 33a916c0f48a0bf8f9d03425b536cbad71a4a867 Mon Sep 17 00:00:00 2001 From: Samir Patel Date: Sat, 16 Jul 2022 20:27:20 -0500 Subject: [PATCH] Revert "increase time range for timestamp by using specified granularity" This PR reverts the timestamp work. The timestamp work requires changes to FB and IDK; and there are circular dependencies between tests in either repo preventing merging of either. The work here is pretty stable, but required bypassing the smoketest. Meanwhile, I found some additional things in IDK that need addressing which means I merged this work in pre-maturely. Once I get that worked out, I'll re-commit these commits. This reverts following commits related to timestamp work: bypass of smoke test b/c of circular dep with IDK: 0676790 update codec to reflect changes to timstamp range: bd5dc76 fix few bugs regarding timestamp: 33fce8a increase time range for timestamp by using specified granularity: 5939923. --- .gitlab/.gitlab-ci.yml | 2 - api.go | 3 +- executor.go | 117 +------- executor_test.go | 625 +---------------------------------------- field.go | 147 ++++------ field_test.go | 12 +- fragment.go | 2 - ingest/codec.go | 54 +--- ingest/codec_test.go | 8 +- server/pg.go | 8 +- util.go | 6 + util_test.go | 18 ++ 12 files changed, 120 insertions(+), 882 deletions(-) diff --git a/.gitlab/.gitlab-ci.yml b/.gitlab/.gitlab-ci.yml index 23dedc822..04f4b2850 100644 --- a/.gitlab/.gitlab-ci.yml +++ b/.gitlab/.gitlab-ci.yml @@ -518,7 +518,6 @@ smoke test auth: - report.xml reports: junit: report.xml - allow_failure: true smoke test: stage: integration @@ -577,7 +576,6 @@ smoke test: - report.xml reports: junit: report.xml - allow_failure: true tremor-delete-test: stage: integration diff --git a/api.go b/api.go index 1fca0fa99..66fdfb8f5 100644 --- a/api.go +++ b/api.go @@ -2006,7 +2006,8 @@ func (api *API) IngestOperations(ctx context.Context, qcx *Qcx, indexName string return fmt.Errorf("adding decimal field to codec: %w", err) } case "timestamp": - if err = codec.AddTimestampField(field.name, field.options.TimeUnit, field.options.Base); err != nil { + nanos := TimeUnitNanos(field.options.TimeUnit) + if err = codec.AddTimestampField(field.name, time.Duration(nanos), field.options.Base); err != nil { return fmt.Errorf("adding timestamp field to codec: %w", err) } default: diff --git a/executor.go b/executor.go index a407a5e0d..de2bf1f0d 100644 --- a/executor.go +++ b/executor.go @@ -939,11 +939,7 @@ func (e *executor) executeFieldValueCallShard(ctx context.Context, qcx *Qcx, fie other.FloatVal = 0 other.Val = 0 } else if field.Type() == FieldTypeTimestamp { - ts, err := ValToTimestamp(field.Options().TimeUnit, value) - if err != nil { - return ValCount{}, err - } - other.TimestampVal = ts + other.TimestampVal = time.Unix(0, value*int64(TimeUnitNanos(field.Options().TimeUnit))) } return other, nil @@ -1580,11 +1576,7 @@ func (e *executor) executeDistinctShard(ctx context.Context, qcx *Qcx, index str cols := r.Pos.Columns() results := make([]string, len(cols)) for i, val := range cols { - t, err := ValToTimestamp(field.options.TimeUnit, int64(val)+bsig.Base) - if err != nil { - return nil, errors.Wrap(err, "translating value to timestamp") - } - results[i] = t.Format(time.RFC3339Nano) + results[i] = FormatTimestampNano(int64(val), bsig.Base, field.options.TimeUnit) } result = DistinctTimestamp{Name: fieldName, Values: results} return result, nil @@ -1625,11 +1617,6 @@ func (d DistinctTimestamp) ToRows(callback func(*proto.RowResponse) error) error return nil } -// ToTable implements the ToTabler interface for DistinctTimestamp -func (d DistinctTimestamp) ToTable() (*proto.TableResponse, error) { - return proto.RowsToTable(&d, len(d.Values)) -} - // Union returns the union of the values of `d` and `other` func (d *DistinctTimestamp) Union(other DistinctTimestamp) DistinctTimestamp { both := map[string]struct{}{} @@ -3197,16 +3184,13 @@ func (fr *FieldRow) Clone() (clone *FieldRow) { func (fr FieldRow) MarshalJSON() ([]byte, error) { if fr.Value != nil { if fr.FieldOptions.Type == FieldTypeTimestamp { - ts, err := ValToTimestamp(fr.FieldOptions.TimeUnit, int64(*fr.Value)+fr.FieldOptions.Base) - if err != nil { - return nil, errors.Wrap(err, "translating value to timestamp") - } + ts := FormatTimestampNano(int64(*fr.Value), fr.FieldOptions.Base, fr.FieldOptions.TimeUnit) return json.Marshal(struct { Field string `json:"field"` Value string `json:"value"` }{ Field: fr.Field, - Value: ts.Format(time.RFC3339Nano), + Value: ts, }) } else { return json.Marshal(struct { @@ -7587,17 +7571,12 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index } case FieldTypeTimestamp: datatype = "timestamp" - unit := field.Options().TimeUnit mapper = func(ids []uint64) (_ interface{}, err error) { switch len(ids) { case 0: return nil, nil case 1: - ts, err := ValToTimestamp(unit, int64(ids[0])) - if err != nil { - return nil, err - } - return ts, nil + return time.Unix(0, int64(ids[0])*int64(TimeUnitNanos(field.Options().TimeUnit))).UTC(), nil default: return nil, errors.Errorf("BSI field %q has too many values: %v", field.Name(), ids) } @@ -7658,37 +7637,6 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index return result, nil } -// ValToTimestamp takes a timeunit and an integer value and converts it to time.Time -func ValToTimestamp(unit string, val int64) (time.Time, error) { - switch unit { - case TimeUnitSeconds: - return time.Unix(val, 0).UTC(), nil - case TimeUnitMilliseconds: - return time.UnixMilli(val).UTC(), nil - case TimeUnitMicroseconds, TimeUnitUSeconds: - return time.UnixMicro(val).UTC(), nil - case TimeUnitNanoseconds: - return time.Unix(0, val).UTC(), nil - default: - return time.Time{}, errors.Errorf("Unknown time unit: '%v'", unit) - } -} - -// TimestampToVal takes a time unit and a time.Time and converts it to an integer value -func TimestampToVal(unit string, ts time.Time) int64 { - switch unit { - case TimeUnitSeconds: - return ts.Unix() - case TimeUnitMilliseconds: - return ts.UnixMilli() - case TimeUnitMicroseconds, TimeUnitUSeconds: - return ts.UnixMicro() - case TimeUnitNanoseconds: - return ts.UnixNano() - } - return 0 -} - // detectRangeCall returns true if the call or one of its children contains a Range call // TODO: Remove at version 2.0 func (e *executor) detectRangeCall(c *pql.Call) bool { @@ -7999,8 +7947,6 @@ func (vc *ValCount) smaller(other ValCount) ValCount { return vc.decimalSmaller(other) } else if vc.FloatVal != 0 || other.FloatVal != 0 { return vc.floatSmaller(other) - } else if !vc.TimestampVal.IsZero() || !other.TimestampVal.IsZero() { - return vc.timestampSmaller(other) } if vc.Count == 0 || (other.Val < vc.Val && other.Count > 0) { return other @@ -8018,28 +7964,6 @@ func (vc *ValCount) smaller(other ValCount) ValCount { } } -// timestampSmaller returns the smaller of the two (vc or other), while merging the count -// if they are equal. -func (vc *ValCount) timestampSmaller(other ValCount) ValCount { - if other.TimestampVal.Equal(time.Time{}) { - return *vc - } - if vc.Count == 0 || vc.TimestampVal.Equal(time.Time{}) || (other.TimestampVal.Before(vc.TimestampVal) && other.Count > 0) { - return other - } - extra := int64(0) - if vc.TimestampVal.Equal(other.TimestampVal) { - extra += other.Count - } - return ValCount{ - Val: vc.Val, - TimestampVal: vc.TimestampVal, - Count: vc.Count + extra, - } -} - -// decimalSmaller returns the smaller of the two (vc or other), while merging the count -// if they are equal. func (vc *ValCount) decimalSmaller(other ValCount) ValCount { if other.DecimalVal == nil { return *vc @@ -8057,8 +7981,6 @@ func (vc *ValCount) decimalSmaller(other ValCount) ValCount { } } -// floatSmaller returns the smaller of the two (vc or other), while merging the count -// if they are equal. func (vc *ValCount) floatSmaller(other ValCount) ValCount { if vc.Count == 0 || (other.FloatVal < vc.FloatVal && other.Count > 0) { return other @@ -8079,8 +8001,6 @@ func (vc *ValCount) larger(other ValCount) ValCount { return vc.decimalLarger(other) } else if vc.FloatVal != 0 || other.FloatVal != 0 { return vc.floatLarger(other) - } else if !vc.TimestampVal.Equal(time.Time{}) || !other.TimestampVal.Equal(time.Time{}) { - return vc.timestampLarger(other) } if vc.Count == 0 || (other.Val > vc.Val && other.Count > 0) { return other @@ -8098,28 +8018,6 @@ func (vc *ValCount) larger(other ValCount) ValCount { } } -// timestampLarger returns the larger of the two (vc or other), while merging the count -// if they are equal. -func (vc *ValCount) timestampLarger(other ValCount) ValCount { - if other.TimestampVal.Equal(time.Time{}) { - return *vc - } - if vc.Count == 0 || vc.TimestampVal.Equal(time.Time{}) || (other.TimestampVal.After(vc.TimestampVal) && other.Count > 0) { - return other - } - extra := int64(0) - if vc.TimestampVal.Equal(other.TimestampVal) { - extra += other.Count - } - return ValCount{ - Val: vc.Val, - TimestampVal: vc.TimestampVal, - Count: vc.Count + extra, - } -} - -// decimalLarger returns the larger of the two (vc or other), while merging the count -// if they are equal. func (vc *ValCount) decimalLarger(other ValCount) ValCount { if other.DecimalVal == nil { return *vc @@ -8137,8 +8035,6 @@ func (vc *ValCount) decimalLarger(other ValCount) ValCount { } } -// floatLarger returns the larger of the two (vc or other), while merging the count -// if they are equal. func (vc *ValCount) floatLarger(other ValCount) ValCount { if vc.Count == 0 || (other.FloatVal > vc.FloatVal && other.Count > 0) { return other @@ -8556,8 +8452,7 @@ func getScaledInt(f *Field, v interface{}) (int64, error) { } else if opt.Type == FieldTypeTimestamp { switch tv := v.(type) { case time.Time: - v := TimestampToVal(f.options.TimeUnit, tv) - value = v + value = tv.UnixNano() / TimeUnitNanos(f.options.TimeUnit) case int64: value = tv default: diff --git a/executor_test.go b/executor_test.go index bce250b41..1ad1e9eb4 100644 --- a/executor_test.go +++ b/executor_test.go @@ -7242,8 +7242,6 @@ func TestVariousQueries(t *testing.T) { variousQueriesOnPercentiles(t, c) variousQueriesCountDistinctTimestamp(t, c) variousQueriesOnIntFields(t, c) - variousQueriesOnTimestampFields(t, c) - variousQueriesOnLargeEpoch(t, c) backupTest(t, c, "") // test backup/restore of all indexes }) } @@ -7494,10 +7492,6 @@ func variousQueriesOnPercentiles(t *testing.T, c *test.Cluster) { } } -func lineBreaker(s string) string { - return strings.Join(strings.Split(s, " "), "\n") + "\n" -} - // tests for abbreviating time values in queries func variousQueriesOnTimeFields(t *testing.T, c *test.Cluster) { ts := func(t time.Time) int64 { @@ -7544,6 +7538,10 @@ func variousQueriesOnTimeFields(t *testing.T, c *test.Cluster) { return strings.Join(ss, "\n") + "\n" } + toCSV := func(s string) string { + return strings.Join(strings.Split(s, " "), "\n") + "\n" + } + type testCase struct { query string qrVerifier func(t *testing.T, resp pilosa.QueryResponse) @@ -7554,44 +7552,44 @@ func variousQueriesOnTimeFields(t *testing.T, c *test.Cluster) { // Rows { query: `Rows(f1, from='2019-08-04T14:36', to='2019-08-04T16:00')`, - csvVerifier: lineBreaker("R4 R5"), + csvVerifier: toCSV("R4 R5"), }, { query: `Rows(f1, from='2019-08-04T14', to='2019-08-04T17:00')`, - csvVerifier: lineBreaker("R4 R5 R6"), + csvVerifier: toCSV("R4 R5 R6"), }, { query: `Rows(f1, from='2019-08-04', to='2019-08-05')`, - csvVerifier: lineBreaker("R3 R4 R5 R6"), + csvVerifier: toCSV("R3 R4 R5 R6"), }, { query: `Rows(f1, from='2019-08', to='2019-12')`, - csvVerifier: lineBreaker("R2 R3 R4 R5 R6 R7"), + csvVerifier: toCSV("R2 R3 R4 R5 R6 R7"), }, { query: `Rows(f1, from='2019', to='2020')`, - csvVerifier: lineBreaker("R1 R2 R3 R4 R5 R6 R7 R8"), + csvVerifier: toCSV("R1 R2 R3 R4 R5 R6 R7 R8"), }, // Row { query: `Row(f2='R', from='2019-08-04T14:36', to='2019-08-04T16:00')`, - csvVerifier: lineBreaker("C4 C5"), + csvVerifier: toCSV("C4 C5"), }, { query: `Row(f2='R', from='2019-08-04T14', to='2019-08-04T17:00')`, - csvVerifier: lineBreaker("C4 C5 C6"), + csvVerifier: toCSV("C4 C5 C6"), }, { query: `Row(f2='R', from='2019-08-04', to='2019-08-05')`, - csvVerifier: lineBreaker("C3 C4 C5 C6"), + csvVerifier: toCSV("C3 C4 C5 C6"), }, { query: `Row(f2='R', from='2019-08', to='2019-12')`, - csvVerifier: lineBreaker("C2 C3 C4 C5 C6 C7"), + csvVerifier: toCSV("C2 C3 C4 C5 C6 C7"), }, { query: `Row(f2='R', from='2019', to='2020')`, - csvVerifier: lineBreaker("C1 C2 C3 C4 C5 C6 C7 C8"), + csvVerifier: toCSV("C1 C2 C3 C4 C5 C6 C7 C8"), }, } @@ -7704,598 +7702,6 @@ userG,-1,10,10,10 } } -// Constants used in TimestampField testing -var ( - minTime = pilosa.MinTimestamp - maxTime = pilosa.MaxTimestamp - minSec = minTime.Unix() - maxSec = maxTime.Unix() - minMilli = minTime.UnixMilli() - maxMilli = maxTime.UnixMilli() - minMicro = minTime.UnixMicro() - maxMicro = maxTime.UnixMicro() - minNano = pilosa.MinTimestampNano.UnixNano() - maxNano = pilosa.MaxTimestampNano.UnixNano() -) - -// variousQueriesOnTimestampFields tests queries on Timestamp Fields at various granularities using the default epoch -func variousQueriesOnTimestampFields(t *testing.T, c *test.Cluster) { - index := "ts_test01" - - // Testing whether the max and min timestamps can be represented for seconds - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_sec", pilosa.OptFieldTypeTimestamp(time.Unix(0, 0), "s")) - c.ImportIntKey(t, index, "unix_sec", []test.IntKey{ - {Val: minSec, Key: "userA"}, - {Val: maxSec, Key: "userB"}, - }) - - // Testing whether the max and min timestamps can be represented for milliseconds - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_milli", pilosa.OptFieldTypeTimestamp(time.Unix(0, 0), "ms")) - c.ImportIntKey(t, index, "unix_milli", []test.IntKey{ - {Val: minMilli, Key: "userA"}, - {Val: maxMilli, Key: "userB"}, - }) - - // Testing whether the max and min timestamps can be represented for microseconds - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_micro", pilosa.OptFieldTypeTimestamp(time.Unix(0, 0), "us")) - c.ImportIntKey(t, index, "unix_micro", []test.IntKey{ - {Val: minMicro, Key: "userA"}, - {Val: maxMicro, Key: "userB"}, - }) - - // Note that min and max values that can be represented for Nanos is a much smaller range than any of the above granularities. - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_nano", pilosa.OptFieldTypeTimestamp(time.Unix(0, 0), "ns")) - c.ImportIntKey(t, index, "unix_nano", []test.IntKey{ - {Val: minNano, Key: "userA"}, - {Val: maxNano, Key: "userB"}, - }) - - splitSortBackToCSV := func(csvStr string) string { - if len(csvStr) == 0 { - return "" - } - ss := strings.Split(csvStr[:len(csvStr)-1], "\n") - sort.Strings(ss) - return strings.Join(ss, "\n") + "\n" - } - - type testCase struct { - query string - qrVerifier func(t *testing.T, resp pilosa.QueryResponse) - csvVerifier string - } - - tests := []testCase{ - { - query: "extract(All(), Rows(unix_sec))", - csvVerifier: `userA,0001-01-01T00:00:01Z -userB,9999-12-31T23:59:59Z -`, - }, - { - query: "extract(All(), Rows(unix_milli))", - csvVerifier: `userA,0001-01-01T00:00:01Z -userB,9999-12-31T23:59:59Z -`, - }, - { - query: "extract(All(), Rows(unix_micro))", - csvVerifier: `userA,0001-01-01T00:00:01Z -userB,9999-12-31T23:59:59Z -`, - }, - { - query: "extract(All(), Rows(unix_nano))", - csvVerifier: `userA,1833-11-24T17:31:44Z -userB,2106-02-07T06:28:16Z -`, - }, - { - query: "All()", - csvVerifier: `userA -userB -`, - }, - { - query: "count(All())", - csvVerifier: `2 -`, - }, - // Found an existing bug: Disticnt on timestamp does not return values prior to the epoch. - // These tests need to be uncommented when issue is fixed. - // { - // query: "Distinct(field=unix_sec)", - // csvVerifier: `0001-01-01T00:00:01Z - // 9999-12-31T23:59:59Z - // `, - // }, - // { - // query: "Distinct(field=unix_milli)", - // csvVerifier: `0001-01-01T00:00:01Z - // 9999-12-31T23:59:59Z - // `, - // }, - { - query: "Min(unix_sec)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_sec)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_milli)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_milli)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_micro)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_micro)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_nano)", - csvVerifier: `1833-11-24T17:31:44Z,1 -`, - }, - { - query: "Max(unix_nano)", - csvVerifier: `2106-02-07T06:28:16Z,1 -`, - }, - { - query: "GroupBy(Rows(unix_micro))", - csvVerifier: `-62135596799000000,1 -253402300799000000,1 -`, - }, - { - query: `Row(unix_sec="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_sec="9999-12-31T23:59:59Z")`, - csvVerifier: lineBreaker("userB"), - }, - { - query: `Row(unix_milli="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_milli="9999-12-31T23:59:59Z")`, - csvVerifier: lineBreaker("userB"), - }, - { - query: `Row(unix_micro="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_micro="9999-12-31T23:59:59Z")`, - csvVerifier: lineBreaker("userB"), - }, - { - query: `Row(unix_nano="1833-11-24T17:31:44Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_nano="2106-02-07T06:28:16Z")`, - csvVerifier: lineBreaker("userB"), - }, - { - query: `Union(Row(unix_nano="2106-02-07T06:28:16Z"), Row(unix_micro="0001-01-01T00:00:01Z"))`, - csvVerifier: lineBreaker("userA\nuserB"), - }, - { - query: `Set("userA", unix_milli="2000-12-31T23:59:59.999Z")`, - csvVerifier: `true -`, - }, - { - query: `Clear("userA", unix_milli="2000-12-31T23:59:59.999Z")`, - csvVerifier: `true -`, - }, - // Not Supported for Timestamp - // { - // query: `ClearRow(unix_milli="2000-12-31T23:59:59.999Z")`, - // csvVerifier: `true`, - // }, - } - - for i, tst := range tests { - t.Run(fmt.Sprintf("%d-%s", i, tst.query), func(t *testing.T) { - tr := c.QueryGRPC(t, index, tst.query) - csvString, err := tableResponseToCSVString(tr) - if err != nil { - t.Fatal(err) - } - // verify everything after header - got := splitSortBackToCSV(csvString[strings.Index(csvString, "\n")+1:]) - if got != tst.csvVerifier { - t.Errorf("expected:\n%s\ngot:\n%s", tst.csvVerifier, got) - } - }) - } -} - -// variousQueriesOnLargeEpoch tests queries on TimestampFields when the epoch is set either to -// the min or max timestamp allowed. -func variousQueriesOnLargeEpoch(t *testing.T, c *test.Cluster) { - index := "ts_epoch_test20321" - // mag := int64(10000000000) - - // These large constants are close to min and max int64 but not quite since go has to account for the difference between Unix Epoch and Go's Epoch - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_sec_min", pilosa.OptFieldTypeTimestamp(minTime, "s")) - c.ImportIntKey(t, index, "unix_sec_min", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -minSec, Key: "userB"}, - {Val: -minSec + maxSec, Key: "userC"}, - }) - - // vprint.VV("-MaxSec: %+v", -maxSec) - // vprint.VV("-MaxSec + minsec: %+v", -maxSec+minSec) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_sec_max", pilosa.OptFieldTypeTimestamp(maxTime, "s")) - c.ImportIntKey(t, index, "unix_sec_max", []test.IntKey{ - {Val: 0, Key: "userA"}, // 9999-12-31 - // {Val: 1, Key: "userE"}, - // {Val: -1, Key: "userD"}, - {Val: -maxSec, Key: "userB"}, //1970.... - {Val: -maxSec + minSec, Key: "userC"}, //0001 - }) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_milli_min", pilosa.OptFieldTypeTimestamp(minTime, "ms")) - c.ImportIntKey(t, index, "unix_milli_min", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -minMilli, Key: "userB"}, - {Val: -minMilli + maxMilli, Key: "userC"}, - }) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_milli_max", pilosa.OptFieldTypeTimestamp(maxTime, "ms")) - c.ImportIntKey(t, index, "unix_milli_max", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -maxMilli, Key: "userB"}, - {Val: -maxMilli + minMilli, Key: "userC"}, - }) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_micro_min", pilosa.OptFieldTypeTimestamp(minTime, "us")) - c.ImportIntKey(t, index, "unix_micro_min", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -minMicro, Key: "userB"}, - {Val: -minMicro + maxMicro, Key: "userC"}, - }) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_micro_max", pilosa.OptFieldTypeTimestamp(maxTime, "us")) - c.ImportIntKey(t, index, "unix_micro_max", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -maxMicro, Key: "userB"}, - {Val: -maxMicro + minMicro, Key: "userC"}, - }) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_nano_min", pilosa.OptFieldTypeTimestamp(pilosa.MinTimestampNano, "ns")) - c.ImportIntKey(t, index, "unix_nano_min", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -minNano - 1, Key: "userB"}, - {Val: -minNano, Key: "userC"}, - }) - - c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_nano_max", pilosa.OptFieldTypeTimestamp(pilosa.MaxTimestampNano, "ns")) - c.ImportIntKey(t, index, "unix_nano_max", []test.IntKey{ - {Val: 0, Key: "userA"}, - {Val: -maxNano + 1, Key: "userB"}, - {Val: -maxNano, Key: "userC"}, - }) - - splitSortBackToCSV := func(csvStr string) string { - if csvStr == "" { - return "" - } - ss := strings.Split(csvStr[:len(csvStr)-1], "\n") - sort.Strings(ss) - return strings.Join(ss, "\n") + "\n" - } - - type testCase struct { - query string - qrVerifier func(t *testing.T, resp pilosa.QueryResponse) - csvVerifier string - } - - tests := []testCase{ - { - query: "extract(All(), Rows(unix_sec_min))", - csvVerifier: `userA,0001-01-01T00:00:01Z -userB,1970-01-01T00:00:00Z -userC,9999-12-31T23:59:59Z -`, - }, - { - query: "extract(All(), Rows(unix_sec_max))", - csvVerifier: `userA,9999-12-31T23:59:59Z -userB,1970-01-01T00:00:00Z -userC,0001-01-01T00:00:01Z -`, - }, - { - query: "extract(All(), Rows(unix_milli_min))", - csvVerifier: `userA,0001-01-01T00:00:01Z -userB,1970-01-01T00:00:00Z -userC,9999-12-31T23:59:59Z -`, - }, - { - query: "extract(All(), Rows(unix_milli_max))", - csvVerifier: `userA,9999-12-31T23:59:59Z -userB,1970-01-01T00:00:00Z -userC,0001-01-01T00:00:01Z -`, - }, - { - query: "extract(All(), Rows(unix_micro_min))", - csvVerifier: `userA,0001-01-01T00:00:01Z -userB,1970-01-01T00:00:00Z -userC,9999-12-31T23:59:59Z -`, - }, - { - query: "extract(All(), Rows(unix_micro_max))", - csvVerifier: `userA,9999-12-31T23:59:59Z -userB,1970-01-01T00:00:00Z -userC,0001-01-01T00:00:01Z -`, - }, - { - query: "extract(All(), Rows(unix_nano_min))", - csvVerifier: `userA,1833-11-24T17:31:44Z -userB,1969-12-31T23:59:59.999999999Z -userC,1970-01-01T00:00:00Z -`, - }, - { - query: "extract(All(), Rows(unix_nano_max))", - csvVerifier: `userA,2106-02-07T06:28:16Z -userB,1970-01-01T00:00:00.000000001Z -userC,1970-01-01T00:00:00Z -`, - }, - { - query: "All()", - csvVerifier: `userA -userB -userC -`, - }, - { - query: "count(All())", - csvVerifier: `3 -`, - }, - { - query: "Min(unix_sec_min)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_sec_min)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_sec_max)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_sec_max)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_milli_min)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_milli_min)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Max(unix_milli_max)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_micro_min)", - csvVerifier: `0001-01-01T00:00:01Z,1 -`, - }, - { - query: "Max(unix_micro_max)", - csvVerifier: `9999-12-31T23:59:59Z,1 -`, - }, - { - query: "Min(unix_nano_min)", - csvVerifier: `1833-11-24T17:31:44Z,1 -`, - }, - { - query: "Max(unix_nano_min)", - csvVerifier: `1970-01-01T00:00:00Z,1 -`, - }, - { - query: "Min(unix_nano_max)", - csvVerifier: `1970-01-01T00:00:00Z,1 -`, - }, - { - query: "Max(unix_nano_max)", - csvVerifier: `2106-02-07T06:28:16Z,1 -`, - }, - { - query: "GroupBy(Rows(unix_micro_min))", - csvVerifier: `-62135596799000000,1 -0,1 -253402300799000000,1 -`, - }, - { - query: `Row(unix_sec_min="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_sec_max="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userC"), - }, - { - query: `Row(unix_sec_max="9999-12-31T23:59:59Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_milli_min="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_milli_min="9999-12-31T23:59:59Z")`, - csvVerifier: lineBreaker("userC"), - }, - { - query: `Row(unix_micro_max="0001-01-01T00:00:01Z")`, - csvVerifier: lineBreaker("userC"), - }, - { - query: `Row(unix_micro_max="9999-12-31T23:59:59Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_nano_min="1833-11-24T17:31:44Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_nano_min="1969-12-31T23:59:59.999999999Z")`, - csvVerifier: lineBreaker("userB"), - }, - { - query: `Row(unix_nano_min="1970-01-01T00:00:00Z")`, - csvVerifier: lineBreaker("userC"), - }, - { - query: `Row(unix_nano_max="2106-02-07T06:28:16Z")`, - csvVerifier: lineBreaker("userA"), - }, - { - query: `Row(unix_nano_max="1970-01-01T00:00:00.000000001Z")`, - csvVerifier: lineBreaker("userB"), - }, - { - query: `Row(unix_nano_max="1970-01-01T00:00:00Z")`, - csvVerifier: lineBreaker("userC"), - }, - { - query: `Union(Row(unix_nano_max="2106-02-07T06:28:16Z"), Row(unix_micro_max="0001-01-01T00:00:01Z"))`, - csvVerifier: lineBreaker("userA\nuserC"), - }, - { - query: `Set("userA", unix_sec_min="2000-12-31T23:59:59.999Z")`, - csvVerifier: `true -`, - }, - { - query: "extract(All(), Rows(unix_sec_min))", - csvVerifier: `userA,2000-12-31T23:59:59Z -userB,1970-01-01T00:00:00Z -userC,9999-12-31T23:59:59Z -`, - }, - { - query: `Clear("userA", unix_sec_min="2000-12-31T23:59:59.999Z")`, - csvVerifier: `true -`, - }, - { - query: `Set("userA", unix_milli_max="2000-12-31T23:59:59.999Z")`, - csvVerifier: `true -`, - }, - { - query: "extract(All(), Rows(unix_milli_max))", - csvVerifier: `userA,2000-12-31T23:59:59.999Z -userB,1970-01-01T00:00:00Z -userC,0001-01-01T00:00:01Z -`, - }, - { - query: `Clear("userA", unix_milli_max="2000-12-31T23:59:59.999Z")`, - csvVerifier: `true -`, - }, - { - query: `Set("userA", unix_nano_max="2050-02-01T00:00:00.000000002Z")`, - csvVerifier: `true -`, - }, - { - query: "extract(All(), Rows(unix_nano_max))", - csvVerifier: `userA,2050-02-01T00:00:00.000000002Z -userB,1970-01-01T00:00:00.000000001Z -userC,1970-01-01T00:00:00Z -`, - }, - { - query: `Clear("userA", unix_nano_max="1970-01-01T00:00:00.000000002Z")`, - csvVerifier: `true -`, - }, - { - query: `Set("userA", unix_nano_min="1969-12-31T23:59:59.999999998Z")`, - csvVerifier: `true -`, - }, - { - query: "extract(All(), Rows(unix_nano_min))", - csvVerifier: `userA,1969-12-31T23:59:59.999999998Z -userB,1969-12-31T23:59:59.999999999Z -userC,1970-01-01T00:00:00Z -`, - }, - { - query: `Clear("userA", unix_nano_min="1969-12-31T23:59:59.999999998Z")`, - csvVerifier: `true -`, - }, - } - - for i, tst := range tests { - t.Run(fmt.Sprintf("%d-%s", i, tst.query), func(t *testing.T) { - tr := c.QueryGRPC(t, index, tst.query) - csvString, err := tableResponseToCSVString(tr) - if err != nil { - t.Fatal(err) - } - // verify everything after header - got := splitSortBackToCSV(csvString[strings.Index(csvString, "\n")+1:]) - if got != tst.csvVerifier { - t.Errorf("expected:\n%s\ngot:\n%s", tst.csvVerifier, got) - } - }) - } -} - var usersIndex = "users" func populateTestData(t *testing.T, c *test.Cluster) { @@ -8921,9 +8327,6 @@ func tableResponseToCSV(m *proto.TableResponse, w io.Writer) error { record = append(record, fmt.Sprintf("%v", col.GetBoolVal())) case "int64": record = append(record, fmt.Sprintf("%v", col.GetInt64Val())) - case "timestamp": - record = append(record, fmt.Sprintf("%v", col.GetTimestampVal())) - } } err := writer.Write(record) diff --git a/field.go b/field.go index 3fc7e18e2..d2bd17223 100644 --- a/field.go +++ b/field.go @@ -195,59 +195,26 @@ func OptFieldTypeInt(min, max int64) FieldOption { // provide any respective configuration values. func OptFieldTypeTimestamp(epoch time.Time, timeUnit string) FieldOption { return func(fo *FieldOptions) error { + // Check if the epoch will overflow when converted to nano. + if err := CheckUnixNanoOverflow(epoch); err != nil { + return err + } + epochValue := epoch.UnixNano() / TimeUnitNanos(timeUnit) if fo.Type != "" { return errors.Errorf("field type is already set to: %s", fo.Type) } - - minTime := MinTimestamp - maxTime := MaxTimestamp - - var base, minInt, maxInt int64 - switch timeUnit { - case TimeUnitSeconds: - base = epoch.Unix() - minInt = minTime.Unix() - base - maxInt = maxTime.Unix() - base - case TimeUnitMilliseconds: - base = epoch.UnixMilli() - minInt = minTime.UnixMilli() - base - maxInt = maxTime.UnixMilli() - base - case TimeUnitMicroseconds, TimeUnitUSeconds: - base = epoch.UnixMicro() - minInt = minTime.UnixMicro() - base - maxInt = maxTime.UnixMicro() - base - case TimeUnitNanoseconds: - // Note: For nano, the min and max values are also the min and max integer - // values we support. Also, keep in mind that MinNano is a negative - // number. So if base is positive and we do MinNano - base...it would increase minInt - // beyond what we support. This isn't an issue with larger granularities. - base = epoch.UnixNano() - if base > 0 { - maxInt = MaxTimestampNano.UnixNano() - base - minInt = MinTimestampNano.UnixNano() - } else { - maxInt = MaxTimestampNano.UnixNano() - minInt = MinTimestampNano.UnixNano() - base - } - minTime = MinTimestampNano - maxTime = MaxTimestampNano - default: - return errors.Errorf("invalid time unit: '%q'", fo.TimeUnit) + if timeUnit == "" { + return errors.Errorf("time unit required for timestamp field") + } else if !IsValidTimeUnit(timeUnit) { + return errors.Errorf("invalid time unit: %q", fo.TimeUnit) } - - if err := CheckEpochOutOfRange(epoch, minTime, maxTime); err != nil { - return err - } - fo.Type = FieldTypeTimestamp fo.TimeUnit = timeUnit - fo.Base = base - fo.Min = pql.NewDecimal(minInt, 0) - fo.Max = pql.NewDecimal(maxInt, 0) - + fo.Min = pql.NewDecimal(MinTimestamp.UnixNano()/TimeUnitNanos(timeUnit), 0) + fo.Max = pql.NewDecimal(MaxTimestamp.UnixNano()/TimeUnitNanos(timeUnit), 0) + fo.Base = epochValue return nil } - } // OptFieldTypeDecimal is a functional option for creating a `decimal` field. @@ -1435,20 +1402,15 @@ func (f *Field) SetValue(tx Tx, columnID uint64, value int64) (changed bool, err bsig := f.bsiGroup(f.name) if bsig == nil { return false, ErrBSIGroupNotFound - } - - // Determine base value to store. - baseValue := int64(value - bsig.Base) - //Timestamp expects incoming value to already be relative to epoch - if f.Type() == FieldTypeTimestamp { - value = baseValue - } - if value < bsig.Min { + } else if value < bsig.Min { return false, errors.Wrapf(ErrBSIGroupValueTooLow, "index = %v, field = %v, column ID = %v, value %v is smaller than min allowed %v", f.index, f.name, columnID, value, bsig.Min) } else if value > bsig.Max { return false, errors.Wrapf(ErrBSIGroupValueTooHigh, "index = %v, field = %v, column ID = %v, value %v is larger than max allowed %v", f.index, f.name, columnID, value, bsig.Max) } + // Determine base value to store. + baseValue := int64(value - bsig.Base) + requiredBitDepth := bitDepthInt64(baseValue) // Increase bit depth value if the unsigned value is greater. @@ -1589,14 +1551,8 @@ func (f *Field) valCountize(val int64, cnt uint64, bsig *bsiGroup) (ValCount, er dec := pql.NewDecimal(val+bsig.Base, bsig.Scale) valCount.DecimalVal = &dec } else if f.Options().Type == FieldTypeTimestamp { - ts, err := ValToTimestamp(f.options.TimeUnit, val+bsig.Base) - if err != nil { - return ValCount{}, errors.Wrap(err, "translating value to timestamp") - } - valCount.TimestampVal = ts - // valCount.TimestampVal = time.Unix(0, (val+bsig.Base)*TimeUnitNanos(f.options.TimeUnit)).UTC() + valCount.TimestampVal = time.Unix(0, (val+bsig.Base)*TimeUnitNanos(f.options.TimeUnit)).UTC() } - valCount.Val = val + bsig.Base return valCount, nil } @@ -1783,7 +1739,7 @@ func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time } for i, t := range values { - ivalues[i] = TimestampToVal(f.options.TimeUnit, t) + ivalues[i] = t.UnixNano() / TimeUnitNanos(f.options.TimeUnit) } return f.importValue(qcx, columnIDs, ivalues, shard, options) } @@ -1826,23 +1782,23 @@ func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, shard if value < min { min = value } - - } - - // Timestamps differ from other BSI fields in that integer representations - // of timestamps are already relative to the epoch (base). - // So a user may set an epoch to 2022-03-01 as the start of a race - // and import finishing times in seconds. - // Timestamps ingested as timestamps are of course absolute, but by the time - // we get here it would be a relative integer. - if f.Type() != FieldTypeTimestamp { - min -= bsig.Base - max -= bsig.Base + if f.Type() == FieldTypeTimestamp { + scale := (TimeUnitNanos(f.options.TimeUnit)) + offset := f.options.Base * scale + dur := value * scale + if offset > 0 { + if dur > math.MaxInt64-offset { + return errors.Wrap(ErrBSIGroupValueTooHigh, "value + epoch is too far from Unix epoch") + } + } else if dur < math.MinInt64-offset { + return errors.Wrap(ErrBSIGroupValueTooLow, "value + epoch is too far from Unix epoch") + } + } } // Determine the highest bit depth required by the min & max. - requiredDepth := bitDepthInt64(min) - if v := bitDepthInt64(max); v > requiredDepth { + requiredDepth := bitDepthInt64(min - bsig.Base) + if v := bitDepthInt64(max - bsig.Base); v > requiredDepth { requiredDepth = v } // Increase bit depth if required. @@ -2108,11 +2064,6 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) { o.Keys, }) case FieldTypeTimestamp: - epoch, err := ValToTimestamp(o.TimeUnit, o.Base) - if err != nil { - return nil, errors.Wrap(err, "translating val to timestamp") - } - return json.Marshal(struct { Type string `json:"type"` Epoch time.Time `json:"epoch"` @@ -2122,7 +2073,7 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) { TimeUnit string `json:"timeUnit"` }{ o.Type, - epoch, + time.Unix(0, o.Base*TimeUnitNanos(o.TimeUnit)).UTC(), o.BitDepth, o.Min, o.Max, @@ -2164,6 +2115,16 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) { return nil, errors.Errorf("invalid field type: '%s'", o.Type) } +// MinTimestamp returns the minimum value for a timestamp field. +func (o FieldOptions) MinTimestamp() time.Time { + return time.Unix(0, o.Min.ToInt64(0)*int64(TimeUnitNanos(o.TimeUnit))) +} + +// MaxTimestamp returns the maxnimum value for a timestamp field. +func (o FieldOptions) MaxTimestamp() time.Time { + return time.Unix(0, o.Max.ToInt64(0)*int64(TimeUnitNanos(o.TimeUnit))) +} + // List of bsiGroup types. const ( bsiGroupTypeInt = "int" @@ -2329,16 +2290,15 @@ func (f *Field) persistView(ctx context.Context, cvm *CreateViewMessage) error { return f.schemator.CreateView(ctx, cvm.Index, cvm.Field, cvm.View) } -// Timestamp field ranges. +// Timestamp field range. var ( - DefaultEpoch = time.Unix(0, 0).UTC() // 1970-01-01T00:00:00Z - MinTimestampNano = time.Unix(-1<<32, 0).UTC() // 1833-11-24T17:31:44Z - MaxTimestampNano = time.Unix(1<<32, 0).UTC() // 2106-02-07T06:28:16Z - MinTimestamp = time.Unix(-62135596799, 0).UTC() // 0001-01-01T00:00:01Z - MaxTimestamp = time.Unix(253402300799, 0).UTC() // 9999-12-31T23:59:59Z + DefaultEpoch = time.Unix(0, 0).UTC() // 1970-01-01T00:00:00Z + + MinTimestamp = time.Unix(-1<<32, 0).UTC() // 1833-11-24T17:31:44Z + MaxTimestamp = time.Unix(1<<32, 0).UTC() // 2106-02-07T06:28:16Z ) -// Constants related to timestamp. +// List of time units. const ( TimeUnitSeconds = "s" TimeUnitMilliseconds = "ms" @@ -2371,9 +2331,12 @@ func TimeUnitNanos(unit string) int64 { } } -// CheckEpochOutOfRange checks if the epoch is after max or before min -func CheckEpochOutOfRange(epoch, min, max time.Time) error { - if epoch.After(max) || epoch.Before(min) { +func CheckUnixNanoOverflow(epoch time.Time) error { + if time.Unix(0, 0).After(epoch) { + if epoch.UnixNano() > 0 { + return errors.Errorf("custom epoch too far from Unix epoch: %s", epoch) + } + } else if epoch.UnixNano() < 0 { return errors.Errorf("custom epoch too far from Unix epoch: %s", epoch) } return nil diff --git a/field_test.go b/field_test.go index 85299d4ea..1c602416a 100644 --- a/field_test.go +++ b/field_test.go @@ -307,8 +307,6 @@ func TestFieldInfoMarshal(t *testing.T) { } func TestCheckUnixNanoOverflow(t *testing.T) { - minNano = pilosa.MinTimestampNano.UnixNano() - maxNano = pilosa.MaxTimestampNano.UnixNano() tests := []struct { name string epoch time.Time @@ -316,28 +314,28 @@ func TestCheckUnixNanoOverflow(t *testing.T) { }{ { name: "too small", - epoch: time.Unix(-1, minNano), + epoch: time.Unix(-1, math.MinInt64), wantErr: true, }, { name: "just right-1", - epoch: time.Unix(0, minNano), + epoch: time.Unix(0, math.MinInt64), wantErr: false, }, { name: "just right-2", - epoch: time.Unix(0, maxNano), + epoch: time.Unix(0, math.MaxInt64), wantErr: false, }, { name: "too large", - epoch: time.Unix(1, maxNano), + epoch: time.Unix(1, math.MaxInt64), wantErr: true, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - if err := pilosa.CheckEpochOutOfRange(tt.epoch, pilosa.MinTimestampNano, pilosa.MaxTimestampNano); (err != nil) != tt.wantErr { + if err := pilosa.CheckUnixNanoOverflow(tt.epoch); (err != nil) != tt.wantErr { t.Errorf("checkUnixNanoOverflow() error = %v, wantErr %v", err, tt.wantErr) } }) diff --git a/fragment.go b/fragment.go index 37518afa5..97dcf1be6 100644 --- a/fragment.go +++ b/fragment.go @@ -828,7 +828,6 @@ func (f *fragment) min(tx Tx, filter *Row, bitDepth uint64) (min int64, count ui // minUnsigned the lowest value without considering the sign bit. Filter is required. func (f *fragment) minUnsigned(tx Tx, filter *Row, bitDepth uint64) (min int64, count uint64, err error) { - count = filter.Count() for i := int(bitDepth - 1); i >= 0; i-- { row, err := f.row(tx, uint64(bsiOffsetBit+i)) if err != nil { @@ -880,7 +879,6 @@ func (f *fragment) max(tx Tx, filter *Row, bitDepth uint64) (max int64, count ui // maxUnsigned the highest value without considering the sign bit. Filter is required. func (f *fragment) maxUnsigned(tx Tx, filter *Row, bitDepth uint64) (max int64, count uint64, err error) { - count = filter.Count() for i := int(bitDepth - 1); i >= 0; i-- { row, err := f.row(tx, uint64(bsiOffsetBit+i)) if err != nil { diff --git a/ingest/codec.go b/ingest/codec.go index 8c4abf5a8..c1e11f375 100644 --- a/ingest/codec.go +++ b/ingest/codec.go @@ -44,7 +44,7 @@ type Codec interface { AddBoolField(name string) error AddIntField(name string, keys KeyTranslator) error AddDecimalField(name string, scale int64) error - AddTimestampField(name string, scale string, epoch int64) error + AddTimestampField(name string, scale time.Duration, epoch int64) error // Parse data from a reader into the vectors. // This must only be called once on a codec. @@ -70,14 +70,14 @@ type fieldCodec struct { currentOp *FieldOperation decode jsonDecFn encode jsonEncFn + // For timestamp: scale-in-nanoseconds; for instance, if scaleUnit is + // 1,000,000,000, we are storing numbers-of-seconds since the Unix epoch. + // The actual value recorded in BSI will be offset by the field's + // epoch, but we don't need to know that. // For decimal: Decimal digits of precision. So for instance, with // scale 2, scaleUnit is 100, "1" is stored as 100 and "1.2" is stored as // 120. scaleUnit int64 - // For timestamp: we use timeUnit to determine the scale at which - // to store a timestamp. For example if the timeUnit is milliseconds - // we store the number of milliseconds from the given epoch. - timeUnit string scale int64 epoch int64 // used only by Timestamp fields scratch []uint64 // reusable scratch space for sets of values @@ -282,10 +282,10 @@ func (codec *JSONCodec) AddBoolField(name string) error { // passed to this function should be the offset from the Unix epoch to the // desired epoch, in the same scale. (So if the scale is milliseconds, // it should be the Unix timestamp in seconds, times 1000.) -func (codec *JSONCodec) AddTimestampField(name string, timeScale string, epoch int64) error { +func (codec *JSONCodec) AddTimestampField(name string, timeScale time.Duration, epoch int64) error { fieldCodec := &fieldCodec{ fieldType: FieldTypeTimeStamp, - timeUnit: timeScale, + scaleUnit: int64(timeScale), epoch: epoch, } fieldCodec.decode = fieldCodec.DecodeTimeValue @@ -590,7 +590,7 @@ func (j *fieldCodec) DecodeTimeValue(recID uint64, dataType jsonparser.ValueType if err != nil { return fmt.Errorf("parsing timestamp: %w", err) } - j.currentOp.AddSignedPair(recID, TimestampToVal(j.timeUnit, stamp)-j.epoch) + j.currentOp.AddSignedPair(recID, (stamp.UnixNano()/j.scaleUnit)-j.epoch) case jsonparser.Number: // We could in theory convert this to a time, then convert it // back, by multiplying by scaleUnit, then dividing. Or... not. @@ -609,11 +609,7 @@ func (j *fieldCodec) EncodeTimeValue(dst *jsonBuffer, values []uint64, signed [] if len(signed) == 0 { return errors.New("encoding time value: no value provided") } - t, err := ValToTimestamp(j.timeUnit, signed[0]+j.epoch) - if err != nil { - return errors.Wrap(err, "translating value to timestamp") - } - dst.EncodeTime(t) + dst.EncodeTime(time.Unix(0, (signed[0]+j.epoch)*j.scaleUnit).UTC()) return nil } @@ -1136,35 +1132,3 @@ type errFieldNotFound struct { func (err errFieldNotFound) Error() string { return fmt.Sprintf("field not found: %q", err.field) } - -// TimestampToVal takes a time unit and a time.Time and converts it to an integer value -func TimestampToVal(unit string, ts time.Time) int64 { - switch unit { - case "s": - return ts.Unix() - case "ms": - return ts.UnixMilli() - case "us": - return ts.UnixMicro() - case "ns": - return ts.UnixNano() - } - return 0 - -} - -// ValToTimestamp takes a timeunit and an integer value and converts it to time.Time -func ValToTimestamp(unit string, val int64) (time.Time, error) { - switch unit { - case "s": - return time.Unix(val, 0).UTC(), nil - case "ms": - return time.UnixMilli(val).UTC(), nil - case "us", "μs": - return time.UnixMicro(val).UTC(), nil - case "ns": - return time.Unix(0, val).UTC(), nil - default: - return time.Time{}, errors.Errorf("Unknown time unit: '%v'", unit) - } -} diff --git a/ingest/codec_test.go b/ingest/codec_test.go index 838328981..d36a0b757 100644 --- a/ingest/codec_test.go +++ b/ingest/codec_test.go @@ -58,7 +58,7 @@ func TestEncode(t *testing.T) { if err != nil { t.Fatalf("can't parse sample epoch time: %v", err) } - _ = codec.AddTimestampField("ts", "ms", epoch.Unix()*1000) + _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") @@ -74,7 +74,7 @@ func TestEncode(t *testing.T) { _ = codec.AddTimeQuantumField("tqkeys", newStableTranslator()) _ = codec.AddIntField("int", nil) _ = codec.AddIntField("intkeys", newStableTranslator()) - _ = codec.AddTimestampField("ts", "ms", epoch.Unix()*1000) + _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") @@ -291,7 +291,7 @@ func TestCodecErrors(t *testing.T) { if err != nil { t.Fatalf("can't parse sample epoch time: %v", err) } - _ = codec.AddTimestampField("ts", "ms", epoch.Unix()*1000) + _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") @@ -475,7 +475,7 @@ func TestSimpleCodec(t *testing.T) { if err != nil { t.Fatalf("can't parse sample epoch time: %v", err) } - _ = codec.AddTimestampField("ts", "ms", epoch.Unix()*1000) + _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") var nextShard = uint64(1<