Compare commits

...

10 Commits

Author SHA1 Message Date
nikhilmantri0902
f5bd7416d4 chore: test case both columns contain host.name -> resources take precedence 2026-07-24 02:27:56 +05:30
nikhilmantri0902
85f1bf29ce chore: more simplified comment 2026-07-24 02:19:01 +05:30
nikhilmantri0902
6d72efde81 chore: simplified comment 2026-07-24 00:37:44 +05:30
nikhilmantri0902
f66ea90b35 chore: added integration tests 2026-07-23 21:56:37 +05:30
nikhilmantri0902
9f34bcf5b3 chore: added data_source condition 2026-07-23 21:06:18 +05:30
nikhilmantri0902
e0a8f41ac6 chore: flip AND to or in outside condition 2026-07-23 19:46:08 +05:30
Nikhil Mantri
63e890a20e Merge branch 'main' into querier/telemetrymetadata_fix_map_contains_fallback 2026-07-23 18:52:54 +05:30
Nikhil Mantri
1783b942a2 Merge branch 'main' into querier/telemetrymetadata_fix_map_contains_fallback 2026-07-23 14:32:41 +05:30
nikhilmantri0902
ab3ed06fc0 chore: changed to distributed table 2026-07-23 14:27:44 +05:30
nikhilmantri0902
b483cb3545 chore: update fallback to false for positive operators 2026-07-23 13:28:27 +05:30
10 changed files with 511 additions and 10 deletions

View File

@@ -70,7 +70,7 @@ func newProvider(
telemetryaudit.LogAttributeKeysTblName,
telemetryaudit.LogResourceKeysTblName,
telemetrymetadata.DBName,
telemetrymetadata.AttributesMetadataLocalTableName,
telemetrymetadata.AttributesMetadataTableName,
telemetrymetadata.ColumnEvolutionMetadataTableName,
flagger,
)

View File

@@ -447,7 +447,7 @@ func New(
telemetryaudit.LogAttributeKeysTblName,
telemetryaudit.LogResourceKeysTblName,
telemetrymetadata.DBName,
telemetrymetadata.AttributesMetadataLocalTableName,
telemetrymetadata.AttributesMetadataTableName,
telemetrymetadata.ColumnEvolutionMetadataTableName,
flagger,
)

View File

@@ -96,8 +96,11 @@ func (c *conditionBuilder) conditionForKey(
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(key, value, fieldExpression, operator)
// key must exists to apply main filter
expr := `if(mapContains(%s, %s), %s, true)`
// key must exist to apply the main filter. for positive operators the
// absent-key rows are excluded (fallback false); for negative operators
// they are kept (fallback true) so rows legitimately lacking the key match.
keyMissingFallback := operator.IsNegativeOperator()
expr := `if(mapContains(%s, %s), %s, %t)`
var cond string
@@ -171,5 +174,5 @@ func (c *conditionBuilder) conditionForKey(
}
}
return fmt.Sprintf(expr, columns[0].Name, sb.Var(key.Name), cond), nil
return fmt.Sprintf(expr, columns[0].Name, sb.Var(key.Name), cond, keyMissingFallback), nil
}

View File

@@ -34,7 +34,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorILike,
value: "%admin%",
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) LIKE LOWER(?), true)",
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) LIKE LOWER(?), false)",
expectedError: nil,
},
{
@@ -49,6 +49,150 @@ func TestConditionFor(t *testing.T) {
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) NOT LIKE LOWER(?), true)",
expectedError: nil,
},
{
name: "Equal operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorEqual,
value: "admin",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] = ?, false)",
expectedError: nil,
},
{
name: "Not Equal operator - negative fallback true",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorNotEqual,
value: "admin",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] <> ?, true)",
expectedError: nil,
},
{
name: "In operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorIn,
value: []any{"admin", "root"},
expectedSQL: "WHERE if(mapContains(attributes, ?), (attributes['user.id'] = ? OR attributes['user.id'] = ?), false)",
expectedError: nil,
},
{
name: "Not In operator - negative fallback true",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorNotIn,
value: []any{"admin", "root"},
expectedSQL: "WHERE if(mapContains(attributes, ?), (attributes['user.id'] <> ? AND attributes['user.id'] <> ?), true)",
expectedError: nil,
},
{
name: "Like operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorLike,
value: "%admin%",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] LIKE ?, false)",
expectedError: nil,
},
{
name: "Not Like operator - negative fallback true",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorNotLike,
value: "%admin%",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] NOT LIKE ?, true)",
expectedError: nil,
},
{
name: "Contains operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorContains,
value: "admin",
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) LIKE LOWER(?), false)",
expectedError: nil,
},
{
name: "Not Contains operator - negative fallback true",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorNotContains,
value: "admin",
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) NOT LIKE LOWER(?), true)",
expectedError: nil,
},
{
name: "Regexp operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorRegexp,
value: "adm.*",
expectedSQL: "WHERE if(mapContains(attributes, ?), match(attributes['user.id'], ?), false)",
expectedError: nil,
},
{
name: "Not Regexp operator - negative fallback true",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorNotRegexp,
value: "adm.*",
expectedSQL: "WHERE if(mapContains(attributes, ?), NOT match(attributes['user.id'], ?), true)",
expectedError: nil,
},
{
name: "Exists operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorExists,
value: nil,
expectedSQL: "WHERE if(mapContains(attributes, ?), mapContains(attributes, 'user.id') = ?, false)",
expectedError: nil,
},
{
name: "Not Exists operator - negative fallback true",
key: telemetrytypes.TelemetryFieldKey{
Name: "user.id",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorNotExists,
value: nil,
expectedSQL: "WHERE if(mapContains(attributes, ?), mapContains(attributes, 'user.id') <> ?, true)",
expectedError: nil,
},
}
for _, tc := range testCases {

View File

@@ -1436,6 +1436,11 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
sb.Where(sb.LE("unix_milli", fieldValueSelector.EndUnixMilli))
}
// scope to the requested signal's rows;
if fieldValueSelector.Signal != telemetrytypes.SignalUnspecified {
sb.Where(sb.E("data_source", fieldValueSelector.Signal.StringValue()))
}
if fieldValueSelector.Value != "" {
var conds []string
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextAttribute &&
@@ -1464,9 +1469,9 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
}
if len(conds) != 0 {
// see `expr` in condition_builder.go, if key doesn't exist we don't check for value
// hence, this is join of conditions on resource and attributes
sb.Where(sb.And(conds...))
// the key may sit in the resource or attribute map (or both), so OR the
// two conditions — match if the key's value in either map contains searchText.
sb.Where(sb.Or(conds...))
}
}

View File

@@ -61,7 +61,7 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
telemetryaudit.LogAttributeKeysTblName,
telemetryaudit.LogResourceKeysTblName,
DBName,
AttributesMetadataLocalTableName,
AttributesMetadataTableName,
ColumnEvolutionMetadataTableName,
flaggertest.New(t),
)

View File

@@ -16,6 +16,7 @@ pytest_plugins = [
"fixtures.logs",
"fixtures.traces",
"fixtures.metrics",
"fixtures.metadata",
"fixtures.meter",
"fixtures.browser",
"fixtures.keycloak",

121
tests/fixtures/metadata.py vendored Normal file
View File

@@ -0,0 +1,121 @@
import datetime
import json
from abc import ABC
from collections.abc import Callable, Generator
from typing import Any
import numpy as np
import pytest
from fixtures import types
from fixtures.fingerprint import LogsOrTracesFingerprint
class AttributesMetadata(ABC):
"""Represents a row in signoz_metadata.attributes_metadata.
This is the table `getRelatedValues` reads for the related-values
suggestions returned by /api/v1/fields/values. `data_source` is one of
logs/traces/metrics; fingerprints mirror the collector's FNV-1a hashes so
distinct rows are not collapsed by the ReplacingMergeTree ORDER BY key.
"""
unix_milli: np.int64
data_source: str
resource_fingerprint: np.uint64
attrs_fingerprint: np.uint64
resource_attributes: dict[str, str]
attributes: dict[str, str]
def __init__(
self,
data_source: str,
resource_attributes: dict[str, Any] = {},
attributes: dict[str, Any] = {},
timestamp: datetime.datetime | None = None,
) -> None:
if timestamp is None:
timestamp = datetime.datetime.now()
self.data_source = data_source
self.resource_attributes = {k: str(v) for k, v in resource_attributes.items()}
self.attributes = {k: str(v) for k, v in attributes.items()}
self.unix_milli = np.int64(int(timestamp.timestamp() * 1e3))
self.resource_fingerprint = np.uint64(LogsOrTracesFingerprint.hash(self.resource_attributes))
self.attrs_fingerprint = np.uint64(LogsOrTracesFingerprint.hash(self.attributes))
def to_row(self) -> list:
return [
self.unix_milli,
self.data_source,
self.resource_fingerprint,
self.attrs_fingerprint,
self.resource_attributes,
self.attributes,
]
@classmethod
def from_dict(cls, data: dict, timestamp: datetime.datetime | None = None) -> "AttributesMetadata":
return cls(
data_source=data["data_source"],
resource_attributes=data.get("resource_attributes", {}),
attributes=data.get("attributes", {}),
timestamp=timestamp,
)
@classmethod
def load_from_file(cls, file_path: str, timestamp: datetime.datetime | None = None) -> list["AttributesMetadata"]:
"""Load rows from a JSONL file, stamping each with `timestamp`.
Each line is a JSON object with `data_source` and optional
`resource_attributes` / `attributes` maps.
"""
rows: list[AttributesMetadata] = []
with open(file_path, encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
rows.append(cls.from_dict(json.loads(line), timestamp=timestamp))
return rows
def insert_attributes_metadata_to_clickhouse(conn, rows: list[AttributesMetadata]) -> None:
"""Insert rows into signoz_metadata.distributed_attributes_metadata.
Pure function so it can be reused outside the pytest fixture. `conn` is a
clickhouse-connect Client.
"""
if not rows:
return
conn.insert(
database="signoz_metadata",
table="distributed_attributes_metadata",
column_names=[
"unix_milli",
"data_source",
"resource_fingerprint",
"attrs_fingerprint",
"resource_attributes",
"attributes",
],
data=[row.to_row() for row in rows],
)
def truncate_attributes_metadata_table(conn, cluster: str) -> None:
conn.query(f"TRUNCATE TABLE signoz_metadata.attributes_metadata ON CLUSTER '{cluster}' SYNC")
@pytest.fixture(name="insert_attributes_metadata", scope="function")
def insert_attributes_metadata(
clickhouse: types.TestContainerClickhouse,
) -> Generator[Callable[[list[AttributesMetadata]], None], Any]:
def _insert(rows: list[AttributesMetadata]) -> None:
insert_attributes_metadata_to_clickhouse(clickhouse.conn, rows)
yield _insert
truncate_attributes_metadata_table(
clickhouse.conn,
clickhouse.env["SIGNOZ_TELEMETRYSTORE_CLICKHOUSE_CLUSTER"],
)

View File

@@ -0,0 +1,18 @@
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-a", "k8s.statefulset.name": "sts-a", "k8s.pod.name": "podA1-logs"}}
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-a", "k8s.deployment.name": "dep-a", "k8s.pod.name": "podA2-logs"}}
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-b", "k8s.pod.name": "podB1-logs"}}
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.pod.name": "podOrphan-logs"}}
{"data_source": "traces", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-a", "k8s.statefulset.name": "sts-a", "k8s.pod.name": "podA1-traces"}}
{"data_source": "traces", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-a", "k8s.deployment.name": "dep-a", "k8s.pod.name": "podA2-traces"}}
{"data_source": "traces", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-b", "k8s.pod.name": "podB1-traces"}}
{"data_source": "traces", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.pod.name": "podOrphan-traces"}}
{"data_source": "metrics", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-a", "k8s.statefulset.name": "sts-a", "k8s.pod.name": "podA1-metrics"}}
{"data_source": "metrics", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-a", "k8s.deployment.name": "dep-a", "k8s.pod.name": "podA2-metrics"}}
{"data_source": "metrics", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "ns-b", "k8s.pod.name": "podB1-metrics"}}
{"data_source": "metrics", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.pod.name": "podOrphan-metrics"}}
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1", "k8s.namespace.name": "", "k8s.pod.name": "podEmpty-logs"}}
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1", "host.name": "host-res-logs"}}
{"data_source": "logs", "resource_attributes": {"k8s.cluster.name": "c1"}, "attributes": {"host.name": "host-attr-logs"}}
{"data_source": "metrics", "resource_attributes": {"k8s.cluster.name": "c1", "host.name": "host-res-metrics"}}
{"data_source": "metrics", "resource_attributes": {"k8s.cluster.name": "c1"}, "attributes": {"host.name": "host-attr-metrics"}}
{"data_source": "traces", "resource_attributes": {"k8s.cluster.name": "c1", "host.name": "host-rest"}, "attributes": {"host.name": "host-attributes"}}

View File

@@ -0,0 +1,209 @@
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
import pytest
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.fs import get_testdata_file_path
from fixtures.logs import Logs
from fixtures.metadata import AttributesMetadata
from fixtures.metrics import Metrics
from fixtures.traces import Traces
# Related values (/api/v1/fields/values with existingQuery) are served from the
# shared signoz_metadata.attributes_metadata table, scoped by data_source =
# signal. One dataset (queriercommon/related_values.jsonl) drives every scenario:
# - Filter attributes (namespace/statefulset/cluster) are identical across
# data_sources, so the same existingQuery matches under any signal.
# - Target values (k8s.pod.name / host.name) are suffixed per data_source, so
# the returned set differs by signal -> every scenario also proves the
# data_source scoping (no cross-signal leakage).
DATASET = "queriercommon/related_values.jsonl"
@pytest.fixture(name="seed_related_values", scope="function")
def seed_related_values(
insert_attributes_metadata: Callable[[list[AttributesMetadata]], None],
insert_logs: Callable[[list[Logs]], None],
insert_traces: Callable[[list[Traces]], None],
insert_metrics: Callable[[list[Metrics]], None],
) -> datetime:
now = datetime.now(tz=UTC)
# values queried by related-values come from attributes_metadata
insert_attributes_metadata(AttributesMetadata.load_from_file(get_testdata_file_path(DATASET), timestamp=now))
# existingQuery key resolution reads each signal's own keys tables, not
# attributes_metadata; register the filter keys (namespace/statefulset/
# cluster) as resource keys for all three signals. These rows do not write
# to attributes_metadata, so they never leak into the related values.
reg_res = {
"k8s.cluster.name": "c1",
"k8s.namespace.name": "reg-ns",
"k8s.statefulset.name": "reg-sts",
}
insert_logs([Logs(timestamp=now, resources=reg_res, body="reg")])
insert_traces([Traces(timestamp=now, resources=reg_res)])
insert_metrics([Metrics(metric_name="reg_metric", timestamp=now, value=1.0, resource_attributes=reg_res)])
return now
@pytest.mark.parametrize(
"signal,name,existing_query,search_text,expected",
[
pytest.param(
"logs",
"k8s.pod.name",
"k8s.namespace.name = 'ns-a'",
"",
{"podA1-logs", "podA2-logs"},
id="positive_equal_excludes_missing_key",
),
pytest.param(
"metrics",
"k8s.pod.name",
"k8s.namespace.name IN ['ns-a', 'ns-b']",
"",
{"podA1-metrics", "podA2-metrics", "podB1-metrics"},
id="positive_in",
),
pytest.param(
"traces",
"k8s.pod.name",
"k8s.statefulset.name NOT IN ['sts-a']",
"",
{"podA2-traces", "podB1-traces", "podOrphan-traces"},
id="negative_not_in_keeps_missing_key",
),
pytest.param(
"metrics",
"k8s.pod.name",
"k8s.namespace.name != 'ns-a'",
"",
{"podB1-metrics", "podOrphan-metrics"},
id="negative_not_equal_keeps_missing_key",
),
pytest.param(
"logs",
"k8s.pod.name",
"k8s.namespace.name IN ['ns-a'] AND k8s.statefulset.name NOT IN ['sts-a']",
"",
{"podA2-logs"},
id="mixed_positive_and_negative_and",
),
pytest.param(
"traces",
"k8s.pod.name",
"k8s.statefulset.name EXISTS",
"",
{"podA1-traces"},
id="positive_exists",
),
pytest.param(
"logs",
"k8s.pod.name",
"k8s.statefulset.name NOT EXISTS",
"",
{"podA2-logs", "podB1-logs", "podOrphan-logs", "podEmpty-logs"},
id="negative_not_exists_keeps_missing_key",
),
pytest.param(
"logs",
"host.name",
"k8s.cluster.name = 'c1'",
"host",
{"host-res-logs", "host-attr-logs"},
id="dual_context_search_text_or",
),
pytest.param(
"metrics",
"host.name",
"k8s.cluster.name = 'c1'",
"host",
{"host-res-metrics", "host-attr-metrics"},
id="dual_context_search_text_signal_scoped",
),
pytest.param(
"logs",
"host.name",
"k8s.cluster.name = 'c1'",
"HOST-RES",
{"host-res-logs"},
id="search_text_case_insensitive",
),
pytest.param(
"traces",
"host.name",
"k8s.cluster.name = 'c1'",
"host",
# host.name is in both maps; the SELECT prefers the resource value,
# so only host-rest surfaces (host-attributes is shadowed).
{"host-rest"},
id="key_in_both_maps_resource_value_wins",
),
pytest.param(
"logs",
"k8s.pod.name",
"k8s.namespace.name != 'ns-a'",
"",
{"podB1-logs", "podOrphan-logs", "podEmpty-logs"},
id="empty_value_and_absent_key_both_kept_on_negative",
),
pytest.param(
None,
"k8s.pod.name",
"k8s.namespace.name = 'ns-a'",
"",
{
"podA1-logs",
"podA2-logs",
"podA1-traces",
"podA2-traces",
"podA1-metrics",
"podA2-metrics",
},
id="unspecified_signal_unions_all_data_sources",
),
],
)
def test_related_values(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
seed_related_values: datetime,
signal: str | None,
name: str,
existing_query: str,
search_text: str,
expected: set[str],
) -> None:
now = seed_related_values
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
params = {
"name": name,
"searchText": search_text,
"existingQuery": existing_query,
"startUnixMilli": int((now - timedelta(hours=1)).timestamp() * 1000),
"endUnixMilli": int((now + timedelta(hours=1)).timestamp() * 1000),
}
# omit signal entirely for the unspecified-signal scenario
if signal is not None:
params["signal"] = signal
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/fields/values"),
timeout=5,
headers={"authorization": f"Bearer {token}"},
params=params,
)
assert response.status_code == HTTPStatus.OK
assert response.json()["status"] == "success"
related = response.json()["data"]["values"].get("relatedValues") or []
assert set(related) == expected