mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-06 19:07:50 +00:00
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:0676790update codec to reflect changes to timstamp range:bd5dc76fix few bugs regarding timestamp:33fce8aincrease time range for timestamp by using specified granularity:5939923.
This commit is contained in:
parent
d348cc65e9
commit
33a916c0f4
12 changed files with 120 additions and 882 deletions
|
|
@ -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
|
||||
|
|
|
|||
3
api.go
3
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:
|
||||
|
|
|
|||
117
executor.go
117
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:
|
||||
|
|
|
|||
625
executor_test.go
625
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)
|
||||
|
|
|
|||
147
field.go
147
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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
})
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<<shardwidth.Exponent) + 5
|
||||
|
|
|
|||
|
|
@ -319,13 +319,7 @@ func pgWriteGroupCount(w pg.QueryResultWriter, counts *pilosa.GroupCounts) error
|
|||
switch {
|
||||
case g.Value != nil:
|
||||
if g.FieldOptions.Type == pilosa.FieldTypeTimestamp {
|
||||
ts, err := pilosa.ValToTimestamp(g.FieldOptions.TimeUnit, int64(*g.Value)+g.FieldOptions.Base)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "translating val to timestamp")
|
||||
}
|
||||
|
||||
v = ts.Format(time.RFC3339Nano)
|
||||
// v = pilosa.FormatTimestampNano(int64(*g.Value), g.FieldOptions.Base, g.FieldOptions.TimeUnit)
|
||||
v = pilosa.FormatTimestampNano(int64(*g.Value), g.FieldOptions.Base, g.FieldOptions.TimeUnit)
|
||||
} else {
|
||||
v = strconv.FormatInt(*g.Value, 10)
|
||||
}
|
||||
|
|
|
|||
6
util.go
6
util.go
|
|
@ -60,6 +60,12 @@ func GetLoopProgress(start time.Time, now time.Time, iteration uint, total uint)
|
|||
return time.Duration(avgItemTime * float64(itemsLeft)), pctDone
|
||||
}
|
||||
|
||||
// FormatTimestampNano returns the string representation of a timestamp given:
|
||||
// an epoch value, base, and time unit
|
||||
func FormatTimestampNano(value, base int64, timeUnit string) string {
|
||||
return time.Unix(0, (value+base)*TimeUnitNanos(timeUnit)).UTC().Format(time.RFC3339Nano)
|
||||
}
|
||||
|
||||
type MemoryUsage struct {
|
||||
Capacity uint64 `json:"capacity"`
|
||||
TotalUse uint64 `json:"totalUsed"`
|
||||
|
|
|
|||
18
util_test.go
18
util_test.go
|
|
@ -72,6 +72,24 @@ func TestGetLoopProgress(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestFormatTimestampNano(t *testing.T) {
|
||||
if FormatTimestampNano(0, 69, "s") != "1970-01-01T00:01:09Z" {
|
||||
t.Fatal("Timestamp not formatted properly")
|
||||
}
|
||||
if FormatTimestampNano(0, 420, "ms") != "1970-01-01T00:00:00.42Z" {
|
||||
t.Fatal("Timestamp not formatted properly")
|
||||
}
|
||||
if FormatTimestampNano(420, 0, "μs") != "1970-01-01T00:00:00.00000042Z" {
|
||||
t.Fatal("Timestamp not formatted properly")
|
||||
}
|
||||
if FormatTimestampNano(420, 69, "us") != "1970-01-01T00:00:00.000489Z" {
|
||||
t.Fatal("Timestamp not formatted properly")
|
||||
}
|
||||
if FormatTimestampNano(69, 420, "ns") != "1970-01-01T00:00:00.000000489Z" {
|
||||
t.Fatal("Timestamp not formatted properly")
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetMemoryUsage(t *testing.T) {
|
||||
if _, err := GetMemoryUsage(); err != nil {
|
||||
t.Fatalf("unexpected error getting memory usage: %v", err)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue