mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-30 23:30:40 +01:00
Compare commits
17 Commits
feat/updat
...
ns/read-at
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4870fb92de | ||
|
|
ec060ad882 | ||
|
|
ed17f8e60a | ||
|
|
be1ec7102c | ||
|
|
9211b5e829 | ||
|
|
9f79083e99 | ||
|
|
e9748b9672 | ||
|
|
9a2afc3cab | ||
|
|
06a3f09aa4 | ||
|
|
3b8dd1119d | ||
|
|
20cfeeab3b | ||
|
|
a0fe8b5c6b | ||
|
|
e7cb50aa08 | ||
|
|
94eeeb20f1 | ||
|
|
1825a605e2 | ||
|
|
2f35a9059e | ||
|
|
7badbada4e |
2
go.mod
2
go.mod
@@ -12,7 +12,6 @@ require (
|
||||
github.com/SigNoz/signoz-otel-collector v0.144.6
|
||||
github.com/antlr4-go/antlr/v4 v4.13.1
|
||||
github.com/antonmedv/expr v1.15.3
|
||||
github.com/bytedance/sonic v1.14.1
|
||||
github.com/cespare/xxhash/v2 v2.3.0
|
||||
github.com/coreos/go-oidc/v3 v3.17.0
|
||||
github.com/dgraph-io/ristretto/v2 v2.3.0
|
||||
@@ -112,6 +111,7 @@ require (
|
||||
github.com/aws/aws-sdk-go-v2/service/sts v1.41.9 // indirect
|
||||
github.com/aws/smithy-go v1.24.2 // indirect
|
||||
github.com/bytedance/gopkg v0.1.3 // indirect
|
||||
github.com/bytedance/sonic v1.14.1 // indirect
|
||||
github.com/bytedance/sonic/loader v0.3.0 // indirect
|
||||
github.com/cloudwego/base64x v0.1.6 // indirect
|
||||
github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6 // indirect
|
||||
|
||||
@@ -17,7 +17,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
@@ -40,11 +39,18 @@ func stripKeyAlias(name string) string {
|
||||
return keyAliasRe.ReplaceAllString(name, "")
|
||||
}
|
||||
|
||||
// unwrapVariant returns the concrete value inside the chcol.Variant envelope the driver scans a
|
||||
// Dynamic column — a JSON path such as body_v2.level — into.
|
||||
func unwrapVariant(val any) any {
|
||||
if v, ok := val.(chcol.Variant); ok {
|
||||
// unwrapVariant decodes a scan envelope: a chcol.Variant to its value, a chcol.JSON column to a map.
|
||||
// The attributes bag decodes flat so its dotted keys merge with the legacy attribute maps and a
|
||||
// scalar-and-object key stays two keys; every other JSON column decodes nested.
|
||||
func unwrapVariant(name string, val any) any {
|
||||
switch v := val.(type) {
|
||||
case chcol.Variant:
|
||||
return v.Any()
|
||||
case chcol.JSON:
|
||||
if name == "attributes" {
|
||||
return v.ValuesByPath()
|
||||
}
|
||||
return v.NestedMap()
|
||||
}
|
||||
return val
|
||||
}
|
||||
@@ -54,11 +60,11 @@ func unwrapVariant(val any) any {
|
||||
// series. JSON goes through encoding/json for its sorted map keys: ClickHouse groups documents by
|
||||
// structure, so two rows it considers equal have to produce the same label.
|
||||
func labelValue(val any) string {
|
||||
val = unwrapVariant(val)
|
||||
val = unwrapVariant("", val)
|
||||
if val == nil {
|
||||
return ""
|
||||
}
|
||||
if v, ok := val.(telemetrystoretypes.JSONValue); ok {
|
||||
if v, ok := val.(map[string]any); ok {
|
||||
if raw, err := json.Marshal(v); err == nil {
|
||||
return string(raw)
|
||||
}
|
||||
@@ -204,7 +210,7 @@ func readAsTimeSeries(rows driver.Rows, queryWindow *qbtypes.TimeRange, step qbt
|
||||
Value: *val,
|
||||
})
|
||||
|
||||
case *telemetrystoretypes.JSONValue, *chcol.Variant:
|
||||
case *chcol.JSON, *chcol.Variant:
|
||||
val := labelValue(derefValue(ptr))
|
||||
lblVals = append(lblVals, val)
|
||||
lblObjs = append(lblObjs, &qbtypes.Label{
|
||||
@@ -478,7 +484,7 @@ func readAsScalar(rows driver.Rows, queryName string) (*qbtypes.ScalarData, erro
|
||||
// 2. deref each slot into the output row
|
||||
row := make([]any, len(scan))
|
||||
for i, cell := range scan {
|
||||
row[i] = unwrapVariant(derefValue(cell))
|
||||
row[i] = unwrapVariant(cd[i].Name, derefValue(cell))
|
||||
}
|
||||
data = append(data, row)
|
||||
}
|
||||
@@ -536,7 +542,7 @@ func readAsRaw(rows driver.Rows, queryName string) (*qbtypes.RawData, error) {
|
||||
name := stripKeyAlias(colNames[i])
|
||||
|
||||
// de-reference the typed pointer to any
|
||||
val := unwrapVariant(reflect.ValueOf(cellPtr).Elem().Interface())
|
||||
val := unwrapVariant(name, reflect.ValueOf(cellPtr).Elem().Interface())
|
||||
|
||||
// special-case: timestamp column
|
||||
if name == "timestamp" || name == "timestamp_datetime" {
|
||||
@@ -576,8 +582,6 @@ func flattenJSONPaths(prefix string, m map[string]any, out map[string]any) {
|
||||
switch child := v.(type) {
|
||||
case map[string]any:
|
||||
flattenJSONPaths(key, child, out)
|
||||
case telemetrystoretypes.JSONValue:
|
||||
flattenJSONPaths(key, child, out)
|
||||
default:
|
||||
out[key] = v
|
||||
}
|
||||
@@ -593,7 +597,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
|
||||
attrStr, hasStr := data["attributes_string"]
|
||||
attrNum, hasNum := data["attributes_number"]
|
||||
attrBool, hasBool := data["attributes_bool"]
|
||||
attrJSON, _ := data["attributes"].(telemetrystoretypes.JSONValue)
|
||||
attrJSON, _ := data["attributes"].(map[string]any)
|
||||
// todo(nitya): move to resource json
|
||||
resStr, hasRes := data["resources_string"]
|
||||
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {
|
||||
|
||||
@@ -3,14 +3,9 @@ package querier
|
||||
import (
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
|
||||
cmock "github.com/SigNoz/clickhouse-go-mock"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
@@ -83,101 +78,47 @@ func TestMergeSpanAttributeColumns_ParsesEventsAndLinks(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A ClickHouse query can put a JSON column in the result of any request type — e.g.
|
||||
// `select * from signoz_logs.logs_v2` on a body_v2 stack, where `*` covers body_v2.
|
||||
func TestConsume_JSONColumn(t *testing.T) {
|
||||
ts := time.Date(2026, 8, 14, 10, 0, 0, 0, time.UTC)
|
||||
body := `{"level":"error","attrs":{"code":500}}`
|
||||
wantBody := telemetrystoretypes.JSONValue{
|
||||
"level": "error",
|
||||
"attrs": map[string]any{"code": float64(500)},
|
||||
func TestUnwrapVariant(t *testing.T) {
|
||||
j := chcol.NewJSON()
|
||||
j.SetValueAtPath("level", "error")
|
||||
j.SetValueAtPath("attrs.code", int64(500))
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
column string
|
||||
input any
|
||||
want any
|
||||
}{
|
||||
{name: "VariantScalar", column: "", input: chcol.NewDynamicWithType("error", "String"), want: "error"},
|
||||
{name: "EmptyDynamic", column: "", input: chcol.Dynamic{}, want: nil},
|
||||
{name: "PlainValuePassthrough", column: "", input: uint64(3), want: uint64(3)},
|
||||
{name: "JSONColumnNested", column: "body_v2", input: *j, want: map[string]any{"level": "error", "attrs": map[string]any{"code": int64(500)}}},
|
||||
{name: "AttributesColumnFlat", column: "attributes", input: *j, want: map[string]any{"level": "error", "attrs.code": int64(500)}},
|
||||
}
|
||||
|
||||
// the scalar reader reuses its scan slots across rows, so each row must still carry its own body
|
||||
t.Run("scalar", func(t *testing.T) {
|
||||
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
|
||||
{Name: "body_v2", Type: "JSON"},
|
||||
{Name: "__result_0", Type: "UInt64"},
|
||||
}, [][]any{{body, uint64(3)}, {`{"level":"warn"}`, uint64(1)}}))
|
||||
|
||||
payload, err := consume(rows, qbtypes.RequestTypeScalar, nil, qbtypes.Step{}, "A")
|
||||
require.NoError(t, err)
|
||||
|
||||
data := payload.(*qbtypes.ScalarData)
|
||||
require.Len(t, data.Data, 2)
|
||||
assert.Equal(t, wantBody, data.Data[0][0])
|
||||
assert.Equal(t, uint64(3), data.Data[0][1])
|
||||
assert.Equal(t, telemetrystoretypes.JSONValue{"level": "warn"}, data.Data[1][0])
|
||||
assert.Equal(t, uint64(1), data.Data[1][1])
|
||||
})
|
||||
|
||||
t.Run("time series", func(t *testing.T) {
|
||||
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
|
||||
{Name: "ts", Type: "DateTime"},
|
||||
{Name: "body_v2", Type: "JSON"},
|
||||
{Name: "__result_0", Type: "UInt64"},
|
||||
}, [][]any{{ts, body, uint64(3)}}))
|
||||
|
||||
payload, err := consume(rows, qbtypes.RequestTypeTimeSeries, nil, qbtypes.Step{}, "A")
|
||||
require.NoError(t, err)
|
||||
|
||||
data := payload.(*qbtypes.TimeSeriesData)
|
||||
require.Len(t, data.Aggregations, 1)
|
||||
require.Len(t, data.Aggregations[0].Series, 1)
|
||||
require.Len(t, data.Aggregations[0].Series[0].Values, 1)
|
||||
assert.Equal(t, float64(3), data.Aggregations[0].Series[0].Values[0].Value)
|
||||
})
|
||||
|
||||
// grouping by a JSON column is legal in ClickHouse, so each document has to label its own
|
||||
// series rather than being dropped, which would merge every group into one
|
||||
t.Run("time series grouped by the JSON column", func(t *testing.T) {
|
||||
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
|
||||
{Name: "ts", Type: "DateTime"},
|
||||
{Name: "body_v2", Type: "JSON"},
|
||||
{Name: "__result_0", Type: "UInt64"},
|
||||
}, [][]any{
|
||||
{ts, `{"level":"error"}`, uint64(7)},
|
||||
{ts, `{"level":"warn"}`, uint64(2)},
|
||||
}))
|
||||
|
||||
payload, err := consume(rows, qbtypes.RequestTypeTimeSeries, nil, qbtypes.Step{}, "A")
|
||||
require.NoError(t, err)
|
||||
|
||||
data := payload.(*qbtypes.TimeSeriesData)
|
||||
require.Len(t, data.Aggregations, 1)
|
||||
require.Len(t, data.Aggregations[0].Series, 2)
|
||||
|
||||
got := map[string]float64{}
|
||||
for _, series := range data.Aggregations[0].Series {
|
||||
require.Len(t, series.Labels, 1)
|
||||
require.Len(t, series.Values, 1)
|
||||
got[series.Labels[0].Value.(string)] = series.Values[0].Value
|
||||
}
|
||||
assert.Equal(t, map[string]float64{`{"level":"error"}`: 7, `{"level":"warn"}`: 2}, got)
|
||||
})
|
||||
|
||||
t.Run("raw", func(t *testing.T) {
|
||||
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
|
||||
{Name: "timestamp", Type: "DateTime"},
|
||||
{Name: "body_v2", Type: "JSON"},
|
||||
}, [][]any{{ts, body}}))
|
||||
|
||||
payload, err := consume(rows, qbtypes.RequestTypeRaw, nil, qbtypes.Step{}, "A")
|
||||
require.NoError(t, err)
|
||||
|
||||
data := payload.(*qbtypes.RawData)
|
||||
require.Len(t, data.Rows, 1)
|
||||
assert.Equal(t, ts, data.Rows[0].Timestamp.UTC())
|
||||
assert.Equal(t, wantBody, data.Rows[0].Data["body_v2"])
|
||||
})
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
assert.Equal(t, testCase.want, unwrapVariant(testCase.column, testCase.input))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// A JSON path (e.g. `body_v2.level`) comes back as a Dynamic column, which the driver scans
|
||||
// into a chcol.Variant envelope rather than the value itself.
|
||||
func TestUnwrapVariant(t *testing.T) {
|
||||
assert.Equal(t, "error", unwrapVariant(chcol.NewDynamicWithType("error", "String")))
|
||||
assert.Nil(t, unwrapVariant(chcol.Dynamic{}))
|
||||
assert.Equal(t, uint64(3), unwrapVariant(uint64(3)))
|
||||
func TestLabelValue(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
input any
|
||||
want string
|
||||
}{
|
||||
{name: "Nil", input: nil, want: ""},
|
||||
{name: "Variant", input: chcol.NewDynamicWithType("error", "String"), want: "error"},
|
||||
{name: "JSONMap", input: map[string]any{"level": "error", "attrs": map[string]any{"code": 500}}, want: `{"attrs":{"code":500},"level":"error"}`},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
assert.Equal(t, testCase.want, labelValue(testCase.input))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestMergeSpanAttributeColumns_EmptyEventsAndLinks(t *testing.T) {
|
||||
@@ -207,7 +148,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
|
||||
{
|
||||
name: "JSONOnly_FlattensNestedPaths_PreservesTypes",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{
|
||||
"attributes": map[string]any{
|
||||
"http": map[string]any{"route": "/api/pay", "retry": map[string]any{"count": float64(3)}},
|
||||
"cache.hit": true,
|
||||
},
|
||||
@@ -219,7 +160,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
|
||||
data: map[string]any{
|
||||
"attributes_string": map[string]string{"http.route": "/old", "only.map": "m"},
|
||||
"attributes_number": map[string]float64{"http.status": 500},
|
||||
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"route": "/new"}, "only.json": "j"},
|
||||
"attributes": map[string]any{"http": map[string]any{"route": "/new"}, "only.json": "j"},
|
||||
},
|
||||
want: map[string]any{"http.route": "/old", "only.map": "m", "http.status": float64(500), "only.json": "j"},
|
||||
},
|
||||
@@ -229,7 +170,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
|
||||
"attributes_string": map[string]string{"http.route": "/map"},
|
||||
"attributes_number": map[string]float64{"http.status": 200},
|
||||
"attributes_bool": map[string]bool{"cache.hit": true},
|
||||
"attributes": telemetrystoretypes.JSONValue{},
|
||||
"attributes": map[string]any{},
|
||||
},
|
||||
want: map[string]any{"http.route": "/map", "http.status": float64(200), "cache.hit": true},
|
||||
},
|
||||
@@ -237,28 +178,28 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
|
||||
name: "MapOnly_NilJSON_BehavesAsAbsent",
|
||||
data: map[string]any{
|
||||
"attributes_string": map[string]string{"http.route": "/map"},
|
||||
"attributes": telemetrystoretypes.JSONValue(nil),
|
||||
"attributes": map[string]any(nil),
|
||||
},
|
||||
want: map[string]any{"http.route": "/map"},
|
||||
},
|
||||
{
|
||||
name: "Arrays_StayLeafValues",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"tags": []any{"a", "b"}, "codes": []any{float64(1), float64(2)}}},
|
||||
"attributes": map[string]any{"http": map[string]any{"tags": []any{"a", "b"}, "codes": []any{float64(1), float64(2)}}},
|
||||
},
|
||||
want: map[string]any{"http.tags": []any{"a", "b"}, "http.codes": []any{float64(1), float64(2)}},
|
||||
},
|
||||
{
|
||||
name: "TopLevelArrayOfMaps_StaysNativeLeaf",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
|
||||
"attributes": map[string]any{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
|
||||
},
|
||||
want: map[string]any{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
|
||||
},
|
||||
{
|
||||
name: "NestedArrayOfMaps_StaysNativeLeaf_NoIndexPaths",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
|
||||
"attributes": map[string]any{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
|
||||
},
|
||||
want: map[string]any{"http.items": []any{map[string]any{"a": float64(1)}}},
|
||||
},
|
||||
@@ -266,28 +207,28 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
|
||||
name: "DualWritten_NestedArray_IndexKeysAndJSONArrayCoexist",
|
||||
data: map[string]any{
|
||||
"attributes_number": map[string]float64{"http.items.0.a": 1},
|
||||
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
|
||||
"attributes": map[string]any{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
|
||||
},
|
||||
want: map[string]any{"http.items.0.a": float64(1), "http.items": []any{map[string]any{"a": float64(1)}}},
|
||||
},
|
||||
{
|
||||
name: "JSONNull_KeptAsNil",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{"k": nil},
|
||||
"attributes": map[string]any{"k": nil},
|
||||
},
|
||||
want: map[string]any{"k": nil},
|
||||
},
|
||||
{
|
||||
name: "KeyIsLeafValue_NotFlattened",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{"http": "plaintext"},
|
||||
"attributes": map[string]any{"http": "plaintext"},
|
||||
},
|
||||
want: map[string]any{"http": "plaintext"},
|
||||
},
|
||||
{
|
||||
name: "KeyIsParent_FlattensToDottedPath",
|
||||
data: map[string]any{
|
||||
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"route": "/a"}},
|
||||
"attributes": map[string]any{"http": map[string]any{"route": "/a"}},
|
||||
},
|
||||
want: map[string]any{"http.route": "/a"},
|
||||
},
|
||||
|
||||
@@ -16,7 +16,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
"github.com/SigNoz/signoz/pkg/types/featuretypes"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
@@ -1202,7 +1201,7 @@ func (q *querier) postProcessLogBody(ctx context.Context, orgID valuer.UUID, res
|
||||
// carried one. Anything that is not a decoded document — the legacy string body, a NULL cell —
|
||||
// is legal under these names and left alone.
|
||||
func stripEmptyBodyMessage(val any) {
|
||||
bodyMap, ok := val.(telemetrystoretypes.JSONValue)
|
||||
bodyMap, ok := val.(map[string]any)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -97,8 +97,8 @@ func New(ctx context.Context, providerSettings factory.ProviderSettings, config
|
||||
options.MaxIdleConns = config.Connection.MaxIdleConns
|
||||
options.MaxOpenConns = config.Connection.MaxOpenConns
|
||||
options.DialTimeout = config.Connection.DialTimeout
|
||||
// This is to avoid the driver decoding issues with JSON columns
|
||||
options.Settings["output_format_native_write_json_as_string"] = 1
|
||||
// Native flattened JSON serialization (CH 25.6+): clickhouse-go mis-decodes JSON(max_dynamic_paths=0) columns without it.
|
||||
options.Settings["output_format_native_use_flattened_dynamic_and_json_serialization"] = 1
|
||||
|
||||
chConn, err := clickhouse.Open(options)
|
||||
if err != nil {
|
||||
@@ -184,7 +184,7 @@ func (p *provider) Query(ctx context.Context, query string, args ...interface{})
|
||||
}
|
||||
|
||||
return &rowsWithHooks{
|
||||
Rows: telemetrystore.WrapRows(rows),
|
||||
Rows: rows,
|
||||
ctx: ctx,
|
||||
event: event,
|
||||
onClose: func() { telemetrystore.WrapAfterQuery(p.hooks, ctx, event) },
|
||||
|
||||
@@ -1,39 +0,0 @@
|
||||
package telemetrystore
|
||||
|
||||
import (
|
||||
"reflect"
|
||||
"strings"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
)
|
||||
|
||||
// WrapRows reports JSONValue as the scan type of every JSON column. Nested JSON — Array(JSON),
|
||||
// Map(String, JSON) — is not covered.
|
||||
func WrapRows(rows driver.Rows) driver.Rows {
|
||||
return &rowsWithJSONScanType{Rows: rows}
|
||||
}
|
||||
|
||||
type rowsWithJSONScanType struct {
|
||||
driver.Rows
|
||||
}
|
||||
|
||||
func (r *rowsWithJSONScanType) ColumnTypes() []driver.ColumnType {
|
||||
colTypes := r.Rows.ColumnTypes()
|
||||
wrapped := make([]driver.ColumnType, len(colTypes))
|
||||
for i, colType := range colTypes {
|
||||
wrapped[i] = colType
|
||||
if strings.HasPrefix(strings.ToUpper(colType.DatabaseTypeName()), "JSON") {
|
||||
wrapped[i] = jsonColumnType{ColumnType: colType}
|
||||
}
|
||||
}
|
||||
return wrapped
|
||||
}
|
||||
|
||||
type jsonColumnType struct {
|
||||
driver.ColumnType
|
||||
}
|
||||
|
||||
func (jsonColumnType) ScanType() reflect.Type {
|
||||
return reflect.TypeFor[telemetrystoretypes.JSONValue]()
|
||||
}
|
||||
@@ -1,23 +0,0 @@
|
||||
package telemetrystoretest
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2"
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
)
|
||||
|
||||
// conn wraps rows the way the clickhouse provider does, so mocked JSON columns report the scan
|
||||
// type they do in production.
|
||||
type conn struct {
|
||||
clickhouse.Conn
|
||||
}
|
||||
|
||||
func (c conn) Query(ctx context.Context, query string, args ...any) (driver.Rows, error) {
|
||||
rows, err := c.Conn.Query(ctx, query, args...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return telemetrystore.WrapRows(rows), nil
|
||||
}
|
||||
@@ -32,7 +32,7 @@ func New(_ telemetrystore.Config, matcher sqlmock.QueryMatcher) *Provider {
|
||||
|
||||
// ClickhouseDB returns the mock Clickhouse connection.
|
||||
func (p *Provider) ClickhouseDB() clickhouse.Conn {
|
||||
return conn{Conn: p.clickhouseDB.(clickhouse.Conn)}
|
||||
return p.clickhouseDB.(clickhouse.Conn)
|
||||
}
|
||||
|
||||
// Cluster returns the cluster name.
|
||||
|
||||
@@ -1,37 +0,0 @@
|
||||
package telemetrystoretypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/bytedance/sonic"
|
||||
)
|
||||
|
||||
var ErrCodeUnmarshalJSONColumn = errors.MustNewCode("fail_unmarshal_json_column")
|
||||
|
||||
// JSONValue is the scan target for a ClickHouse JSON column: the connection sets
|
||||
// output_format_native_write_json_as_string, so the column arrives as a raw document rather than
|
||||
// the chcol.JSON the driver reports as its scan type.
|
||||
type JSONValue map[string]any
|
||||
|
||||
// Scan decodes into a fresh map every time: a scan target is reused across rows, and unmarshalling
|
||||
// into the map already there would both keep its keys and hand every row the same map.
|
||||
func (v *JSONValue) Scan(src any) error {
|
||||
var raw []byte
|
||||
switch value := src.(type) {
|
||||
case nil:
|
||||
*v = nil
|
||||
return nil
|
||||
case string:
|
||||
raw = []byte(value)
|
||||
case []byte:
|
||||
raw = value
|
||||
default:
|
||||
return errors.NewInternalf(ErrCodeUnmarshalJSONColumn, "cannot decode %T as a JSON column", src)
|
||||
}
|
||||
|
||||
decoded := JSONValue{}
|
||||
if err := sonic.Unmarshal(raw, &decoded); err != nil {
|
||||
return errors.WrapInternalf(err, ErrCodeUnmarshalJSONColumn, "failed to unmarshal JSON column")
|
||||
}
|
||||
*v = decoded
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
import json
|
||||
from collections.abc import Callable
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.logs import Logs
|
||||
from fixtures.querier import RequestType, build_raw_query, get_rows, make_query_request
|
||||
|
||||
SERVICE = "scalar-object-body"
|
||||
|
||||
|
||||
def test_scalar_and_object_body_key_collapses(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_logs: Callable[[list[Logs]], None],
|
||||
export_json_types: Callable[[list[Logs]], None],
|
||||
) -> None:
|
||||
"""A log body key stored as both a scalar and an object (db.function and db.function.arg) keeps
|
||||
only one: the body is returned nested, not flattened to dotted keys, so it cannot hold both."""
|
||||
now = datetime.now(tz=UTC)
|
||||
logs = [
|
||||
Logs(
|
||||
timestamp=now - timedelta(seconds=1),
|
||||
resources={"service.name": SERVICE},
|
||||
body_v2=json.dumps({"db.function": "refresh", "db.function.arg": 2}),
|
||||
body_promoted="",
|
||||
),
|
||||
]
|
||||
export_json_types(logs)
|
||||
insert_logs(logs)
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = make_query_request(
|
||||
signoz,
|
||||
token,
|
||||
start_ms=int((now - timedelta(minutes=5)).timestamp() * 1000),
|
||||
end_ms=int(now.timestamp() * 1000),
|
||||
request_type=RequestType.RAW,
|
||||
queries=[build_raw_query("A", "logs", limit=10, filter_expression=f"service.name = '{SERVICE}'")],
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
rows = get_rows(response)
|
||||
assert len(rows) == 1
|
||||
body = rows[0]["data"]["body"]
|
||||
# The scalar survives as a string; the sibling object path db.function.arg is not preserved
|
||||
# under it, so db.function is a plain value rather than an object.
|
||||
assert body["db"]["function"] == "refresh"
|
||||
assert isinstance(body["db"]["function"], str)
|
||||
@@ -0,0 +1,72 @@
|
||||
from collections.abc import Callable
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.querier import (
|
||||
BuilderQuery,
|
||||
OrderBy,
|
||||
RequestType,
|
||||
TelemetryFieldKey,
|
||||
get_rows,
|
||||
make_query_request,
|
||||
)
|
||||
from fixtures.traces import TraceIdGenerator, Traces
|
||||
|
||||
|
||||
def test_traces_scalar_and_object_key_both_survive(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
seed_attribute_evolution: Callable[[str, datetime], None],
|
||||
) -> None:
|
||||
"""A span attribute stored as both a scalar and an object prefix (db.function and
|
||||
db.function.arg) keeps both dotted keys through the native JSON read. json_only writes only the
|
||||
JSON column, so no legacy attribute map backfills the keys."""
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
evolution_time = datetime.now(tz=UTC).replace(second=0, microsecond=0) - timedelta(minutes=30)
|
||||
seed_attribute_evolution("traces", evolution_time)
|
||||
|
||||
service = "scalar-object-service"
|
||||
span_time = evolution_time + timedelta(minutes=5)
|
||||
insert_traces(
|
||||
[
|
||||
Traces(
|
||||
timestamp=span_time,
|
||||
trace_id=TraceIdGenerator.trace_id(),
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
name="scalar and object",
|
||||
resources={"service.name": service},
|
||||
attributes={"db.function": "refresh", "db.function.arg": 2},
|
||||
attribute_write_mode="json_only",
|
||||
),
|
||||
]
|
||||
)
|
||||
|
||||
response = make_query_request(
|
||||
signoz,
|
||||
token,
|
||||
start_ms=int((span_time - timedelta(minutes=1)).timestamp() * 1000),
|
||||
end_ms=int((span_time + timedelta(minutes=1)).timestamp() * 1000),
|
||||
request_type=RequestType.RAW,
|
||||
queries=[
|
||||
BuilderQuery(
|
||||
signal="traces",
|
||||
name="A",
|
||||
limit=10,
|
||||
filter_expression=f"resource.service.name = '{service}'",
|
||||
order=[OrderBy(TelemetryFieldKey("timestamp"), "asc")],
|
||||
).to_dict()
|
||||
],
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
rows = get_rows(response)
|
||||
assert len(rows) == 1
|
||||
attributes = rows[0]["data"]["attributes"]
|
||||
assert attributes["db.function"] == "refresh"
|
||||
assert isinstance(attributes["db.function"], str)
|
||||
assert attributes["db.function.arg"] == 2
|
||||
Reference in New Issue
Block a user