mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-14 00:40:35 +01:00
Compare commits
2 Commits
fix/clear-
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0cf3988867 | ||
|
|
abff2aefd8 |
@@ -3104,7 +3104,19 @@ func (aH *APIHandler) PreviewLogsPipelinesHandler(w http.ResponseWriter, r *http
|
||||
return
|
||||
}
|
||||
|
||||
resultLogs, err := aH.LogsParsingPipelineController.PreviewLogsPipelines(r.Context(), &req)
|
||||
claims, errv2 := authtypes.ClaimsFromContext(r.Context())
|
||||
if errv2 != nil {
|
||||
render.Error(w, errv2)
|
||||
return
|
||||
}
|
||||
|
||||
orgID, errv2 := valuer.NewUUID(claims.OrgID)
|
||||
if errv2 != nil {
|
||||
render.Error(w, errv2)
|
||||
return
|
||||
}
|
||||
|
||||
resultLogs, err := aH.LogsParsingPipelineController.PreviewLogsPipelines(r.Context(), orgID, &req)
|
||||
if err != nil {
|
||||
render.Error(w, err)
|
||||
return
|
||||
|
||||
@@ -342,6 +342,7 @@ type PipelinesPreviewResponse struct {
|
||||
|
||||
func (ic *LogParsingPipelineController) PreviewLogsPipelines(
|
||||
ctx context.Context,
|
||||
orgID valuer.UUID,
|
||||
request *PipelinesPreviewRequest,
|
||||
) (*PipelinesPreviewResponse, error) {
|
||||
pipelines, err := ic.enrichPipelinesFilters(ctx, request.Pipelines)
|
||||
@@ -349,6 +350,11 @@ func (ic *LogParsingPipelineController) PreviewLogsPipelines(
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// The collector gets the same pipeline prepended over opamp; see RecommendAgentConfig.
|
||||
if ic.fl.BooleanOrEmpty(ctx, flagger.FeatureUseJSONBody, featuretypes.NewFlaggerEvaluationContext(orgID)) {
|
||||
pipelines = append([]pipelinetypes.GettablePipeline{ic.getNormalizePipeline()}, pipelines...)
|
||||
}
|
||||
|
||||
result, collectorLogs, err := SimulatePipelinesProcessing(ctx, pipelines, request.Logs)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@@ -6,11 +6,14 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
|
||||
"github.com/SigNoz/signoz/pkg/query-service/model"
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/types/pipelinetypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/google/uuid"
|
||||
"github.com/open-telemetry/opentelemetry-collector-contrib/pkg/stanza/entry"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
@@ -18,39 +21,7 @@ func TestPipelinePreview(t *testing.T) {
|
||||
require := require.New(t)
|
||||
|
||||
testPipelines := []pipelinetypes.GettablePipeline{
|
||||
{
|
||||
StoreablePipeline: pipelinetypes.StoreablePipeline{
|
||||
OrderID: 1,
|
||||
Name: "pipeline1",
|
||||
Alias: "pipeline1",
|
||||
Enabled: true,
|
||||
},
|
||||
Filter: &v3.FilterSet{
|
||||
Operator: "AND",
|
||||
Items: []v3.FilterItem{
|
||||
{
|
||||
Key: v3.AttributeKey{
|
||||
Key: "method",
|
||||
DataType: v3.AttributeKeyDataTypeString,
|
||||
Type: v3.AttributeKeyTypeTag,
|
||||
},
|
||||
Operator: "=",
|
||||
Value: "GET",
|
||||
},
|
||||
},
|
||||
},
|
||||
Config: []pipelinetypes.PipelineOperator{
|
||||
{
|
||||
OrderId: 1,
|
||||
ID: "add",
|
||||
Type: "add",
|
||||
Field: "attributes.test",
|
||||
Value: "val",
|
||||
Enabled: true,
|
||||
Name: "test add",
|
||||
},
|
||||
},
|
||||
},
|
||||
makeTestAddAttributePipeline(),
|
||||
{
|
||||
StoreablePipeline: pipelinetypes.StoreablePipeline{
|
||||
OrderID: 2,
|
||||
@@ -148,6 +119,93 @@ func TestPipelinePreview(t *testing.T) {
|
||||
|
||||
}
|
||||
|
||||
func TestPipelinePreviewNormalizesBodyWithJSONBodyEnabled(t *testing.T) {
|
||||
controller := &LogParsingPipelineController{fl: flaggertest.WithUseJSONBody(t, true)}
|
||||
|
||||
result, err := controller.PreviewLogsPipelines(
|
||||
context.Background(),
|
||||
valuer.GenerateUUID(),
|
||||
&PipelinesPreviewRequest{
|
||||
Pipelines: []pipelinetypes.GettablePipeline{makeTestAddAttributePipeline()},
|
||||
Logs: []model.SignozLog{
|
||||
makeTestSignozLog("test log body", map[string]interface{}{"method": "GET"}),
|
||||
makeTestSignozLog(
|
||||
`{"level":"error","msg":"json log body"}`,
|
||||
map[string]interface{}{"method": "GET"},
|
||||
),
|
||||
},
|
||||
},
|
||||
)
|
||||
|
||||
require.NoError(t, err)
|
||||
require.Len(t, result.OutputLogs, 2)
|
||||
|
||||
assert.Equal(t, `{"message":"test log body"}`, result.OutputLogs[0].Body)
|
||||
assert.Equal(
|
||||
t,
|
||||
`{"level":"error","message":"json log body"}`,
|
||||
result.OutputLogs[1].Body,
|
||||
)
|
||||
assert.Equal(t, "val", result.OutputLogs[0].Attributes_string["test"])
|
||||
}
|
||||
|
||||
func TestPipelinePreviewKeepsBodyAsIsWithJSONBodyDisabled(t *testing.T) {
|
||||
controller := &LogParsingPipelineController{fl: flaggertest.WithUseJSONBody(t, false)}
|
||||
|
||||
result, err := controller.PreviewLogsPipelines(
|
||||
context.Background(),
|
||||
valuer.GenerateUUID(),
|
||||
&PipelinesPreviewRequest{
|
||||
Pipelines: []pipelinetypes.GettablePipeline{makeTestAddAttributePipeline()},
|
||||
Logs: []model.SignozLog{
|
||||
makeTestSignozLog("test log body", map[string]interface{}{"method": "GET"}),
|
||||
},
|
||||
},
|
||||
)
|
||||
|
||||
require.NoError(t, err)
|
||||
require.Len(t, result.OutputLogs, 1)
|
||||
|
||||
assert.Equal(t, "test log body", result.OutputLogs[0].Body)
|
||||
assert.Equal(t, "val", result.OutputLogs[0].Attributes_string["test"])
|
||||
}
|
||||
|
||||
func makeTestAddAttributePipeline() pipelinetypes.GettablePipeline {
|
||||
return pipelinetypes.GettablePipeline{
|
||||
StoreablePipeline: pipelinetypes.StoreablePipeline{
|
||||
OrderID: 1,
|
||||
Name: "pipeline1",
|
||||
Alias: "pipeline1",
|
||||
Enabled: true,
|
||||
},
|
||||
Filter: &v3.FilterSet{
|
||||
Operator: "AND",
|
||||
Items: []v3.FilterItem{
|
||||
{
|
||||
Key: v3.AttributeKey{
|
||||
Key: "method",
|
||||
DataType: v3.AttributeKeyDataTypeString,
|
||||
Type: v3.AttributeKeyTypeTag,
|
||||
},
|
||||
Operator: "=",
|
||||
Value: "GET",
|
||||
},
|
||||
},
|
||||
},
|
||||
Config: []pipelinetypes.PipelineOperator{
|
||||
{
|
||||
OrderId: 1,
|
||||
ID: "add",
|
||||
Type: "add",
|
||||
Field: "attributes.test",
|
||||
Value: "val",
|
||||
Enabled: true,
|
||||
Name: "test add",
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func TestGrokParsingProcessor(t *testing.T) {
|
||||
require := require.New(t)
|
||||
|
||||
|
||||
@@ -123,23 +123,6 @@ func (b *StatementBuilder) Build(
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// TODO(srikanthccv): move the missing-key detection into the where clause
|
||||
// visitor. Doing it here over the lexer-derived selectors can't tell a key
|
||||
// from a value, so dashboard variables and bare literals in value position
|
||||
// (e.g. `service.name = $service`) get flagged as missing keys. We still add
|
||||
// a labels fallback for any unresolved selector so the query can be built,
|
||||
// but we no longer emit a warning until the visitor can classify keys.
|
||||
for _, sel := range keySelectors {
|
||||
if _, ok := keys[sel.Name]; !ok {
|
||||
keys[sel.Name] = []*telemetrytypes.TelemetryFieldKey{{
|
||||
Name: sel.Name,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Signal: telemetrytypes.SignalMetrics,
|
||||
}}
|
||||
}
|
||||
}
|
||||
|
||||
start, end = querybuilder.AdjustedMetricTimeRange(start, end, uint64(query.StepInterval.Seconds()), query)
|
||||
|
||||
return b.buildPipelineStatement(ctx, orgID, start, end, query, keys, variables)
|
||||
@@ -179,9 +162,10 @@ func (b *StatementBuilder) buildPipelineStatement(
|
||||
|
||||
var timeSeriesCTE string
|
||||
var timeSeriesCTEArgs []any
|
||||
var filterWarnings []string
|
||||
var err error
|
||||
|
||||
if timeSeriesCTE, timeSeriesCTEArgs, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, variables, tsTable); err != nil {
|
||||
if timeSeriesCTE, timeSeriesCTEArgs, filterWarnings, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, variables, tsTable); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -236,6 +220,7 @@ func (b *StatementBuilder) buildPipelineStatement(
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
mainStmt.Warnings = append(mainStmt.Warnings, filterWarnings...)
|
||||
if reducedFragments == nil {
|
||||
return mainStmt, nil
|
||||
}
|
||||
@@ -483,7 +468,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
|
||||
keys map[string][]*telemetrytypes.TelemetryFieldKey,
|
||||
variables map[string]qbtypes.VariableItem,
|
||||
tsTable string,
|
||||
) (string, []any, error) {
|
||||
) (string, []any, []string, error) {
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
var preparedWhereClause querybuilder.PreparedWhereClause
|
||||
@@ -503,7 +488,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
|
||||
EndNs: end,
|
||||
})
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
return "", nil, nil, err
|
||||
}
|
||||
}
|
||||
|
||||
@@ -513,7 +498,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
|
||||
for i, g := range query.GroupBy {
|
||||
col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
return "", nil, nil, err
|
||||
}
|
||||
sb.SelectMore(fmt.Sprintf("%s AS `%s`", sqlbuilder.Escape(col), GroupByColumnAlias(i, g.Name)))
|
||||
}
|
||||
@@ -542,7 +527,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
|
||||
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
||||
|
||||
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
return fmt.Sprintf("(%s) AS filtered_time_series", q), args, nil
|
||||
return fmt.Sprintf("(%s) AS filtered_time_series", q), args, preparedWhereClause.Warnings, nil
|
||||
}
|
||||
|
||||
func (b *StatementBuilder) buildTemporalAggregationCTE(
|
||||
|
||||
@@ -316,8 +316,9 @@ func TestStatementBuilder(t *testing.T) {
|
||||
},
|
||||
},
|
||||
expected: qbtypes.Statement{
|
||||
Query: "WITH __temporal_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`, multiIf(row_number() OVER rate_window = 1, nan, (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) < 0, per_series_value / (ts - lagInFrame(ts, 1) OVER rate_window), (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) / (ts - lagInFrame(ts, 1) OVER rate_window)) AS per_series_value FROM (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`, max(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, JSONExtractString(labels, 'k8s.statefulset.name') AS `__GROUP_BY_KEY_0_k8s.statefulset.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND JSONExtractString(labels, 'k8s.statefulset.name') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_k8s.statefulset.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_k8s.statefulset.name` ORDER BY fingerprint, ts) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`, sum(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_k8s.statefulset.name`, ts",
|
||||
Args: []any{"signoz_calls_total", uint64(1747936800000), uint64(1747983420000), "cumulative", "my-statefulset", "signoz_calls_total", uint64(1747947360000), uint64(1747983420000), 0},
|
||||
Query: "WITH __temporal_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`, multiIf(row_number() OVER rate_window = 1, nan, (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) < 0, per_series_value / (ts - lagInFrame(ts, 1) OVER rate_window), (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) / (ts - lagInFrame(ts, 1) OVER rate_window)) AS per_series_value FROM (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`, max(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, JSONExtractString(labels, 'k8s.statefulset.name') AS `__GROUP_BY_KEY_0_k8s.statefulset.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND JSONExtractString(labels, 'k8s.statefulset.name') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_k8s.statefulset.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_k8s.statefulset.name` ORDER BY fingerprint, ts) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`, sum(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_k8s.statefulset.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_k8s.statefulset.name`, ts",
|
||||
Args: []any{"signoz_calls_total", uint64(1747936800000), uint64(1747983420000), "cumulative", "my-statefulset", "signoz_calls_total", uint64(1747947360000), uint64(1747983420000), 0},
|
||||
Warnings: []string{"label `k8s.statefulset.name` not found in metadata; check the label name for typos"},
|
||||
},
|
||||
expectedErr: nil,
|
||||
},
|
||||
|
||||
@@ -175,8 +175,9 @@ def test_metrics_filter_unknown_label_matches_nothing(
|
||||
insert_metrics: Callable[[list[Metrics]], None],
|
||||
) -> None:
|
||||
"""A filter on a label no metric carries resolves to JSONExtractString(labels,'<missing>')
|
||||
= '' and matches nothing: metrics returns 200 with an empty result and — unlike the
|
||||
logs/traces synthesize path — emits no key-not-found warning."""
|
||||
= '' and matches nothing: metrics returns 200 with an empty result, and warns that the
|
||||
label is absent from metadata. Only keys are flagged — a value or dashboard variable in
|
||||
value position never reaches the condition builder, so it cannot be mistaken for one."""
|
||||
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
|
||||
insert_metrics(
|
||||
[
|
||||
@@ -208,7 +209,7 @@ def test_metrics_filter_unknown_label_matches_nothing(
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert querier.get_scalar_table_data(response.json()) == []
|
||||
assert querier.get_all_warnings(response.json()) == []
|
||||
assert [w["message"] for w in querier.get_all_warnings(response.json())] == ["label `does_not_exist_label` not found in metadata; check the label name for typos"]
|
||||
|
||||
|
||||
def test_metrics_full_text_filter_does_not_error(
|
||||
|
||||
Reference in New Issue
Block a user