Compare commits

...

5 Commits

Author SHA1 Message Date
srikanthccv
e28cf8aaaa feat: add semconv migration product surfaces 2026-08-08 15:56:15 +05:30
srikanthccv
a0e5c7c126 feat: resolve semconv families across logs and metrics 2026-08-08 15:28:19 +05:30
srikanthccv
78c4f5a7e4 test: add semantic convention phase one closure gate 2026-08-08 15:24:54 +05:30
srikanthccv
09bdda1280 feat: support semconv evolution in services 2026-08-08 15:19:46 +05:30
srikanthccv
7b9e9207ea feat: resolve semantic convention names in trace queries 2026-08-08 15:19:46 +05:30
99 changed files with 3525 additions and 221 deletions

View File

@@ -220,6 +220,18 @@ py-test-teardown: ## Tear down the shared SigNoz backend
py-test: ## Runs integration tests
@cd tests && uv run pytest --basetemp=./tmp/ -vv --capture=no integration/tests/
.PHONY: py-test-semconv-phase1
py-test-semconv-phase1: py-test-setup ## Rebuild the shared stack and run the semantic-convention Phase 1 matrix
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --capture=no integration/tests/queriertraces/13_semconv_evolution.py
.PHONY: py-test-semconv-phase2
py-test-semconv-phase2: py-test-setup ## Rebuild the shared stack and run the Phase 1-2 cross-signal matrices
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --capture=no integration/tests/queriertraces/13_semconv_evolution.py integration/tests/queriersemconv/02_cross_signal.py
.PHONY: py-test-semconv-phase3
py-test-semconv-phase3: py-test-setup ## Rebuild the shared stack and run the Phase 1-3 compatibility and migration-report matrices
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --capture=no integration/tests/queriertraces/13_semconv_evolution.py integration/tests/queriersemconv/02_cross_signal.py
.PHONY: py-clean
py-clean: ## Clear all pycache and pytest cache from tests directory recursively
@echo ">> cleaning python cache files from tests directory"

View File

@@ -0,0 +1,9 @@
import axios from 'api';
import { SemconvMigrationReport } from 'types/api/semconvMigration';
async function getSemconvMigrationReport(): Promise<SemconvMigrationReport> {
const response = await axios.get('/fields/semconv-migration');
return response.data.data;
}
export default getSemconvMigrationReport;

View File

@@ -5,7 +5,7 @@ import { ArrowUpRight } from '@signozhq/icons';
const QUICK_FILTER_DOC_PATHS: Record<string, string> = {
severity_text: 'severity-text',
'deployment.environment': 'environment',
'deployment.environment.name': 'environment',
'service.name': 'service-name',
'host.name': 'hostname',
'k8s.cluster.name': 'k8s-cluster-name',

View File

@@ -0,0 +1,32 @@
import { Alert } from 'antd';
import { findOldSemconvNames } from 'utils/semconv';
interface SemconvEditorWarningProps {
value: unknown;
editor: string;
}
function SemconvEditorWarning({
value,
editor,
}: SemconvEditorWarningProps): JSX.Element | null {
const text = typeof value === 'string' ? value : JSON.stringify(value ?? '');
const renames = findOldSemconvNames(text);
if (renames.length === 0) {
return null;
}
return (
<Alert
type="warning"
showIcon
data-testid="semconv-editor-warning"
message={`${editor} contains renamed OpenTelemetry fields`}
description={renames
.map(({ old, current }) => `${old}${current}`)
.join(', ')}
/>
);
}
export default SemconvEditorWarning;

View File

@@ -0,0 +1,23 @@
import { Badge } from '@signozhq/ui/badge';
import { getSemconvRename } from 'utils/semconv';
interface SemconvOldNameBadgeProps {
name: string;
}
function SemconvOldNameBadge({
name,
}: SemconvOldNameBadgeProps): JSX.Element | null {
const rename = getSemconvRename(name);
if (!rename || rename.family.kind !== 'attribute') {
return null;
}
return (
<Badge color="amber" variant="outline" data-testid="semconv-old-name-badge">
old name, renamed to {rename.current}
</Badge>
);
}
export default SemconvOldNameBadge;

View File

@@ -0,0 +1,39 @@
import { render, screen } from '@testing-library/react';
import SemconvEditorWarning from '../SemconvEditorWarning';
import SemconvOldNameBadge from '../SemconvOldNameBadge';
describe('semantic convention product hints', () => {
it('badges an old raw attribute with its current name', () => {
render(<SemconvOldNameBadge name="deployment.environment" />);
expect(screen.getByTestId('semconv-old-name-badge')).toHaveTextContent(
'old name, renamed to deployment.environment.name',
);
});
it('does not badge a current raw attribute', () => {
render(<SemconvOldNameBadge name="deployment.environment.name" />);
expect(
screen.queryByTestId('semconv-old-name-badge'),
).not.toBeInTheDocument();
});
it('shows an informational editor warning without disabling the editor', () => {
render(
<>
<input aria-label="query" defaultValue="db.system = 'postgresql'" />
<SemconvEditorWarning
value="db.system = 'postgresql'"
editor="ClickHouse SQL"
/>
</>,
);
expect(screen.getByLabelText('query')).not.toBeDisabled();
expect(screen.getByTestId('semconv-editor-warning')).toHaveTextContent(
'db.system → db.system.name',
);
});
});

View File

@@ -0,0 +1,2 @@
export { default as SemconvEditorWarning } from './SemconvEditorWarning';
export { default as SemconvOldNameBadge } from './SemconvOldNameBadge';

View File

@@ -11,12 +11,21 @@ export type SemconvFamily = {
};
export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
{
current: 'container.cpu.usage',
old: ['container.cpu.utilization'],
kind: 'metric',
contexts: ['metric'],
signals: ['metrics'],
applyToMetrics: [],
valueMap: {},
},
{
current: 'db.system.name',
old: ['db.system'],
kind: 'attribute',
contexts: [],
signals: [],
contexts: ['attribute', 'resource'],
signals: ['logs', 'metrics', 'traces'],
applyToMetrics: [],
valueMap: {},
},
@@ -24,8 +33,26 @@ export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
current: 'deployment.environment.name',
old: ['deployment.environment'],
kind: 'attribute',
contexts: [],
signals: [],
contexts: ['attribute', 'resource'],
signals: ['logs', 'metrics', 'traces'],
applyToMetrics: [],
valueMap: {},
},
{
current: 'k8s.node.cpu.usage',
old: ['k8s.node.cpu.utilization'],
kind: 'metric',
contexts: ['metric'],
signals: ['metrics'],
applyToMetrics: [],
valueMap: {},
},
{
current: 'k8s.pod.cpu.usage',
old: ['k8s.pod.cpu.utilization'],
kind: 'metric',
contexts: ['metric'],
signals: ['metrics'],
applyToMetrics: [],
valueMap: {},
},

View File

@@ -155,7 +155,7 @@ function DomainList(): JSX.Element {
dataSource={DataSource.TRACES}
queryData={query}
onChange={handleSearchChange}
placeholder="Enter your filter query (e.g., deployment.environment = 'otel-demo' AND service.name = 'frontend')"
placeholder="Enter your filter query (e.g., deployment.environment.name = 'otel-demo' AND service.name = 'frontend')"
hardcodedAttributeKeys={ApiMonitoringHardcodedAttributeKeys}
/>
</div>

View File

@@ -6,9 +6,9 @@ import { SPAN_ATTRIBUTES } from './Explorer/Domains/DomainDetails/constants';
export const ApiMonitoringHardcodedAttributeKeys: QueryKeyDataSuggestionsProps[] =
[
{
label: 'deployment.environment',
label: 'deployment.environment.name',
type: 'resource',
name: 'deployment.environment',
name: 'deployment.environment.name',
signal: 'traces',
fieldDataType: QUERY_BUILDER_KEY_TYPES.STRING,
},

View File

@@ -87,7 +87,7 @@ export const ApiMonitoringQuickFiltersConfig: IQuickFiltersConfig[] = [
title: 'Environment',
attributeKey: {
key: 'deployment.environment',
key: 'deployment.environment.name',
dataType: DataTypes.String,
type: 'resource',
},

View File

@@ -112,7 +112,7 @@ export const INFRA_MONITORING_ATTR_KEYS = {
K8S_OBJECT_NAME: 'k8s.object.name',
// Environment
DEPLOYMENT_ENVIRONMENT: 'deployment.environment',
DEPLOYMENT_ENVIRONMENT: 'deployment.environment.name',
// Host System
OS_TYPE: 'os.type',
@@ -733,7 +733,7 @@ export const ENTITY_FILTER_PLACEHOLDERS: Record<InfraMonitoringEntity, string> =
[InfraMonitoringEntity.NAMESPACES]:
"Enter your filter query (e.g., k8s.namespace.name = 'production' AND k8s.cluster.name = 'prod-cluster')",
[InfraMonitoringEntity.CLUSTERS]:
"Enter your filter query (e.g., k8s.cluster.name = 'prod-cluster' AND deployment.environment = 'production')",
"Enter your filter query (e.g., k8s.cluster.name = 'prod-cluster' AND deployment.environment.name = 'production')",
[InfraMonitoringEntity.DEPLOYMENTS]:
"Enter your filter query (e.g., k8s.deployment.name = 'api-server' AND k8s.namespace.name = 'production')",
[InfraMonitoringEntity.STATEFULSETS]:

View File

@@ -2,6 +2,17 @@
color: white;
}
.semconv-migration-report {
margin-top: 32px;
display: flex;
flex-direction: column;
gap: 12px;
.ant-table-wrapper {
margin-top: 4px;
}
}
.ingestion-key-container {
margin-top: 24px;
display: flex;

View File

@@ -5,6 +5,8 @@ import getIngestionData from 'api/settings/getIngestionData';
import { useAppContext } from 'providers/App/App';
import { IngestionDataType } from 'types/api/settings/ingestion';
import SemconvMigrationReport from './SemconvMigrationReport';
import './IngestionSettings.styles.scss';
export default function IngestionSettings(): JSX.Element {
@@ -84,6 +86,7 @@ export default function IngestionSettings(): JSX.Element {
dataSource={data}
bordered
/>
<SemconvMigrationReport />
</div>
);
}

View File

@@ -83,6 +83,8 @@ import { MeterAggregateOperator } from 'types/common/queryBuilder';
import { USER_ROLES } from 'types/roles';
import { getDaysUntilExpiry } from 'utils/timeUtils';
import SemconvMigrationReport from './SemconvMigrationReport';
import './IngestionSettings.styles.scss';
const { Option } = Select;
@@ -1705,6 +1707,7 @@ function MultiIngestionSettings(): JSX.Element {
}}
className="ingestion-keys-table"
/>
<SemconvMigrationReport />
</div>
{/* Delete Key Modal */}

View File

@@ -0,0 +1,72 @@
import { useQuery } from 'react-query';
import { Alert, Table, TableColumnsType } from 'antd';
import { Typography } from '@signozhq/ui/typography';
import getSemconvMigrationReport from 'api/semconv/getMigrationReport';
import dayjs from 'dayjs';
import { SemconvMigrationReportEntry } from 'types/api/semconvMigration';
function SemconvMigrationReport(): JSX.Element {
const { data, isLoading, isError } = useQuery({
queryKey: ['semconv-migration-report'],
queryFn: getSemconvMigrationReport,
});
const columns: TableColumnsType<SemconvMigrationReportEntry> = [
{
title: 'Old name',
dataIndex: 'old',
key: 'old',
},
{
title: 'Current name',
dataIndex: 'current',
key: 'current',
},
{
title: 'Signal',
dataIndex: 'signal',
key: 'signal',
},
{
title: 'Services still sending only the old name',
dataIndex: 'services',
key: 'services',
render: (services: string[]): string => services.join(', '),
},
{
title: 'Last seen',
dataIndex: 'lastSeenUnixMilli',
key: 'lastSeenUnixMilli',
render: (value: number): string =>
dayjs(value).format('YYYY-MM-DD HH:mm:ss'),
},
];
return (
<section className="semconv-migration-report">
<Typography.Title level={4}>Semantic convention migration</Typography.Title>
<Typography.Text>
Services in this report sent an old OpenTelemetry field during the last 24
hours without sending its current replacement. Update their SDK or
instrumentation when practical; SigNoz queries remain backward compatible.
</Typography.Text>
{isError && (
<Alert
type="error"
showIcon
message="Could not load the semantic convention migration report"
/>
)}
<Table
loading={isLoading}
columns={columns}
dataSource={data?.entries ?? []}
rowKey={(entry): string => `${entry.current}-${entry.old}-${entry.signal}`}
pagination={false}
locale={{ emptyText: 'No old-only services found in the last 24 hours' }}
/>
</section>
);
}
export default SemconvMigrationReport;

View File

@@ -15,7 +15,7 @@ export const SAMPLE_SPAN_JSON = `{
},
"resource": {
"service.name": "llm-gateway",
"deployment.environment": "production"
"deployment.environment.name": "production"
}
}`;

View File

@@ -1120,7 +1120,7 @@
"plugin": {
"kind": "signoz/QueryVariable",
"spec": {
"queryValue": "SELECT DISTINCT resources_string['deployment.environment'] AS environment FROM signoz_traces.distributed_signoz_index_v3 WHERE mapContains(resources_string, 'deployment.environment') AND timestamp >= now() - INTERVAL 1 DAY"
"queryValue": "SELECT DISTINCT resources_string['deployment.environment.name'] AS environment FROM signoz_traces.distributed_signoz_index_v3 WHERE mapContains(resources_string, 'deployment.environment.name') AND timestamp >= now() - INTERVAL 1 DAY"
}
}
}

View File

@@ -1,6 +1,7 @@
import { Divider } from '@signozhq/ui/divider';
import { TooltipSimple } from '@signozhq/ui/tooltip';
import { Typography } from '@signozhq/ui/typography';
import { SemconvOldNameBadge } from 'components/Semconv';
import { TagContainer, TagLabel, TagValue } from './FieldRenderer.styles';
import { FieldRendererProps } from './LogDetailedView.types';
@@ -28,6 +29,7 @@ function FieldRenderer({ field }: FieldRendererProps): JSX.Element {
<Typography.Text truncate={1} className="label">
{newField}{' '}
</Typography.Text>
<SemconvOldNameBadge name={newField} />
</TooltipSimple>
<div className="tags">
@@ -47,7 +49,10 @@ function FieldRenderer({ field }: FieldRendererProps): JSX.Element {
</div>
</>
) : (
<span className="label">{field}</span>
<>
<span className="label">{field}</span>
<SemconvOldNameBadge name={field} />
</>
)}
</span>
);

View File

@@ -164,7 +164,9 @@ describe('useInitialQuery - Priority-Based Resource Filtering', () => {
value: 'frontend-service',
}),
expect.objectContaining({
key: expect.objectContaining({ key: 'deployment.environment' }),
key: expect.objectContaining({
key: 'deployment.environment.name',
}),
value: 'production',
}),
expect.objectContaining({
@@ -286,7 +288,9 @@ describe('useInitialQuery - Priority-Based Resource Filtering', () => {
value: 'legacy-app',
}),
expect.objectContaining({
key: expect.objectContaining({ key: 'deployment.environment' }),
key: expect.objectContaining({
key: 'deployment.environment.name',
}),
value: 'production',
}),
expect.objectContaining({

View File

@@ -6,13 +6,14 @@ import {
TagFilterItem,
} from 'types/api/queryBuilder/queryBuilderData';
import { v4 as uuid } from 'uuid';
import { getSemconvRename } from 'utils/semconv';
const FALLBACK_STARTS_WITH_REGEX = /^(k8s|cloud|host|deployment)/; // regex to filter out resources that start with the specified keywords
const FALLBACK_CONTAINS_REGEX = /(env|service|file|container|tenant)/; // regex to filter out resources that contains the specified keywords
// Priority categories for filter selection
// Strategy:
// - Always include: service.name, deployment.environment, env, environment
// - Always include: service.name, deployment.environment.name, env, environment
// - Select ONE category only: stops at the first category with a matching attribute
// - Within category: picks the first available attribute by order
// - Order (highest to lowest priority): Kubernetes > Cloud > Host > Container
@@ -26,27 +27,36 @@ const PRIORITY_CATEGORIES = [
const SERVICE_AND_ENVIRONMENT_KEYS = [
'service.name',
'deployment.environment',
'deployment.environment.name',
'env',
'environment',
];
export const getFiltersFromResources = (
resources: ILog['resources_string'],
): TagFilterItem[] =>
Object.keys(resources).map((key: string) => {
): TagFilterItem[] => {
const items = new Map<string, TagFilterItem>();
Object.keys(resources).forEach((key: string) => {
const currentKey = getSemconvRename(key)?.current ?? key;
const resourceValue = resources[key] as string;
return {
const item = {
id: uuid(),
key: {
key,
key: currentKey,
dataType: DataTypes.String,
type: 'resource',
},
op: OPERATORS['='],
value: resourceValue,
};
// If raw data contains both names, retain the current value just like the
// backend's current-first resolver.
if (!items.has(currentKey) || key === currentKey) {
items.set(currentKey, item);
}
});
return Array.from(items.values());
};
export const isServiceOrEnvironmentAttribute = (key: string): boolean =>
SERVICE_AND_ENVIRONMENT_KEYS.includes(key);

View File

@@ -94,7 +94,7 @@ function DBCall(): JSX.Element {
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const legend = dotMetricsEnabled ? '{{db.system}}' : '{{db_system}}';
const legend = dotMetricsEnabled ? '{{db.system.name}}' : '{{db_system_name}}';
const databaseCallsRPSWidget = useMemo(
() =>

View File

@@ -28,7 +28,7 @@ import { v4 as uuid } from 'uuid';
export const dbSystemTags: Tags[] = [
{
Key: 'db.system.(string)',
Key: 'db.system.name.(string)',
StringValues: [''],
NumberValues: [],
BoolValues: [],

View File

@@ -103,7 +103,7 @@ export enum WidgetKeys {
SignozExternalCallLatencySum = 'signoz_external_call_latency_sum',
Signoz_latency_bucket_norm = 'signoz_latency_bucket',
Signoz_latency_bucket = 'signoz_latency.bucket',
Db_system = 'db.system',
Db_system = 'db.system.name',
Db_system_norm = 'db_system',
}

View File

@@ -2,6 +2,7 @@ import { ChangeEvent, useCallback } from 'react';
import MEditor, { Monaco } from '@monaco-editor/react';
import { Color } from '@signozhq/design-tokens';
import { Input } from 'antd';
import { SemconvEditorWarning } from 'components/Semconv';
import { LEGEND } from 'constants/global';
import { useQueryBuilder } from 'hooks/queryBuilder/useQueryBuilder';
import { useIsDarkMode } from 'hooks/useDarkMode';
@@ -118,6 +119,7 @@ function ClickHouseQueryBuilder({
theme={isDarkMode ? 'my-theme' : 'light'}
beforeMount={setEditorTheme}
/>
<SemconvEditorWarning value={queryData?.query} editor="ClickHouse SQL" />
<Input
onChange={handleUpdateInput}
name="legend"

View File

@@ -1,6 +1,7 @@
import { ChangeEvent, useCallback } from 'react';
import { Input } from 'antd';
import { LEGEND } from 'constants/global';
import { SemconvEditorWarning } from 'components/Semconv';
import { useQueryBuilder } from 'hooks/queryBuilder/useQueryBuilder';
import { IPromQLQuery } from 'types/api/queryBuilder/queryBuilderData';
import { EQueryType } from 'types/common/dashboard';
@@ -66,6 +67,7 @@ function PromQLQueryBuilder({
style={{ marginBottom: '0.5rem' }}
data-testid="promql-query-input"
/>
<SemconvEditorWarning value={queryData?.query} editor="PromQL" />
<Input
onChange={handleUpdateQuery}

View File

@@ -1,6 +1,7 @@
import { useTranslation } from 'react-i18next';
import { Form } from 'antd';
import { initialQueryBuilderFormValuesMap } from 'constants/queryBuilder';
import { SemconvEditorWarning } from 'components/Semconv';
import QueryBuilderSearchV2 from 'container/QueryBuilder/filters/QueryBuilderSearchV2/QueryBuilderSearchV2';
import isEqual from 'lodash-es/isEqual';
import { TagFilter } from 'types/api/queryBuilder/queryBuilderData';
@@ -55,6 +56,7 @@ function TagFilterInputWithLogsResultPreview({
value={value}
onChange={onChange}
/>
<SemconvEditorWarning value={value} editor="Pipeline filter" />
<div className="pipeline-filter-input-preview-container">
<LogsFilterPreview filter={value} />
</div>

View File

@@ -424,7 +424,7 @@ describe('ResourceProvider', () => {
await waitFor(() => {
expect(result.current.queries).toHaveLength(1);
expect(result.current.queries[0]).toMatchObject({
tagKey: 'resource_deployment_environment',
tagKey: 'resource_deployment_environment_name',
operator: 'IN',
tagValue: ['production'],
});
@@ -435,7 +435,7 @@ describe('ResourceProvider', () => {
const seeded = [
{
id: 'env',
tagKey: 'resource_deployment_environment',
tagKey: 'resource_deployment_environment_name',
operator: 'IN',
tagValue: ['production'],
},
@@ -459,7 +459,7 @@ describe('ResourceProvider', () => {
await waitFor(() => {
const tagKeys = result.current.queries.map((q) => q.tagKey);
expect(tagKeys).not.toContain('resource_deployment_environment');
expect(tagKeys).not.toContain('resource_deployment_environment_name');
expect(tagKeys).toContain('resource_service_name');
});
});
@@ -468,7 +468,7 @@ describe('ResourceProvider', () => {
const seeded = [
{
id: 'env',
tagKey: 'resource_deployment_environment',
tagKey: 'resource_deployment_environment_name',
operator: 'IN',
tagValue: ['production'],
},
@@ -486,7 +486,7 @@ describe('ResourceProvider', () => {
await waitFor(() => {
const envQueries = result.current.queries.filter(
(q) => q.tagKey === 'resource_deployment_environment',
(q) => q.tagKey === 'resource_deployment_environment_name',
);
expect(envQueries).toHaveLength(1);
expect(envQueries[0].tagValue).toStrictEqual(['staging']);
@@ -518,7 +518,7 @@ describe('ResourceProvider', () => {
await waitFor(() => {
expect(result.current.queries[0].tagKey).toBe(
'resource_deployment.environment',
'resource_deployment.environment.name',
);
});
});

View File

@@ -6,13 +6,13 @@ import { mappingWithRoutesAndKeys } from '../utils';
describe('useResourceAttribute config', () => {
describe('whilelistedKeys', () => {
it('should include underscore-notation keys (DOT_METRICS_ENABLED=false)', () => {
expect(whilelistedKeys).toContain('resource_deployment_environment');
expect(whilelistedKeys).toContain('resource_deployment_environment_name');
expect(whilelistedKeys).toContain('resource_k8s_cluster_name');
expect(whilelistedKeys).toContain('resource_k8s_cluster_namespace');
});
it('should include dot-notation keys (DOT_METRICS_ENABLED=true)', () => {
expect(whilelistedKeys).toContain('resource_deployment.environment');
expect(whilelistedKeys).toContain('resource_deployment.environment.name');
expect(whilelistedKeys).toContain('resource_k8s.cluster.name');
expect(whilelistedKeys).toContain('resource_k8s.cluster.namespace');
});
@@ -21,8 +21,8 @@ describe('useResourceAttribute config', () => {
describe('mappingWithRoutesAndKeys', () => {
const dotNotationFilters = [
{
label: 'deployment.environment',
value: 'resource_deployment.environment',
label: 'deployment.environment.name',
value: 'resource_deployment.environment.name',
},
{ label: 'k8s.cluster.name', value: 'resource_k8s.cluster.name' },
{ label: 'k8s.cluster.namespace', value: 'resource_k8s.cluster.namespace' },
@@ -30,8 +30,8 @@ describe('useResourceAttribute config', () => {
const underscoreNotationFilters = [
{
label: 'deployment.environment',
value: 'resource_deployment_environment',
label: 'deployment.environment.name',
value: 'resource_deployment_environment_name',
},
{ label: 'k8s.cluster.name', value: 'resource_k8s_cluster_name' },
{ label: 'k8s.cluster.namespace', value: 'resource_k8s_cluster_namespace' },

View File

@@ -1,6 +1,6 @@
export const whilelistedKeys = [
'resource_deployment_environment',
'resource_deployment.environment',
'resource_deployment_environment_name',
'resource_deployment.environment.name',
'resource_k8s_cluster_name',
'resource_k8s.cluster.name',
'resource_k8s_cluster_namespace',

View File

@@ -148,9 +148,9 @@ export const getResourceDeploymentKeys = (
dotMetricsEnabled: boolean,
): string => {
if (dotMetricsEnabled) {
return 'resource_deployment.environment';
return 'resource_deployment.environment.name';
}
return 'resource_deployment_environment';
return 'resource_deployment_environment_name';
};
export const GetTagKeys = async (

View File

@@ -40,7 +40,7 @@ export const LogsQuickFiltersConfig: IQuickFiltersConfig[] = [
type: FiltersType.CHECKBOX,
title: 'Environment',
attributeKey: {
key: 'deployment.environment',
key: 'deployment.environment.name',
dataType: DataTypes.String,
type: 'resource',
},

View File

@@ -11,6 +11,7 @@ import { Skeleton } from 'antd';
import { DetailsHeader, DetailsPanelDrawer } from 'components/DetailsPanel';
import { HeaderAction } from 'components/DetailsPanel/DetailsHeader/DetailsHeader';
import { DetailsPanelState } from 'components/DetailsPanel/types';
import { SemconvOldNameBadge } from 'components/Semconv';
import { QueryParams } from 'constants/query';
import {
initialQueryBuilderFormValuesMap,
@@ -108,6 +109,12 @@ function SpanDetailsContent({
() => getSpanDisplayData(selectedSpan),
[selectedSpan],
);
const semconvLabelSuffix = useCallback(
(fieldKey: string): React.ReactNode => (
<SemconvOldNameBadge name={fieldKey} />
),
[],
);
// Map span attribute actions to PrettyView actions format.
// Use the last key in fieldKeyPath (the actual attribute key), not the full display path.
@@ -329,6 +336,7 @@ function SpanDetailsContent({
visibleActions: VISIBLE_ACTIONS,
pinnedFieldsValue,
onPinnedFieldsChange,
labelSuffixRenderer: semconvLabelSuffix,
}}
/>
</TabsContent>

View File

@@ -29,7 +29,7 @@ export const KEY_ATTRIBUTE_KEYS: Record<string, string[]> = {
traces: [
'service.name',
'service.namespace',
'deployment.environment',
'deployment.environment.name',
'timestamp',
'duration_nano',
'kind_string',

View File

@@ -22,7 +22,7 @@ export const SPAN_CATEGORIES: readonly SpanCategory[] = [
// Map each category to the attribute key it filters on
const CATEGORY_KEYS: Record<Exclude<SpanCategory, 'All'>, string> = {
Database: 'db.system',
Database: 'db.system.name',
HTTP: 'http.method',
Functions: 'kind_string',
Jobs: 'messaging.system',
@@ -34,7 +34,7 @@ const ALL_CATEGORY_KEYS = Object.values(CATEGORY_KEYS);
// The expression clause to add for each category
const CATEGORY_EXPRESSIONS: Record<Exclude<SpanCategory, 'All'>, string> = {
Database: 'db.system exists',
Database: 'db.system.name exists',
HTTP: 'http.method exists',
Functions: "kind_string = 'Internal'",
Jobs: 'messaging.system exists',

View File

@@ -38,7 +38,7 @@ export function Section(props: SectionProps): JSX.Element {
'hasError',
'durationNano',
'serviceName',
'deployment.environment',
'deployment.environment.name',
]),
),
[selectedFilters],

View File

@@ -14,7 +14,7 @@ export const AllTraceFilterKeyValue: Record<string, string> = {
durationNano: 'Duration',
duration_nano: 'Duration',
durationNanoMax: 'Duration',
'deployment.environment': 'Environment',
'deployment.environment.name': 'Environment',
hasError: 'Status',
has_error: 'Status',
serviceName: 'Service Name',
@@ -208,11 +208,11 @@ export const traceFilterKeys: Record<AllTraceFilterKeys, BaseAutocompleteData> =
id: 'serviceName--string--tag--true',
},
'deployment.environment': {
key: 'deployment.environment',
'deployment.environment.name': {
key: 'deployment.environment.name',
dataType: DataTypes.String,
type: 'resource',
id: 'deployment.environment--string--resource--false',
id: 'deployment.environment.name--string--resource--false',
},
name: {
key: 'name',

View File

@@ -223,6 +223,12 @@
padding-left: 6px !important;
}
&__label {
display: inline-flex;
align-items: baseline;
gap: 6px;
}
&__pinned-icon {
flex-shrink: 0;
color: var(--text-robin-400);

View File

@@ -67,6 +67,7 @@ export interface PrettyViewProps {
*/
pinnedFieldsValue?: string[];
onPinnedFieldsChange?: (next: string[]) => void;
labelSuffixRenderer?: (fieldKey: string) => React.ReactNode;
}
function PrettyView({
@@ -78,6 +79,7 @@ function PrettyView({
drawerKey = 'default',
pinnedFieldsValue,
onPinnedFieldsChange,
labelSuffixRenderer,
}: PrettyViewProps): JSX.Element {
const isDarkMode = useIsDarkMode();
const [, setCopy] = useCopyToClipboard();
@@ -305,10 +307,24 @@ function PrettyView({
}}
/>
<span>{displayKey}</span>
{labelSuffixRenderer?.(displayKey)}
</span>
);
},
[togglePin, pinnedEntries],
[togglePin, pinnedEntries, labelSuffixRenderer],
);
const labelRenderer = useCallback(
(keyPath: KeyPath): React.ReactNode => {
const displayKey = String(keyPath[0]);
return (
<span className="pretty-view__label">
<span>{displayKey}</span>
{labelSuffixRenderer?.(displayKey)}
</span>
);
},
[labelSuffixRenderer],
);
return (
@@ -351,6 +367,7 @@ function PrettyView({
shouldExpandNodeInitially={shouldExpandNodeInitially}
valueRenderer={valueRenderer}
getItemString={getItemString}
labelRenderer={labelRenderer}
/>
</div>
);

View File

@@ -0,0 +1,14 @@
export interface SemconvMigrationReportEntry {
current: string;
old: string;
signal: string;
services: string[];
resourceSets: number;
lastSeenUnixMilli: number;
}
export interface SemconvMigrationReport {
startUnixMilli: number;
endUnixMilli: number;
entries: SemconvMigrationReportEntry[];
}

View File

@@ -0,0 +1,29 @@
import { findOldSemconvNames, getSemconvRename } from 'utils/semconv';
describe('semantic convention helpers', () => {
it('returns the current name for an old attribute', () => {
expect(getSemconvRename('deployment.environment')).toMatchObject({
old: 'deployment.environment',
current: 'deployment.environment.name',
});
});
it('finds old names in editor text without matching larger custom names', () => {
expect(
findOldSemconvNames(
"deployment.environment = 'prod' AND custom.db.system.value = 'x'",
),
).toStrictEqual([
expect.objectContaining({
old: 'deployment.environment',
current: 'deployment.environment.name',
}),
]);
});
it('does not warn for current names', () => {
expect(
findOldSemconvNames('deployment.environment.name = prod'),
).toStrictEqual([]);
});
});

View File

@@ -0,0 +1,50 @@
import {
SEMCONV_FAMILIES,
SemconvFamily,
} from 'constants/generated/semconvFamilies.gen';
export type SemconvRename = {
old: string;
current: string;
family: SemconvFamily;
};
const OLD_NAMES = SEMCONV_FAMILIES.flatMap((family) =>
family.old.map((old) => ({ old, current: family.current, family })),
);
const OLD_NAME_INDEX = new Map(OLD_NAMES.map((rename) => [rename.old, rename]));
export function getSemconvRename(name: string): SemconvRename | undefined {
return OLD_NAME_INDEX.get(name);
}
export function findOldSemconvNames(text: string): SemconvRename[] {
if (!text) {
return [];
}
return OLD_NAMES.filter(({ old }) => containsSemconvName(text, old));
}
function containsSemconvName(text: string, name: string): boolean {
let offset = 0;
while (offset < text.length) {
const index = text.indexOf(name, offset);
if (index === -1) {
return false;
}
const before = index === 0 ? '' : text[index - 1];
const afterIndex = index + name.length;
const after = afterIndex === text.length ? '' : text[afterIndex];
if (!isSemconvNameCharacter(before) && !isSemconvNameCharacter(after)) {
return true;
}
offset = index + 1;
}
return false;
}
function isSemconvNameCharacter(value: string): boolean {
return /[A-Za-z0-9_.-]/.test(value);
}

View File

@@ -46,5 +46,23 @@ func (provider *provider) addFieldsRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v1/fields/semconv-migration", handler.New(provider.authzMiddleware.ViewAccess(provider.fieldsHandler.GetSemconvMigrationReport), handler.OpenAPIDef{
ID: "GetSemconvMigrationReport",
Tags: []string{"fields"},
Summary: "Get semantic-convention migration report",
Description: "Returns services that still emit old semantic-convention names without the current family name",
Request: nil,
RequestQuery: new(telemetrytypes.PostableSemconvMigrationReportParams),
RequestContentType: "",
Response: new(telemetrytypes.GettableSemconvMigrationReport),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
})).Methods(http.MethodGet).GetError(); err != nil {
return err
}
return nil
}

View File

@@ -8,4 +8,7 @@ type Handler interface {
// Gets the fields values for the given field value selector
GetFieldsValues(http.ResponseWriter, *http.Request)
// Gets services that still emit only historical semantic-convention names.
GetSemconvMigrationReport(http.ResponseWriter, *http.Request)
}

View File

@@ -2,6 +2,7 @@ package implfields
import (
"net/http"
"time"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/http/binding"
@@ -16,6 +17,43 @@ type handler struct {
telemetryMetadataStore telemetrytypes.MetadataStore
}
func (handler *handler) GetSemconvMigrationReport(rw http.ResponseWriter, req *http.Request) {
ctx := req.Context()
var params telemetrytypes.PostableSemconvMigrationReportParams
if err := binding.Query.BindQuery(req.URL.Query(), &params); err != nil {
render.Error(rw, err)
return
}
now := time.Now()
if params.EndUnixMilli == 0 {
params.EndUnixMilli = now.UnixMilli()
}
if params.StartUnixMilli == 0 {
params.StartUnixMilli = now.Add(-24 * time.Hour).UnixMilli()
}
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
report, err := handler.telemetryMetadataStore.GetSemconvMigrationReport(
ctx,
valuer.MustNewUUID(claims.OrgID),
params.StartUnixMilli,
params.EndUnixMilli,
)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, report)
}
func NewHandler(settings factory.ProviderSettings, telemetryMetadataStore telemetrytypes.MetadataStore) fields.Handler {
return &handler{
telemetryMetadataStore: telemetryMetadataStore,

View File

@@ -786,10 +786,11 @@ func (q *querier) run(
Results: maps.Values(processedResults),
},
Meta: qbtypes.ExecStats{
RowsScanned: stats.RowsScanned,
BytesScanned: stats.BytesScanned,
DurationMS: stats.DurationMS,
StepIntervals: stepIntervals,
RowsScanned: stats.RowsScanned,
BytesScanned: stats.BytesScanned,
DurationMS: stats.DurationMS,
StepIntervals: stepIntervals,
SemconvResolutions: semconvResolutionsForRequest(req),
},
}

View File

@@ -0,0 +1,143 @@
package querier
import (
"encoding/json"
"slices"
"strings"
"unicode"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// semconvResolutionsForRequest reports only query-builder resolutions. Raw
// ClickHouse SQL and PromQL are deliberately excluded: SigNoz does not rewrite
// those languages and their editors surface non-blocking warnings instead.
func semconvResolutionsForRequest(req *qbtypes.QueryRangeRequest) []qbtypes.SemconvResolution {
if req == nil {
return nil
}
resolutions := make([]qbtypes.SemconvResolution, 0)
seen := make(map[string]struct{})
for _, envelope := range req.CompositeQuery.Queries {
signal, applies := semconvResolutionSignal(envelope.Spec)
if !applies {
continue
}
payload, err := json.Marshal(envelope.Spec)
if err != nil {
continue
}
text := string(payload)
for _, family := range semconv.All() {
if _, ok := semconv.Lookup(family.Kind, telemetrytypes.FieldKeySelector{
Name: family.Current,
Signal: signal,
}); !ok {
continue
}
for _, requested := range semconvRequestSpellings(family, signal) {
if !containsSemconvName(text, requested) {
continue
}
identity := family.Kind.StringValue() + "\x00" + requested + "\x00" + family.Current
if _, ok := seen[identity]; ok {
continue
}
seen[identity] = struct{}{}
resolutions = append(resolutions, qbtypes.SemconvResolution{
Requested: requested,
Current: family.Current,
Members: append([]string{family.Current}, family.Old...),
Kind: family.Kind.StringValue(),
})
}
}
}
if len(resolutions) == 0 {
return nil
}
return resolutions
}
func semconvResolutionSignal(spec any) (telemetrytypes.Signal, bool) {
switch spec.(type) {
case qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation], qbtypes.QueryBuilderTraceOperator:
return telemetrytypes.SignalTraces, true
case qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]:
return telemetrytypes.SignalLogs, true
case qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]:
return telemetrytypes.SignalMetrics, true
case qbtypes.QueryBuilderJoin:
// Joins can contain more than one signal. An unspecified signal keeps
// family signal scopes in force while allowing every applicable family.
return telemetrytypes.SignalUnspecified, true
default:
return telemetrytypes.SignalUnspecified, false
}
}
func semconvRequestSpellings(family semconv.Family, signal telemetrytypes.Signal) []string {
spellings := append([]string{family.Current}, family.Old...)
if signal != telemetrytypes.SignalMetrics {
return spellings
}
logical := slices.Clone(spellings)
for _, name := range logical {
normalized := strings.ReplaceAll(name, ".", "_")
if !slices.Contains(spellings, normalized) {
spellings = append(spellings, normalized)
}
if family.Kind == semconv.KindAttribute {
for _, resourceName := range []string{"resource_" + name, "resource_" + normalized} {
if !slices.Contains(spellings, resourceName) {
spellings = append(spellings, resourceName)
}
}
}
}
return spellings
}
func containsSemconvName(text, name string) bool {
for offset := 0; offset < len(text); {
index := strings.Index(text[offset:], name)
if index < 0 {
return false
}
start := offset + index
end := start + len(name)
if hasSemconvNameStartBoundary(text, start) &&
(end == len(text) || !isSemconvNameRune(rune(text[end]))) {
return true
}
offset = start + 1
}
return false
}
func hasSemconvNameStartBoundary(text string, start int) bool {
if start == 0 || !isSemconvNameRune(rune(text[start-1])) {
return true
}
for _, qualifier := range []string{"resource.", "attribute.", "tag.", "point."} {
qualifierStart := start - len(qualifier)
if qualifierStart >= 0 && text[qualifierStart:start] == qualifier &&
(qualifierStart == 0 || !isSemconvNameRune(rune(text[qualifierStart-1]))) {
return true
}
}
return false
}
func isSemconvNameRune(r rune) bool {
return unicode.IsLetter(r) || unicode.IsDigit(r) || r == '_' || r == '.' || r == '-'
}

View File

@@ -0,0 +1,99 @@
package querier
import (
"encoding/json"
"testing"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestSemconvResolutionsReportsOldTraceAttribute(t *testing.T) {
req := &qbtypes.QueryRangeRequest{
CompositeQuery: qbtypes.CompositeQuery{Queries: []qbtypes.QueryEnvelope{{
Spec: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
Filter: &qbtypes.Filter{Expression: "deployment.environment = 'prod'"},
},
}}},
}
assert.Equal(t, []qbtypes.SemconvResolution{{
Requested: "deployment.environment",
Current: "deployment.environment.name",
Members: []string{"deployment.environment.name", "deployment.environment"},
Kind: "attribute",
}}, semconvResolutionsForRequest(req), "old trace attribute should be reported as a family resolution")
}
func TestSemconvResolutionsReportsContextQualifiedOldAttribute(t *testing.T) {
var req qbtypes.QueryRangeRequest
err := json.Unmarshal([]byte(`{
"start": 1,
"end": 2,
"requestType": "raw",
"compositeQuery": {
"queries": [{
"type": "builder_query",
"spec": {
"signal": "traces",
"name": "A",
"filter": {"expression": "resource.deployment.environment EXISTS"}
}
}]
}
}`), &req)
require.NoError(t, err)
assert.Equal(t, []qbtypes.SemconvResolution{{
Requested: "deployment.environment",
Current: "deployment.environment.name",
Members: []string{"deployment.environment.name", "deployment.environment"},
Kind: "attribute",
}}, semconvResolutionsForRequest(&req), "qualified old trace attribute should be reported after request decoding")
}
func TestSemconvResolutionsReportsCurrentLogAttribute(t *testing.T) {
req := &qbtypes.QueryRangeRequest{
CompositeQuery: qbtypes.CompositeQuery{Queries: []qbtypes.QueryEnvelope{{
Spec: qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{
Signal: telemetrytypes.SignalLogs,
SelectFields: []telemetrytypes.TelemetryFieldKey{{
Name: "db.system.name",
}},
},
}}},
}
assert.Equal(t, []qbtypes.SemconvResolution{{
Requested: "db.system.name",
Current: "db.system.name",
Members: []string{"db.system.name", "db.system"},
Kind: "attribute",
}}, semconvResolutionsForRequest(req), "current log attribute should identify its complete family")
}
func TestSemconvResolutionsIgnoresRawSQL(t *testing.T) {
req := &qbtypes.QueryRangeRequest{
CompositeQuery: qbtypes.CompositeQuery{Queries: []qbtypes.QueryEnvelope{{
Spec: qbtypes.ClickHouseQuery{Query: "SELECT attributes_string['deployment.environment']"},
}}},
}
assert.Empty(t, semconvResolutionsForRequest(req), "raw SQL is not rewritten and should not report a resolution")
}
func TestSemconvResolutionsRequiresNameBoundary(t *testing.T) {
req := &qbtypes.QueryRangeRequest{
CompositeQuery: qbtypes.CompositeQuery{Queries: []qbtypes.QueryEnvelope{{
Spec: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
Filter: &qbtypes.Filter{Expression: "custom.deployment.environment = 'prod'"},
},
}}},
}
assert.Empty(t, semconvResolutionsForRequest(req), "a family name embedded in a larger custom key must not match")
}

View File

@@ -52,10 +52,10 @@ import (
"github.com/SigNoz/signoz/pkg/query-service/constants"
chErrors "github.com/SigNoz/signoz/pkg/query-service/errors"
"github.com/SigNoz/signoz/pkg/query-service/metrics"
"github.com/SigNoz/signoz/pkg/query-service/model"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/query-service/utils"
"github.com/SigNoz/signoz/pkg/semconv"
)
const (
@@ -3190,11 +3190,6 @@ func (r *ClickHouseReader) GetMetricAttributeValues(ctx context.Context, orgID v
var rows driver.Rows
var attributeValues v3.FilterAttributeValueResponse
normalized := true
if constants.IsDotMetricsEnabled {
normalized = false
}
reductionEnabled := r.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableMetricsReduction, featuretypes.NewFlaggerEvaluationContext(orgID))
if reductionEnabled {
@@ -3205,8 +3200,7 @@ func (r *ClickHouseReader) GetMetricAttributeValues(ctx context.Context, orgID v
if req.Limit != 0 {
query = query + fmt.Sprintf(" LIMIT %d;", req.Limit)
}
names := []string{req.AggregateAttribute}
names = append(names, metrics.GetTransitionedMetric(req.AggregateAttribute, normalized))
names := semconv.MetricNames(req.AggregateAttribute)
rows, err = r.db.Query(ctx, query, req.FilterAttributeKey, names, req.FilterAttributeKey, fmt.Sprintf("%%%s%%", req.SearchText), common.PastDayRoundOff())

View File

@@ -6,6 +6,8 @@ import (
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/query-service/utils"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
var resourceLogOperators = map[v3.FilterOperator]string{
@@ -29,13 +31,61 @@ var resourceLogOperators = map[v3.FilterOperator]string{
v3.FilterOperatorNotILike: "NOT ILIKE",
}
func resourceSemconvMembers(key string) []string {
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: key,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
})
}
func resourceValueExpression(key string) string {
members := resourceSemconvMembers(key)
if len(members) == 1 {
return fmt.Sprintf("simpleJSONExtractString(labels, '%s')", key)
}
values := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(simpleJSONExtractString(labels, '%s'), '')", member))
}
return "COALESCE(" + strings.Join(values, ", ") + ")"
}
func resourcePresenceExpression(key string, exists bool) string {
members := resourceSemconvMembers(key)
if len(members) == 1 {
if exists {
return fmt.Sprintf("simpleJSONHas(labels, '%s')", key)
}
return fmt.Sprintf("not simpleJSONHas(labels, '%s')", key)
}
conditions := make([]string, 0, len(members))
for _, member := range members {
if exists {
conditions = append(conditions, fmt.Sprintf("simpleJSONHas(labels, '%s')", member))
} else {
conditions = append(conditions, fmt.Sprintf("not simpleJSONHas(labels, '%s')", member))
}
}
separator := " OR "
if !exists {
separator = " AND "
}
return "(" + strings.Join(conditions, separator) + ")"
}
// buildResourceFilter builds a clickhouse filter string for resource labels
func buildResourceFilter(logsOp string, key string, op v3.FilterOperator, value interface{}) string {
// for all operators except contains and like
searchKey := fmt.Sprintf("simpleJSONExtractString(labels, '%s')", key)
searchKey := resourceValueExpression(key)
// for contains and like it will be case insensitive
lowerSearchKey := fmt.Sprintf("simpleJSONExtractString(lower(labels), '%s')", key)
if len(resourceSemconvMembers(key)) > 1 {
lowerSearchKey = "lower(" + searchKey + ")"
}
chFmtVal := utils.ClickHouseFormattedValue(value)
@@ -43,9 +93,9 @@ func buildResourceFilter(logsOp string, key string, op v3.FilterOperator, value
switch op {
case v3.FilterOperatorExists:
return fmt.Sprintf("simpleJSONHas(labels, '%s')", key)
return resourcePresenceExpression(key, true)
case v3.FilterOperatorNotExists:
return fmt.Sprintf("not simpleJSONHas(labels, '%s')", key)
return resourcePresenceExpression(key, false)
case v3.FilterOperatorRegex, v3.FilterOperatorNotRegex:
return fmt.Sprintf(logsOp, searchKey, chFmtVal)
case v3.FilterOperatorContains, v3.FilterOperatorNotContains:
@@ -110,6 +160,38 @@ func buildIndexFilterForInOperator(key string, op v3.FilterOperator, value inter
// we can use lower index for =, in etc but it's difficult to do it for !=, NIN etc
// if as x != "ABC" we cannot predict something like "not lower(labels) like '%%x%%abc%%'". It has it be "not lower(labels) like '%%x%%ABC%%'"
func buildResourceIndexFilter(key string, op v3.FilterOperator, value interface{}) string {
return buildResourceIndexFilterForKey(key, op, value, true)
}
func buildResourceIndexFilterForKey(key string, op v3.FilterOperator, value interface{}, resolveFamily bool) string {
members := []string{key}
if resolveFamily {
members = resourceSemconvMembers(key)
}
if len(members) > 1 {
switch op {
case v3.FilterOperatorNotEqual,
v3.FilterOperatorNotLike,
v3.FilterOperatorNotILike,
v3.FilterOperatorNotContains,
v3.FilterOperatorNotExists,
v3.FilterOperatorNotRegex,
v3.FilterOperatorNotIn:
return ""
}
conditions := make([]string, 0, len(members))
for _, member := range members {
if condition := buildResourceIndexFilterForKey(member, op, value, false); condition != "" {
conditions = append(conditions, condition)
}
}
if len(conditions) == 0 {
return ""
}
return "(" + strings.Join(conditions, " OR ") + ")"
}
// not using clickhouseFormattedValue as we don't wan't the quotes
strVal := fmt.Sprintf("%s", value)
fmtValEscapedForContains := utils.QuoteEscapedStringForContains(strVal, true)
@@ -206,14 +288,31 @@ func buildResourceFiltersFromGroupBy(groupBy []v3.AttributeKey) []string {
if attr.Type != v3.AttributeKeyTypeResource {
continue
}
conditions = append(conditions, fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", attr.Key, attr.Key))
members := resourceSemconvMembers(attr.Key)
if len(members) == 1 {
conditions = append(conditions, fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", attr.Key, attr.Key))
continue
}
indexConditions := make([]string, 0, len(members))
for _, member := range members {
indexConditions = append(indexConditions, fmt.Sprintf("labels like '%%%s%%'", member))
}
conditions = append(conditions, fmt.Sprintf("(%s AND (%s))", resourcePresenceExpression(attr.Key, true), strings.Join(indexConditions, " OR ")))
}
return conditions
}
func buildResourceFiltersFromAggregateAttribute(aggregateAttribute v3.AttributeKey) string {
if aggregateAttribute.Key != "" && aggregateAttribute.Type == v3.AttributeKeyTypeResource {
return fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", aggregateAttribute.Key, aggregateAttribute.Key)
members := resourceSemconvMembers(aggregateAttribute.Key)
if len(members) == 1 {
return fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", aggregateAttribute.Key, aggregateAttribute.Key)
}
indexConditions := make([]string, 0, len(members))
for _, member := range members {
indexConditions = append(indexConditions, fmt.Sprintf("labels like '%%%s%%'", member))
}
return fmt.Sprintf("(%s AND (%s))", resourcePresenceExpression(aggregateAttribute.Key, true), strings.Join(indexConditions, " OR "))
}
return ""

View File

@@ -5,6 +5,8 @@ import (
"testing"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func Test_buildResourceFilter(t *testing.T) {
@@ -552,3 +554,38 @@ func Test_buildResourceSubQuery(t *testing.T) {
})
}
}
func TestSemanticConventionResourceFamily(t *testing.T) {
const resolvedValue = "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''))"
for _, requestedName := range []string{"deployment.environment.name", "deployment.environment"} {
t.Run(requestedName, func(t *testing.T) {
assert.Equal(t, resolvedValue+" = 'production'", buildResourceFilter("=", requestedName, v3.FilterOperatorEqual, "production"))
assert.Equal(t, "(simpleJSONHas(labels, 'deployment.environment.name') OR simpleJSONHas(labels, 'deployment.environment'))", buildResourceFilter("", requestedName, v3.FilterOperatorExists, nil))
assert.Equal(t, "(not simpleJSONHas(labels, 'deployment.environment.name') AND not simpleJSONHas(labels, 'deployment.environment'))", buildResourceFilter("", requestedName, v3.FilterOperatorNotExists, nil))
assert.Equal(t, "(labels like '%deployment.environment.name\":\"production%' OR labels like '%deployment.environment\":\"production%')", buildResourceIndexFilter(requestedName, v3.FilterOperatorEqual, "production"))
assert.Empty(t, buildResourceIndexFilter(requestedName, v3.FilterOperatorNotEqual, "production"), "negative family filter must not use a rejecting index hint")
})
}
filters, err := buildResourceFiltersFromFilterItems(&v3.FilterSet{Items: []v3.FilterItem{{
Key: v3.AttributeKey{
Key: "deployment.environment.name",
DataType: v3.AttributeKeyDataTypeString,
Type: v3.AttributeKeyTypeResource,
},
Operator: v3.FilterOperatorEqual,
Value: "production",
}}})
require.NoError(t, err, "family filter items must build before their output is inspected")
wantFilters := []string{
resolvedValue + " = 'production'",
"(labels like '%deployment.environment.name\":\"production%' OR labels like '%deployment.environment\":\"production%')",
}
assert.Equal(t, wantFilters, filters)
wantPresence := "((simpleJSONHas(labels, 'deployment.environment.name') OR simpleJSONHas(labels, 'deployment.environment')) AND (labels like '%deployment.environment.name%' OR labels like '%deployment.environment%'))"
groupBy := buildResourceFiltersFromGroupBy([]v3.AttributeKey{{Key: "deployment.environment", Type: v3.AttributeKeyTypeResource}})
assert.Equal(t, []string{wantPresence}, groupBy)
assert.Equal(t, wantPresence, buildResourceFiltersFromAggregateAttribute(v3.AttributeKey{Key: "deployment.environment.name", Type: v3.AttributeKeyTypeResource}))
}

View File

@@ -6,16 +6,32 @@ import (
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/SigNoz/signoz/pkg/query-service/model"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
var (
columns = map[string]struct{}{
"deployment_environment": {},
"k8s_cluster_name": {},
"k8s_namespace_name": {},
}
columns = serviceMapColumns()
)
func serviceMapColumns() map[string]string {
columns := map[string]string{
"k8s_cluster_name": "k8s_cluster_name",
"k8s_namespace_name": "k8s_namespace_name",
}
// Dependency-graph rows keep their historical physical column name. Both
// semantic-convention request spellings target that same derived column.
for _, member := range semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
}) {
columns[strings.ReplaceAll(member, ".", "_")] = "deployment_environment"
}
return columns
}
func BuildServiceMapQuery(tags []model.TagQuery) (string, []interface{}) {
var filterQuery string
var namedArgs []interface{}
@@ -24,39 +40,40 @@ func BuildServiceMapQuery(tags []model.TagQuery) (string, []interface{}) {
operator := tag.GetOperator()
value := tag.GetValues()
if _, ok := columns[key]; !ok {
column, ok := columns[key]
if !ok {
continue
}
switch operator {
case model.InOperator:
filterQuery += fmt.Sprintf(" AND %s IN @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s IN @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, value))
case model.NotInOperator:
filterQuery += fmt.Sprintf(" AND %s NOT IN @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s NOT IN @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, value))
case model.EqualOperator:
filterQuery += fmt.Sprintf(" AND %s = @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s = @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, value))
case model.NotEqualOperator:
filterQuery += fmt.Sprintf(" AND %s != @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s != @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, value))
case model.ContainsOperator:
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%%%s%%", value)))
case model.NotContainsOperator:
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%%%s%%", value)))
case model.StartsWithOperator:
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%s%%", value)))
case model.NotStartsWithOperator:
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", key, key)
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", column, key)
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%s%%", value)))
case model.ExistsOperator:
filterQuery += fmt.Sprintf(" AND %s IS NOT NULL", key)
filterQuery += fmt.Sprintf(" AND %s IS NOT NULL", column)
case model.NotExistsOperator:
filterQuery += fmt.Sprintf(" AND %s IS NULL", key)
filterQuery += fmt.Sprintf(" AND %s IS NULL", column)
}
}
return filterQuery, namedArgs

View File

@@ -0,0 +1,35 @@
package services
import (
"testing"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/SigNoz/signoz/pkg/query-service/model"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestBuildServiceMapQueryAcceptsEnvironmentFamily(t *testing.T) {
for _, requestedName := range []string{"deployment.environment.name", "deployment.environment"} {
t.Run(requestedName, func(t *testing.T) {
tags := []model.TagQuery{model.NewTagQueryString(model.TagQueryParam{
Key: requestedName,
StringValues: []string{"production"},
Operator: model.InOperator,
TagType: model.ResourceAttributeTagType,
})}
query, args := BuildServiceMapQuery(tags)
argName := "deployment_environment"
if requestedName == "deployment.environment.name" {
argName = "deployment_environment_name"
}
assert.Equal(t, " AND deployment_environment IN @"+argName, query)
require.Len(t, args, 1)
named, ok := args[0].(driver.NamedValue)
require.True(t, ok)
assert.Equal(t, argName, named.Name)
assert.Equal(t, []interface{}{"production"}, named.Value)
})
}
}

View File

@@ -1,27 +0,0 @@
package metrics
var MetricsUnderTransition = map[string]string{
"k8s_pod_cpu_utilization": "k8s_pod_cpu_usage",
"k8s_node_cpu_utilization": "k8s_node_cpu_usage",
"container_cpu_utilization": "container_cpu_usage",
}
var DotMetricsUnderTransition = map[string]string{
"k8s.pod.cpu.utilization": "k8s.pod.cpu.usage",
"k8s.node.cpu.utilization": "k8s.node.cpu.usage",
"container.cpu.utilization": "container.cpu.usage",
}
func GetTransitionedMetric(metric string, normalized bool) string {
if normalized {
if _, ok := MetricsUnderTransition[metric]; ok {
return MetricsUnderTransition[metric]
}
return metric
} else {
if _, ok := DotMetricsUnderTransition[metric]; ok {
return DotMetricsUnderTransition[metric]
}
return metric
}
}

View File

@@ -10,8 +10,8 @@ import (
"log/slog"
"github.com/SigNoz/signoz/pkg/query-service/constants"
"github.com/SigNoz/signoz/pkg/query-service/metrics"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/semconv"
)
// ValidateAndCastValue validates and casts the value of a key to the corresponding data type of the key
@@ -234,12 +234,12 @@ func ClickHouseFormattedValue(v interface{}) string {
func ClickHouseFormattedMetricNames(v interface{}) string {
if name, ok := v.(string); ok {
transitionedMetrics := metrics.GetTransitionedMetric(name, !constants.IsDotMetricsEnabled)
if transitionedMetrics != name {
return ClickHouseFormattedValue([]interface{}{transitionedMetrics})
} else {
return ClickHouseFormattedValue([]interface{}{name})
members := semconv.MetricNames(name)
values := make([]interface{}, 0, len(members))
for _, member := range members {
values = append(values, member)
}
return ClickHouseFormattedValue(values)
}
return ClickHouseFormattedValue(v)

View File

@@ -2,13 +2,29 @@ package querybuilder
import (
"fmt"
"strings"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
func physicalSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
return []string{key.Name}
}
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: key.Signal,
FieldContext: key.FieldContext,
})
}
// ExistsExpression renders the existence predicate for a key resolved to the given
// columns (negated when exists is false). Comparisons are against constants rendered
// as literals, so the expression carries no bind args and can guard column expressions
@@ -43,11 +59,26 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
if len(evolutionsEntries) > 0 && evolutionsEntries[0] != nil {
columnName = evolutionsEntries[0].ColumnName
}
rawPath := fmt.Sprintf("%s.`%s`", columnName, key.Name)
if exists {
return rawPath + " IS NOT NULL", nil
members := physicalSemconvMembers(key)
paths := make([]string, 0, len(members))
for _, member := range members {
paths = append(paths, fmt.Sprintf("%s.`%s`", columnName, member))
}
return rawPath + " IS NULL", nil
if len(paths) == 1 {
if exists {
return paths[0] + " IS NOT NULL", nil
}
return paths[0] + " IS NULL", nil
}
guards := make([]string, 0, len(paths))
for _, path := range paths {
guards = append(guards, path+" IS NOT NULL")
}
rawPath := "(" + strings.Join(guards, " OR ") + ")"
if exists {
return rawPath, nil
}
return "NOT " + rawPath, nil
case schema.ColumnTypeEnumString,
schema.ColumnTypeEnumFixedString:
if exists {
@@ -88,8 +119,16 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumBool, schema.ColumnTypeEnumFloat64:
leftOperand := fmt.Sprintf("mapContains(%s, '%s')", column.Name, key.Name)
if key.Materialized {
members := physicalSemconvMembers(key)
operands := make([]string, 0, len(members))
for _, member := range members {
operands = append(operands, fmt.Sprintf("mapContains(%s, '%s')", column.Name, member))
}
leftOperand := strings.Join(operands, " OR ")
if len(operands) > 1 {
leftOperand = "(" + leftOperand + ")"
}
if key.Materialized && (len(members) == 1 || key.MaterializedSemconv) {
leftOperand = telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key)
}
if exists {

View File

@@ -0,0 +1,82 @@
package querybuilder_test
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/tracestelemetryschema"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestTraceFamilyUsesMaterializedHistoricalMember(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
historical := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextAttribute,
telemetrytypes.FieldDataTypeString,
)
matches := querybuilder.MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
historical.Name: {historical},
})
require.Len(t, matches, 1, "family metadata should resolve to one logical field")
expression, err := tracestelemetryschema.NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, matches[0])
require.NoError(t, err, "resolved trace family should map to a value expression")
assert.Equal(
t,
"COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(`attribute_string_deployment$$environment`, ''), '')",
expression,
"family expression should retain the promoted historical member",
)
}
func TestTraceFamilyUsesFamilyAwareMaterializedColumn(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
MaterializedColumnName: "attribute_string_db$$system",
MaterializedSemconv: true,
}
historical := &telemetrytypes.TelemetryFieldKey{
Name: "db.system",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextAttribute,
telemetrytypes.FieldDataTypeString,
)
matches := querybuilder.MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
historical.Name: {historical},
})
require.Len(t, matches, 1, "family metadata should resolve to one logical field")
expression, err := tracestelemetryschema.NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, matches[0])
require.NoError(t, err, "resolved trace family should map to a value expression")
assert.Equal(t, "`attribute_string_db$$system`", expression, "family-aware promoted column should serve the complete family")
}

View File

@@ -4,12 +4,14 @@ import (
"context"
"fmt"
"log/slog"
"maps"
"slices"
"strconv"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
grammar "github.com/SigNoz/signoz/pkg/parser/filterquery/grammar"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -38,6 +40,7 @@ type filterExpressionVisitor struct {
fullTextColumn *telemetrytypes.TelemetryFieldKey
skipResourceFilter bool
skipFullTextFilter bool
exactSemconv bool
variables map[string]qbtypes.VariableItem
keysWithWarnings map[string]bool
@@ -58,6 +61,7 @@ type FilterExprVisitorOpts struct {
FullTextColumn *telemetrytypes.TelemetryFieldKey
SkipResourceFilter bool
SkipFullTextFilter bool
ExactSemconv bool
Variables map[string]qbtypes.VariableItem
StartNs uint64
EndNs uint64
@@ -75,6 +79,7 @@ func newFilterExpressionVisitor(opts FilterExprVisitorOpts) *filterExpressionVis
fullTextColumn: opts.FullTextColumn,
skipResourceFilter: opts.SkipResourceFilter,
skipFullTextFilter: opts.SkipFullTextFilter,
exactSemconv: opts.ExactSemconv,
variables: opts.Variables,
keysWithWarnings: make(map[string]bool),
startNs: opts.StartNs,
@@ -379,7 +384,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
// VisitComparison handles all comparison operators.
func (v *filterExpressionVisitor) VisitComparison(ctx *grammar.ComparisonContext) any {
key := v.Visit(ctx.Key()).(*telemetrytypes.TelemetryFieldKey)
matching := MatchingFieldKeys(key, v.fieldKeys)
matching := v.matchingFieldKeys(key)
// Handle EXISTS specially
if ctx.EXISTS() != nil {
@@ -730,7 +735,7 @@ func (v *filterExpressionVisitor) VisitFunctionCall(ctx *grammar.FunctionCallCon
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(key, MatchingFieldKeys(key, v.fieldKeys), operator, value)
conds, ok := v.buildConditions(key, v.matchingFieldKeys(key), operator, value)
if !ok {
return ErrorConditionLiteral
}
@@ -923,7 +928,7 @@ func (v *filterExpressionVisitor) VisitKey(ctx *grammar.KeyContext) any {
// buildConditions invokes the condition builder for a filter term, folding its
// warnings/errors into visitor state; returns false if an error was recorded.
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) ([]string, bool) {
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter}, op, value, v.builder)
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter, ExactSemconv: v.exactSemconv}, op, value, v.builder)
if err != nil {
_, _, _, _, errURL, _ := errors.Unwrapb(err)
assignIfEmpty(&v.mainErrorURL, errURL)
@@ -982,27 +987,159 @@ func assignIfEmpty(s *string, value string) {
// MatchingFieldKeys returns the field keys from the map that match the given key,
// honoring any context/data type the user specified.
func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
fieldKeysForName := []*telemetrytypes.TelemetryFieldKey{}
return matchingFieldKeys(field, fieldKeys, true)
}
// match by name; keep items whose context and data type match (unspecified matches any)
for _, item := range fieldKeys[field.Name] {
if (field.FieldContext == telemetrytypes.FieldContextUnspecified || field.FieldContext == item.FieldContext) &&
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || field.FieldDataType == item.FieldDataType) {
fieldKeysForName = append(fieldKeysForName, item)
// MatchingFieldKeysExact matches only the requested physical spelling.
func MatchingFieldKeysExact(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
return matchingFieldKeys(field, fieldKeys, false)
}
func matchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey, resolveSemconv bool) []*telemetrytypes.TelemetryFieldKey {
members := []string{field.Name}
if resolveSemconv {
signals := []telemetrytypes.Signal{field.Signal}
if field.Signal == telemetrytypes.SignalUnspecified {
signals = []telemetrytypes.Signal{
telemetrytypes.SignalTraces,
telemetrytypes.SignalLogs,
telemetrytypes.SignalMetrics,
}
}
members = nil
for _, signal := range signals {
selector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: signal,
FieldContext: field.FieldContext,
}
for _, member := range semconv.AttributeMembers(selector) {
if !slices.Contains(members, member) {
members = append(members, member)
}
}
}
if len(members) == 0 {
members = []string{field.Name}
}
}
fieldKeysForName := make([]*telemetrytypes.TelemetryFieldKey, 0)
indexByIdentity := make(map[string]int)
appendMatches := func(lookupName string, memberName string, contextAlreadyMatched bool) {
for _, item := range fieldKeys[lookupName] {
if !contextAlreadyMatched && field.FieldContext != telemetrytypes.FieldContextUnspecified && field.FieldContext != item.FieldContext {
continue
}
if field.FieldDataType != telemetrytypes.FieldDataTypeUnspecified && field.FieldDataType != item.FieldDataType {
continue
}
itemMembers := []string{field.Name}
if resolveSemconv {
itemMembers = semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
})
}
familyMatch := resolveSemconv && len(itemMembers) > 1
// A wildcard lookup may find the same name in a signal or context where
// the family does not apply. Exact names remain valid there; aliases do not.
if memberName != field.Name && (!familyMatch || !slices.Contains(itemMembers, memberName)) {
continue
}
physicalMembers := item.SemconvMembers
if len(physicalMembers) == 0 {
physicalMembers = []string{memberName}
}
materializedColumns := maps.Clone(item.SemconvMaterializedColumns)
if item.Materialized {
if materializedColumns == nil {
materializedColumns = make(map[string]string)
}
physicalKey := *item
physicalKey.Name = memberName
materializedColumns[memberName] = strings.Trim(telemetrytypes.FieldKeyToMaterializedColumnName(&physicalKey), "`")
}
identity := item.Signal.StringValue() + ";" + item.FieldContext.StringValue() + ";" + item.FieldDataType.StringValue()
if familyMatch {
if index, found := indexByIdentity[identity]; found {
for _, physicalMember := range physicalMembers {
if !slices.Contains(fieldKeysForName[index].SemconvMembers, physicalMember) {
fieldKeysForName[index].SemconvMembers = append(fieldKeysForName[index].SemconvMembers, physicalMember)
}
}
if len(materializedColumns) > 0 {
if fieldKeysForName[index].SemconvMaterializedColumns == nil {
fieldKeysForName[index].SemconvMaterializedColumns = make(map[string]string)
}
maps.Copy(fieldKeysForName[index].SemconvMaterializedColumns, materializedColumns)
}
if item.Materialized && item.MaterializedSemconv {
fieldKeysForName[index].Materialized = true
fieldKeysForName[index].MaterializedColumnName = item.MaterializedColumnName
fieldKeysForName[index].MaterializedSemconv = true
}
continue
}
indexByIdentity[identity] = len(fieldKeysForName)
}
resolved := *item
// The requested spelling is the response identity. Field mappers use
// it to resolve the available family members current-first.
if familyMatch {
resolved.Name = field.Name
resolved.SemconvMembers = slices.Clone(physicalMembers)
resolved.SemconvMaterializedColumns = materializedColumns
if !resolved.MaterializedSemconv {
// Materialization is member-specific after family keys are merged.
resolved.Materialized = false
}
} else if !resolveSemconv {
resolved.SemconvMembers = []string{field.Name}
}
fieldKeysForName = append(fieldKeysForName, &resolved)
}
}
// A context may have been split off a name that legitimately contained it (e.g.
// `attribute.key`); also look up the context-prefixed name so both readings resolve.
// Members are current-first, so metadata from the current key wins when
// both spellings describe the same signal/context/type.
for _, member := range members {
appendMatches(member, member, false)
}
// A context may have been split off a name that legitimately contained it
// (e.g. `attribute.key`); preserve that historical alternate reading for
// every family member.
if field.FieldContext != telemetrytypes.FieldContextUnspecified {
contextPrefixedFieldName := fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), field.Name)
for _, item := range fieldKeys[contextPrefixedFieldName] {
// Context already matched via the lookup key; only data type needs checking.
if field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || item.FieldDataType == field.FieldDataType {
fieldKeysForName = append(fieldKeysForName, item)
}
for _, member := range members {
appendMatches(fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), member), member, true)
}
}
return fieldKeysForName
}
func (v *filterExpressionVisitor) matchingFieldKeys(field *telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
if v.exactSemconv {
return MatchingFieldKeysExact(field, v.fieldKeys)
}
return MatchingFieldKeys(field, v.fieldKeys)
}
// ExactSemconvKeys returns copies pinned to their physical names, preventing a
// field mapper from expanding a synthesized or metadata-free key into a family.
func ExactSemconvKeys(keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
result := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
for _, key := range keys {
if key == nil {
continue
}
resolved := *key
resolved.SemconvMembers = []string{key.Name}
result = append(result, &resolved)
}
return result
}

View File

@@ -14,6 +14,7 @@ import (
"github.com/antlr4-go/antlr/v4"
sqlbuilder "github.com/huandu/go-sqlbuilder"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// TestPrepareWhereClause_EmptyVariableList ensures PrepareWhereClause errors when a variable has an empty list value.
@@ -685,6 +686,147 @@ func TestVisitKey(t *testing.T) {
}
}
func TestMatchingFieldKeysResolvesCurrentTraceNameFromOldMetadata(t *testing.T) {
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Description: "old metadata",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
"deployment.environment.name",
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{old.Name: {old}})
require.Len(t, matches, 1, "trace family lookup must resolve before inspecting metadata")
assert.Equal(t, "deployment.environment.name", matches[0].Name)
assert.Equal(t, "old metadata", matches[0].Description)
assert.Equal(t, []string{"deployment.environment"}, matches[0].SemconvMembers)
}
func TestMatchingFieldKeysUsesCurrentTraceMetadataForOldName(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Description: "current metadata",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Description: "old metadata",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
old.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "trace family lookup must resolve before inspecting metadata")
assert.Equal(t, old.Name, matches[0].Name)
assert.Equal(t, "current metadata", matches[0].Description)
assert.Equal(t, []string{current.Name, old.Name}, matches[0].SemconvMembers)
}
func TestMatchingFieldKeysResolvesLogSemconvFamily(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "log family lookup must resolve before inspecting metadata")
assert.Equal(t, current.Name, matches[0].Name)
assert.Equal(t, []string{current.Name, old.Name}, matches[0].SemconvMembers)
}
func TestMatchingFieldKeysResolvesMetricStorageSpellings(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "resource_deployment_environment_name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "resource_deployment_environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
"deployment.environment.name",
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "metric family lookup must resolve storage spellings")
assert.Equal(t, requested.Name, matches[0].Name)
assert.Equal(t, []string{current.Name, old.Name}, matches[0].SemconvMembers)
}
func TestMatchingFieldKeysExactKeepsRequestedPhysicalName(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeysExact(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "exact lookup must return the requested physical field")
assert.Equal(t, current.Name, matches[0].Name)
assert.Equal(t, []string{current.Name}, matches[0].SemconvMembers)
}
// ---------------------------------------------------------------------------
// TestVisitComparison
// ---------------------------------------------------------------------------

View File

@@ -2,21 +2,47 @@
package semconv
import "github.com/SigNoz/signoz/pkg/types/telemetrytypes"
var families = []Family{
{
Current: "container.cpu.usage",
Old: []string{"container.cpu.utilization"},
Kind: KindMetric,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextMetric},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalMetrics},
ApplyToMetrics: nil,
},
{
Current: "db.system.name",
Old: []string{"db.system"},
Kind: KindAttribute,
Contexts: nil,
Signals: nil,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextAttribute, telemetrytypes.FieldContextResource},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalMetrics, telemetrytypes.SignalTraces},
ApplyToMetrics: nil,
},
{
Current: "deployment.environment.name",
Old: []string{"deployment.environment"},
Kind: KindAttribute,
Contexts: nil,
Signals: nil,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextAttribute, telemetrytypes.FieldContextResource},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalMetrics, telemetrytypes.SignalTraces},
ApplyToMetrics: nil,
},
{
Current: "k8s.node.cpu.usage",
Old: []string{"k8s.node.cpu.utilization"},
Kind: KindMetric,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextMetric},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalMetrics},
ApplyToMetrics: nil,
},
{
Current: "k8s.pod.cpu.usage",
Old: []string{"k8s.pod.cpu.utilization"},
Kind: KindMetric,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextMetric},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalMetrics},
ApplyToMetrics: nil,
},
}

View File

@@ -2,6 +2,7 @@ package semconv
import (
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -26,6 +27,13 @@ type Family struct {
ValueMap map[string]string
}
type metricSpelling uint8
const (
metricSpellingDotted metricSpelling = iota
metricSpellingNormalized
)
var (
KindAttribute = Kind{String: valuer.NewString("attribute")}
KindMetric = Kind{String: valuer.NewString("metric")}
@@ -59,6 +67,88 @@ func Members(kind Kind, selector telemetrytypes.FieldKeySelector) []string {
return familyMembers[idx]
}
// AttributeMembers returns the physical attribute spellings that may represent
// selector.Name. Metrics have used both dotted and normalized label layouts;
// resource labels have additionally used a resource_ prefix. Keeping that
// storage detail here prevents metrics readers from maintaining local
// transition tables.
func AttributeMembers(selector telemetrytypes.FieldKeySelector) []string {
if selector.Signal != telemetrytypes.SignalMetrics {
return Members(KindAttribute, selector)
}
lookupSelector := selector
lookupSelector.Name = strings.TrimPrefix(selector.Name, "resource_")
family, style, ok := lookupMetricSpelling(KindAttribute, lookupSelector)
if !ok {
return []string{selector.Name}
}
logicalMembers := familyMembers[family]
result := make([]string, 0, len(logicalMembers)*4)
for _, member := range logicalMembers {
dotted := member
normalized := normalizeMetricSpelling(member)
variants := []string{dotted, normalized}
if style == metricSpellingNormalized {
variants[0], variants[1] = variants[1], variants[0]
}
if selector.FieldContext == telemetrytypes.FieldContextResource ||
selector.FieldContext == telemetrytypes.FieldContextUnspecified ||
strings.HasPrefix(selector.Name, "resource_") {
for _, variant := range variants {
result = appendUniqueString(result, "resource_"+variant)
}
}
for _, variant := range variants {
result = appendUniqueString(result, variant)
}
}
return result
}
// MetricNames returns the current and historical storage names for a metric.
// The input's dotted or normalized style is preserved because both layouts are
// valid metric identities and must not be mixed in one query.
func MetricNames(name string) []string {
selector := telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
}
family, style, ok := lookupMetricSpelling(KindMetric, selector)
if !ok {
return []string{name}
}
logicalMembers := familyMembers[family]
result := make([]string, 0, len(logicalMembers))
for _, member := range logicalMembers {
if style == metricSpellingNormalized {
member = normalizeMetricSpelling(member)
}
result = appendUniqueString(result, member)
}
return result
}
// CurrentAttribute returns the canonical dotted name for an attribute
// spelling, or selector.Name if no enabled family matches.
func CurrentAttribute(selector telemetrytypes.FieldKeySelector) string {
if selector.Signal != telemetrytypes.SignalMetrics {
return Current(KindAttribute, selector)
}
lookupSelector := selector
lookupSelector.Name = strings.TrimPrefix(selector.Name, "resource_")
family, _, ok := lookupMetricSpelling(KindAttribute, lookupSelector)
if !ok {
return selector.Name
}
return families[family].Current
}
// Current returns the current name for selector.Name, or the input name when
// it does not belong to an enabled family.
func Current(kind Kind, selector telemetrytypes.FieldKeySelector) string {
@@ -99,6 +189,35 @@ func lookupIndex(kind Kind, selector telemetrytypes.FieldKeySelector) (int, bool
return 0, false
}
func lookupMetricSpelling(kind Kind, selector telemetrytypes.FieldKeySelector) (int, metricSpelling, bool) {
if idx, ok := lookupIndex(kind, selector); ok {
return idx, metricSpellingDotted, true
}
for idx, family := range families {
if !matchesSelector(family, kind, selector) {
continue
}
for _, member := range familyMembers[idx] {
if normalizeMetricSpelling(member) == selector.Name {
return idx, metricSpellingNormalized, true
}
}
}
return 0, metricSpellingDotted, false
}
func normalizeMetricSpelling(name string) string {
return strings.ReplaceAll(name, ".", "_")
}
func appendUniqueString(values []string, value string) []string {
if value == "" || slices.Contains(values, value) {
return values
}
return append(values, value)
}
func matchesSelector(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if family.Kind != kind {
return false

View File

@@ -78,3 +78,57 @@ func TestMembersReturnsInputWhenKindDoesNotMatch(t *testing.T) {
"an attribute family must not match a metric-name lookup",
)
}
func TestAttributeMembersIncludesMetricResourceStorageSpellings(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, []string{
"resource_db.system.name", "resource_db_system_name", "db.system.name", "db_system_name",
"resource_db.system", "resource_db_system", "db.system", "db_system",
}, AttributeMembers(selector), "resource metric attributes should cover every historical storage layout")
}
func TestCurrentAttributeResolvesNormalizedMetricResourceSpelling(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "resource_db_system",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, "db.system.name", CurrentAttribute(selector), "normalized resource spelling should resolve to the dotted current name")
}
func TestAttributeMembersPreservesNormalizedMetricPointStyle(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "db_system",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
}
assert.Equal(t, []string{
"db_system_name", "db.system.name", "db_system", "db.system",
}, AttributeMembers(selector), "normalized point attribute should remain the preferred storage spelling")
}
func TestMetricNamesPreservesDottedStyle(t *testing.T) {
assert.Equal(t,
[]string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"},
MetricNames("k8s.pod.cpu.usage"),
"dotted metric input should produce dotted family names",
)
}
func TestMetricNamesPreservesNormalizedStyle(t *testing.T) {
assert.Equal(t,
[]string{"k8s_pod_cpu_usage", "k8s_pod_cpu_utilization"},
MetricNames("k8s_pod_cpu_utilization"),
"normalized metric input should produce normalized family names",
)
}
func TestMetricNamesReturnsUnknownNameUnchanged(t *testing.T) {
assert.Equal(t, []string{"custom_metric"}, MetricNames("custom_metric"), "unknown metrics should not be expanded")
}

View File

@@ -237,6 +237,7 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddDashboardTuplesFactory(sqlstore),
sqlmigration.NewRestructureSavedViewSpecFactory(sqlstore, sqlschema),
sqlmigration.NewAddSavedViewTuplesFactory(sqlstore),
sqlmigration.NewMigrateDeploymentEnvironmentQuickFilterFactory(),
)
}

View File

@@ -0,0 +1,127 @@
package sqlmigration
import (
"context"
"encoding/json"
"log/slog"
"time"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
)
const deploymentEnvironmentCurrent = "deployment.environment.name"
type migrateDeploymentEnvironmentQuickFilter struct {
logger *slog.Logger
}
type semconvQuickFilterRow struct {
bun.BaseModel `bun:"table:quick_filter"`
ID string `bun:"id"`
Filter string `bun:"filter"`
}
func NewMigrateDeploymentEnvironmentQuickFilterFactory() factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(
factory.MustNewName("migrate_semconv_quick_filter"),
func(_ context.Context, settings factory.ProviderSettings, _ Config) (SQLMigration, error) {
return &migrateDeploymentEnvironmentQuickFilter{logger: settings.Logger}, nil
},
)
}
func (migration *migrateDeploymentEnvironmentQuickFilter) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func deploymentEnvironmentOld() string {
members := semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: deploymentEnvironmentCurrent,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
})
if len(members) < 2 {
return deploymentEnvironmentCurrent
}
return members[1]
}
func rewriteQuickFilterSemconv(filterJSON, from, to string) (string, bool, error) {
var filters []map[string]any
if err := json.Unmarshal([]byte(filterJSON), &filters); err != nil {
return "", false, err
}
changed := false
for _, filter := range filters {
if key, ok := filter["key"].(string); ok && key == from {
filter["key"] = to
changed = true
}
}
if !changed {
return filterJSON, false, nil
}
rewritten, err := json.Marshal(filters)
if err != nil {
return "", false, err
}
return string(rewritten), true, nil
}
func (migration *migrateDeploymentEnvironmentQuickFilter) migrate(ctx context.Context, db *bun.DB, from, to string) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
rows := make([]*semconvQuickFilterRow, 0)
if err := tx.NewSelect().
Model(&rows).
Where("signal IN (?)", bun.In([]string{"traces", "api_monitoring", "exceptions"})).
Scan(ctx); err != nil {
return err
}
for _, row := range rows {
rewritten, changed, err := rewriteQuickFilterSemconv(row.Filter, from, to)
if err != nil {
// Quick filters are user-editable. One malformed legacy row must not
// prevent the application from starting or block every other org's
// migration.
if migration.logger != nil {
migration.logger.WarnContext(ctx, "skipping quick filter with unreadable filter JSON",
slog.String("quick_filter_id", row.ID), slog.Any("error", err))
}
continue
}
if !changed {
continue
}
if _, err := tx.NewUpdate().
Model((*semconvQuickFilterRow)(nil)).
Set("filter = ?", rewritten).
Set("updated_at = ?", time.Now()).
Where("id = ?", row.ID).
Exec(ctx); err != nil {
return err
}
}
return tx.Commit()
}
func (migration *migrateDeploymentEnvironmentQuickFilter) Up(ctx context.Context, db *bun.DB) error {
return migration.migrate(ctx, db, deploymentEnvironmentOld(), deploymentEnvironmentCurrent)
}
func (migration *migrateDeploymentEnvironmentQuickFilter) Down(ctx context.Context, db *bun.DB) error {
return migration.migrate(ctx, db, deploymentEnvironmentCurrent, deploymentEnvironmentOld())
}

View File

@@ -0,0 +1,38 @@
package sqlmigration
import (
"encoding/json"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestRewriteQuickFilterSemconv(t *testing.T) {
oldName := deploymentEnvironmentOld()
input := `[{"key":"service.name","dataType":"string","type":"resource"},{"key":"` + oldName + `","dataType":"string","type":"resource","custom":true}]`
rewritten, changed, err := rewriteQuickFilterSemconv(input, oldName, deploymentEnvironmentCurrent)
require.NoError(t, err)
assert.True(t, changed)
var filters []map[string]any
require.NoError(t, json.Unmarshal([]byte(rewritten), &filters))
assert.Equal(t, "service.name", filters[0]["key"])
assert.Equal(t, deploymentEnvironmentCurrent, filters[1]["key"])
assert.Equal(t, true, filters[1]["custom"], "unknown filter properties must be preserved")
restored, changed, err := rewriteQuickFilterSemconv(rewritten, deploymentEnvironmentCurrent, oldName)
require.NoError(t, err)
assert.True(t, changed)
require.NoError(t, json.Unmarshal([]byte(restored), &filters))
assert.Equal(t, oldName, filters[1]["key"])
}
func TestRewriteQuickFilterSemconvNoop(t *testing.T) {
input := `[{"key":"service.name","dataType":"string","type":"resource"}]`
rewritten, changed, err := rewriteQuickFilterSemconv(input, deploymentEnvironmentOld(), deploymentEnvironmentCurrent)
require.NoError(t, err)
assert.False(t, changed)
assert.Equal(t, input, rewritten)
}

View File

@@ -10,6 +10,7 @@ import (
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
"github.com/SigNoz/signoz/pkg/types/metrictypes"
@@ -31,6 +32,15 @@ const (
OthersMultiTemporality = `IF(LOWER(temporality) LIKE LOWER('delta'), %s, %s) AS per_series_value`
)
func metricNameValues(name string) []any {
members := semconv.MetricNames(name)
values := make([]any, 0, len(members))
for _, member := range members {
values = append(values, member)
}
return values
}
type StatementBuilder struct {
logger *slog.Logger
metadataStore telemetrytypes.MetadataStore
@@ -341,7 +351,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
sb.SelectMore(col)
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -385,7 +395,7 @@ func (b *StatementBuilder) buildReducedSpatialAggFastPath(
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", agg.MetricName),
sb.In("metric_name", metricNameValues(agg.MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -427,7 +437,7 @@ func (b *StatementBuilder) buildReducedTemporalAggregationCTE(
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", agg.MetricName),
sb.In("metric_name", metricNameValues(agg.MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -505,7 +515,7 @@ func (b *StatementBuilder) buildTemporalAggDeltaFastPath(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -560,7 +570,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -638,7 +648,7 @@ func (b *StatementBuilder) buildTemporalAggDelta(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -679,7 +689,7 @@ func (b *StatementBuilder) buildTemporalAggCumulativeOrUnspecified(
baseSb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
baseSb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
baseSb.Where(
baseSb.In("metric_name", query.Aggregations[0].MetricName),
baseSb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
baseSb.GTE("unix_milli", start),
baseSb.LT("unix_milli", end),
)
@@ -770,7 +780,7 @@ func (b *StatementBuilder) buildTemporalAggForMultipleTemporalities(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)

View File

@@ -13,9 +13,21 @@ import (
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestMetricNameValuesResolveFamily(t *testing.T) {
assert.Equal(t,
[]any{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"},
metricNameValues("k8s.pod.cpu.usage"),
)
assert.Equal(t,
[]any{"container_cpu_usage", "container_cpu_utilization"},
metricNameValues("container_cpu_utilization"),
)
}
func TestStatementBuilder(t *testing.T) {
cases := []struct {
name string

View File

@@ -44,6 +44,73 @@ func keyIndexFilter(key *telemetrytypes.TelemetryFieldKey) any {
return fmt.Sprintf(`%%%s%%`, key.Name)
}
func memberKey(key *telemetrytypes.TelemetryFieldKey, name string) *telemetrytypes.TelemetryFieldKey {
member := *key
member.Name = name
return &member
}
func keyIndexCondition(sb *sqlbuilder.SelectBuilder, column string, key *telemetrytypes.TelemetryFieldKey, members []string) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
conditions = append(conditions, sb.Like(column, keyIndexFilter(memberKey(key, member))))
}
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
func valueIndexCondition(
sb *sqlbuilder.SelectBuilder,
column string,
key *telemetrytypes.TelemetryFieldKey,
members []string,
op qbtypes.FilterOperator,
value any,
caseInsensitive bool,
) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
patterns := valueForIndexFilter(op, memberKey(key, member), value)
switch values := patterns.(type) {
case []string:
for _, pattern := range values {
conditions = append(conditions, sb.Like(column, pattern))
}
default:
if caseInsensitive {
conditions = append(conditions, sb.ILike(column, values))
} else {
conditions = append(conditions, sb.Like(column, values))
}
}
}
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
func memberPresenceCondition(sb *sqlbuilder.SelectBuilder, column string, members []string, exists bool) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
field := fmt.Sprintf("simpleJSONHas(%s, '%s')", column, member)
if exists {
conditions = append(conditions, sb.E(field, true))
} else {
conditions = append(conditions, sb.NE(field, true))
}
}
if exists {
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
return sb.And(conditions...)
}
// SkipResourceFilter is not applicable here: the fingerprint table only stores resource attributes.
func (b *defaultConditionBuilder) ConditionFor(
ctx context.Context,
@@ -115,8 +182,10 @@ func (b *defaultConditionBuilder) conditionForKey(
// as we have not changed the resource column in the resource fingerprint table.
column := columns[0]
keyIdxFilter := sb.Like(column.Name, keyIndexFilter(key))
valueForIndexFilter := valueForIndexFilter(op, key, value)
members := resourceSemconvMembers(key)
isFamily := len(members) > 1
keyIdxFilter := keyIndexCondition(sb, column.Name, key, members)
singleValueIndexFilter := valueForIndexFilter(op, memberKey(key, members[0]), value)
fieldName, err := b.fm.FieldFor(ctx, valuer.UUID{}, startNs, endNs, key)
if err != nil {
@@ -128,12 +197,15 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.E(fieldName, formattedValue),
keyIdxFilter,
sb.Like(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, key, members, op, value, false),
), nil
case qbtypes.FilterOperatorNotEqual:
if isFamily {
return sb.NE(fieldName, formattedValue), nil
}
return sb.And(
sb.NE(fieldName, formattedValue),
sb.NotLike(column.Name, valueForIndexFilter),
sb.NotLike(column.Name, singleValueIndexFilter),
), nil
case qbtypes.FilterOperatorGreaterThan:
return sb.And(sb.GT(fieldName, formattedValue), keyIdxFilter), nil
@@ -148,7 +220,7 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.ILike(fieldName, formattedValue),
keyIdxFilter,
sb.ILike(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, key, members, op, value, true),
), nil
case qbtypes.FilterOperatorNotLike, qbtypes.FilterOperatorNotILike:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
@@ -185,13 +257,11 @@ func (b *defaultConditionBuilder) conditionForKey(
inConditions = append(inConditions, sb.E(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.Or(inConditions...)
valConditions := make([]string, 0, len(values))
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.Like(column.Name, v))
}
}
mainCondition = sb.And(mainCondition, keyIdxFilter, sb.Or(valConditions...))
mainCondition = sb.And(
mainCondition,
keyIdxFilter,
valueIndexCondition(sb, column.Name, key, members, op, value, false),
)
return mainCondition, nil
case qbtypes.FilterOperatorNotIn:
@@ -204,8 +274,11 @@ func (b *defaultConditionBuilder) conditionForKey(
notInConditions = append(notInConditions, sb.NE(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.And(notInConditions...)
if isFamily {
return mainCondition, nil
}
valConditions := make([]string, 0, len(values))
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
if valuesForIndexFilter, ok := singleValueIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.NotLike(column.Name, v))
}
@@ -215,13 +288,11 @@ func (b *defaultConditionBuilder) conditionForKey(
case qbtypes.FilterOperatorExists:
return sb.And(
sb.E(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
memberPresenceCondition(sb, column.Name, members, true),
keyIdxFilter,
), nil
case qbtypes.FilterOperatorNotExists:
return sb.And(
sb.NE(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
), nil
return memberPresenceCondition(sb, column.Name, members, false), nil
case qbtypes.FilterOperatorRegexp:
return sb.And(
@@ -237,7 +308,7 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.ILike(fieldName, fmt.Sprintf(`%%%s%%`, formattedValue)),
keyIdxFilter,
sb.ILike(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, key, members, op, value, true),
), nil
case qbtypes.FilterOperatorNotContains:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else

View File

@@ -220,3 +220,230 @@ func TestConditionBuilder(t *testing.T) {
})
}
}
func TestFamilyPositiveFilterExcludesKeylessRows(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)")
assert.Equal(t, []any{
"production",
"%deployment.environment.name%",
"%deployment.environment%",
`%deployment.environment.name":"production%`,
`%deployment.environment":"production%`,
}, args)
}
func TestFamilyNotEqualIncludesKeylessRows(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotEqual, "staging", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ?")
assert.Equal(t, []any{"staging"}, args)
}
func TestFamilyNotInIncludesKeylessRows(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotIn, []any{"staging", "dev"}, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ? AND COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ?)")
assert.Equal(t, []any{"staging", "dev"}, args)
}
func TestFamilyNotLikeIncludesKeylessRows(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotLike, "%stag%", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "LOWER(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')) NOT LIKE LOWER(?)")
assert.Equal(t, []any{"%stag%"}, args)
}
func TestFamilyNotContainsIncludesKeylessRows(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotContains, "stag", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "LOWER(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')) NOT LIKE LOWER(?)")
assert.Equal(t, []any{"%stag%"}, args)
}
func TestFamilyNotRegexpIncludesKeylessRows(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotRegexp, "stag.*", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "NOT match(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), ''), ?)")
assert.Equal(t, []any{"stag.*"}, args)
}
func TestFamilyExistsChecksEveryMember(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorExists, nil, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "(simpleJSONHas(labels, 'deployment.environment.name') = ? OR simpleJSONHas(labels, 'deployment.environment') = ?) AND (labels LIKE ? OR labels LIKE ?)")
assert.Equal(t, []any{true, true, "%deployment.environment.name%", "%deployment.environment%"}, args)
}
func TestFamilyNotExistsChecksEveryMember(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotExists, nil, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "simpleJSONHas(labels, 'deployment.environment.name') <> ? AND simpleJSONHas(labels, 'deployment.environment') <> ?")
assert.Equal(t, []any{true, true}, args)
}
func TestLogSemconvFamilyResolvesEveryStoredMember(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, current,
map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)")
assert.Equal(t, []any{
"production",
"%deployment.environment.name%",
"%deployment.environment%",
`%deployment.environment.name":"production%`,
`%deployment.environment":"production%`,
}, args)
}

View File

@@ -3,8 +3,10 @@ package resourcefilter
import (
"context"
"fmt"
"strings"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -32,6 +34,21 @@ func NewFieldMapper() *defaultFieldMapper {
return &defaultFieldMapper{}
}
func resourceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if (key.Signal != telemetrytypes.SignalTraces && key.Signal != telemetrytypes.SignalLogs) ||
key.FieldContext != telemetrytypes.FieldContextResource {
return []string{key.Name}
}
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
})
}
func (m *defaultFieldMapper) getColumn(
_ context.Context,
_, _ uint64,
@@ -66,7 +83,15 @@ func (m *defaultFieldMapper) FieldFor(
return "", err
}
if key.FieldContext == telemetrytypes.FieldContextResource {
return fmt.Sprintf("simpleJSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
members := resourceSemconvMembers(key)
if len(members) > 1 {
values := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(simpleJSONExtractString(%s, '%s'), '')", columns[0].Name, member))
}
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
}
return fmt.Sprintf("simpleJSONExtractString(%s, '%s')", columns[0].Name, members[0]), nil
}
return columns[0].Name, nil
}

View File

@@ -14,6 +14,7 @@ import (
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/telemetryschema/audittelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema"
@@ -151,6 +152,26 @@ func (t *telemetryMetaStore) tracesTblStatementToFieldKeys(ctx context.Context)
return materialisedKeys, nil
}
func attributeSemconvMembers(name string, signal telemetrytypes.Signal, fieldContext telemetrytypes.FieldContext) []string {
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: name,
Signal: signal,
FieldContext: fieldContext,
})
}
func traceSemconvMembers(name string, fieldContext telemetrytypes.FieldContext) []string {
return attributeSemconvMembers(name, telemetrytypes.SignalTraces, fieldContext)
}
func inStrings(sb *sqlbuilder.SelectBuilder, column string, values []string) string {
args := make([]any, 0, len(values))
for _, value := range values {
args = append(args, value)
}
return sb.In(column, args...)
}
// getTracesKeys returns the keys from the spans that match the field selection criteria.
func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelectors []*telemetrytypes.FieldKeySelector) ([]*telemetrytypes.TelemetryFieldKey, bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
@@ -1312,6 +1333,37 @@ func (t *telemetryMetaStore) GetKeysMulti(ctx context.Context, orgID valuer.UUID
}
}
// GetKeys backs key suggestions and remains literal. Internal multi-key
// lookups expand selectors so query builders can inspect stored family members.
expandSelectors := func(selectors []*telemetrytypes.FieldKeySelector, signal telemetrytypes.Signal) []*telemetrytypes.FieldKeySelector {
expanded := make([]*telemetrytypes.FieldKeySelector, 0, len(selectors))
for _, selector := range selectors {
metricNames := []string{""}
if signal == telemetrytypes.SignalMetrics && selector.MetricContext != nil && selector.MetricContext.MetricName != "" {
metricNames = semconv.MetricNames(selector.MetricContext.MetricName)
}
for _, member := range attributeSemconvMembers(selector.Name, signal, selector.FieldContext) {
for _, metricName := range metricNames {
memberSelector := *selector
memberSelector.Name = member
if selector.MetricContext != nil {
metricContext := *selector.MetricContext
if metricName != "" {
metricContext.MetricName = metricName
}
memberSelector.MetricContext = &metricContext
}
expanded = append(expanded, &memberSelector)
}
}
}
return expanded
}
logsSelectors = expandSelectors(logsSelectors, telemetrytypes.SignalLogs)
tracesSelectors = expandSelectors(tracesSelectors, telemetrytypes.SignalTraces)
metricsSelectors = expandSelectors(metricsSelectors, telemetrytypes.SignalMetrics)
meterSourceMetricsSelectors = expandSelectors(meterSourceMetricsSelectors, telemetrytypes.SignalMetrics)
logsKeys, logsComplete, err := t.getLogsKeys(ctx, orgID, logsSelectors)
if err != nil {
return nil, false, err
@@ -1542,7 +1594,16 @@ func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueS
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(t.tracesDBName + "." + t.tracesFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
members := traceSemconvMembers(fieldValueSelector.Name, fieldValueSelector.FieldContext)
if len(members) == 1 {
sb.Where(sb.E("tag_key", members[0]))
} else {
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
sb.Where(sb.In("tag_key", memberValues...))
}
}
// now look at the field context
@@ -1632,7 +1693,12 @@ func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, fieldValueSe
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(t.logsDBName + "." + t.logsFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
members := attributeSemconvMembers(fieldValueSelector.Name, telemetrytypes.SignalLogs, fieldValueSelector.FieldContext)
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
sb.Where(sb.In("tag_key", memberValues...))
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1830,7 +1896,9 @@ func (t *telemetryMetaStore) getMetricFieldValues(ctx context.Context, orgID val
From(t.metricsDBName + "." + t.metricsFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("attr_name", fieldValueSelector.Name))
sb.Where(inStrings(sb, "attr_name", attributeSemconvMembers(
fieldValueSelector.Name, telemetrytypes.SignalMetrics, fieldValueSelector.FieldContext,
)))
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1842,7 +1910,7 @@ func (t *telemetryMetaStore) getMetricFieldValues(ctx context.Context, orgID val
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
sb.Where(inStrings(sb, "metric_name", semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)))
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -1976,7 +2044,7 @@ func (t *telemetryMetaStore) getIntrinsicMetricFieldValuesForTable(ctx context.C
From(t.metricsDBName + "." + tableName)
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
sb.Where(inStrings(sb, "metric_name", semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)))
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2038,13 +2106,18 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
sb := sqlbuilder.Select("DISTINCT arrayJoin(JSONExtractKeysAndValues(labels, 'String')) AS attr").
From(t.meterDBName + "." + t.meterFieldsTblName)
duplicateFactor := 1
if fieldValueSelector.Name != "" {
sb.Where(sb.E("attr.1", fieldValueSelector.Name))
members := attributeSemconvMembers(fieldValueSelector.Name, telemetrytypes.SignalMetrics, fieldValueSelector.FieldContext)
duplicateFactor *= len(members)
sb.Where(inStrings(sb, "attr.1", members))
}
sb.Where(sb.NotLike("attr.1", "\\_\\_%"))
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
metricNames := semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)
duplicateFactor *= len(metricNames)
sb.Where(inStrings(sb, "metric_name", metricNames))
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2063,8 +2136,10 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
if limit == 0 {
limit = 50
}
// query one extra to check if we hit the limit
sb.Limit(limit + 1)
// A value can be present under several physical spellings; over-fetch and
// de-duplicate after scanning so one family cannot consume the result limit.
dbLimit := limit * duplicateFactor
sb.Limit(dbLimit + 1)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
@@ -2074,11 +2149,13 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
defer rows.Close()
values := &telemetrytypes.TelemetryFieldValues{}
seen := make(map[string]bool)
rowCount := 0
uniqueCount := 0
for rows.Next() {
rowCount++
// reached the limit, we know there are more results
if rowCount > limit {
if rowCount > dbLimit {
break
}
@@ -2087,12 +2164,16 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
return nil, false, errors.Wrap(err, errors.TypeInternal, errors.CodeInternal, ErrFailedToGetMeterValues.Error())
}
if len(attribute) > 1 {
values.StringValues = append(values.StringValues, attribute[1])
if !seen[attribute[1]] && uniqueCount < limit {
values.StringValues = append(values.StringValues, attribute[1])
seen[attribute[1]] = true
uniqueCount++
}
}
}
// hit the limit?
complete := rowCount <= limit
complete := rowCount <= dbLimit && uniqueCount < limit
return values, complete, nil
}

View File

@@ -12,6 +12,7 @@ import (
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -83,3 +84,38 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
t.Errorf("there were unfulfilled expectations: %s", err)
}
}
func TestGetAllValuesReturnsValuesFromEveryTraceSemconvFamilyMember(t *testing.T) {
mockTelemetryStore := telemetrystoretest.New(telemetrystore.Config{}, &regexMatcher{})
mock := mockTelemetryStore.Mock()
metadata := NewTelemetryMetaStore(
instrumentationtest.New().ToProviderSettings(),
mockTelemetryStore,
flaggertest.New(t),
)
mock.ExpectQuery(`SELECT DISTINCT string_value, number_value FROM signoz_traces\.distributed_tag_attributes_v2 WHERE tag_key IN \(\?, \?\) AND tag_type = \? AND tag_data_type = \? LIMIT \?`).
WithArgs("deployment.environment.name", "deployment.environment", "resource", "string", 51).
WillReturnRows(cmock.NewRows([]cmock.ColumnType{
{Name: "string_value", Type: "String"},
{Name: "number_value", Type: "Float64"},
}, [][]any{
{"production", float64(0)},
{"staging", float64(0)},
{"production", float64(0)},
}))
values, complete, err := metadata.GetAllValues(context.Background(), valuer.UUID{}, &telemetrytypes.FieldValueSelector{
FieldKeySelector: &telemetrytypes.FieldKeySelector{
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
Name: "deployment.environment",
},
})
require.NoError(t, err)
assert.True(t, complete)
assert.Equal(t, []string{"production", "staging"}, values.StringValues)
assert.NoError(t, mock.ExpectationsWereMet(), "all expected metadata queries should be executed")
}

View File

@@ -0,0 +1,156 @@
package telemetrymetadata
import (
"context"
"fmt"
"slices"
"strings"
"github.com/huandu/go-sqlbuilder"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
type semconvMigrationRow struct {
current string
old string
signal string
service string
resourceSets uint64
lastSeenUnixMilli int64
}
// GetSemconvMigrationReport derives an old-only service report from the
// generated family registry and attributes_metadata. The latter is already a
// deduplicated set of resource/attribute fingerprints, so this audit avoids a
// scan of the raw telemetry tables.
func (t *telemetryMetaStore) GetSemconvMigrationReport(
ctx context.Context,
_ valuer.UUID,
startUnixMilli, endUnixMilli int64,
) (*telemetrytypes.GettableSemconvMigrationReport, error) {
query, args := t.semconvMigrationReportQuery(startUnixMilli, endUnixMilli)
report := &telemetrytypes.GettableSemconvMigrationReport{
StartUnixMilli: startUnixMilli,
EndUnixMilli: endUnixMilli,
Entries: []*telemetrytypes.SemconvMigrationReportEntry{},
}
if query == "" {
return report, nil
}
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return nil, errors.Wrap(err, errors.TypeInternal, errors.CodeInternal, "failed to build semantic-convention migration report")
}
defer rows.Close()
grouped := make(map[string]*telemetrytypes.SemconvMigrationReportEntry)
serviceSets := make(map[string]map[string]struct{})
for rows.Next() {
var row semconvMigrationRow
if err := rows.Scan(
&row.current,
&row.old,
&row.signal,
&row.service,
&row.resourceSets,
&row.lastSeenUnixMilli,
); err != nil {
return nil, errors.Wrap(err, errors.TypeInternal, errors.CodeInternal, "failed to scan semantic-convention migration report")
}
identity := row.current + "\x00" + row.old + "\x00" + row.signal
entry, ok := grouped[identity]
if !ok {
entry = &telemetrytypes.SemconvMigrationReportEntry{
Current: row.current,
Old: row.old,
Signal: row.signal,
Services: []string{},
}
grouped[identity] = entry
serviceSets[identity] = make(map[string]struct{})
report.Entries = append(report.Entries, entry)
}
serviceSets[identity][row.service] = struct{}{}
entry.ResourceSets += row.resourceSets
entry.LastSeenUnixMilli = max(entry.LastSeenUnixMilli, row.lastSeenUnixMilli)
}
if err := rows.Err(); err != nil {
return nil, errors.Wrap(err, errors.TypeInternal, errors.CodeInternal, "failed to read semantic-convention migration report")
}
for identity, entry := range grouped {
for service := range serviceSets[identity] {
entry.Services = append(entry.Services, service)
}
slices.Sort(entry.Services)
}
slices.SortFunc(report.Entries, func(a, b *telemetrytypes.SemconvMigrationReportEntry) int {
return strings.Compare(a.Current+"\x00"+a.Old+"\x00"+a.Signal, b.Current+"\x00"+b.Old+"\x00"+b.Signal)
})
return report, nil
}
func (t *telemetryMetaStore) semconvMigrationReportQuery(startUnixMilli, endUnixMilli int64) (string, []any) {
builders := make([]sqlbuilder.Builder, 0)
for _, family := range semconv.All() {
if family.Kind != semconv.KindAttribute {
continue
}
for _, old := range family.Old {
sb := sqlbuilder.NewSelectBuilder()
sb.Select(
fmt.Sprintf("%s AS current_name", sb.Var(family.Current)),
fmt.Sprintf("%s AS old_name", sb.Var(old)),
"data_source",
"if(empty(resource_attributes['service.name']), '<unknown>', resource_attributes['service.name']) AS service_name",
"uniqExact(tuple(resource_fingerprint, attrs_fingerprint)) AS resource_sets",
"toInt64(max(unix_milli)) AS last_seen_unix_milli",
)
sb.From(t.relatedMetadataDBName + "." + t.relatedMetadataTblName)
sb.Where(sb.GE("unix_milli", startUnixMilli))
sb.Where(sb.LE("unix_milli", endUnixMilli))
if len(family.Signals) > 0 {
signals := make([]any, 0, len(family.Signals))
for _, signal := range family.Signals {
signals = append(signals, signal.StringValue())
}
sb.Where(sb.In("data_source", signals...))
}
oldPresence := semconvMetadataPresenceConditions(sb, old)
currentPresence := semconvMetadataPresenceConditions(sb, family.Current)
sb.Where(sb.Or(oldPresence...))
sb.Where(fmt.Sprintf("NOT (%s)", sb.Or(currentPresence...)))
sb.GroupBy("data_source", "service_name")
builders = append(builders, sb)
}
}
if len(builders) == 0 {
return "", nil
}
union := sqlbuilder.UnionAll(builders...)
return union.BuildWithFlavor(sqlbuilder.ClickHouse)
}
func semconvMetadataPresenceConditions(sb *sqlbuilder.SelectBuilder, name string) []string {
spellings := []string{name, strings.ReplaceAll(name, ".", "_")}
spellings = append(spellings, "resource_"+name, "resource_"+strings.ReplaceAll(name, ".", "_"))
spellings = slices.Compact(spellings)
conditions := make([]string, 0, len(spellings)*2)
for _, spelling := range spellings {
conditions = append(conditions,
fmt.Sprintf("mapContains(resource_attributes, %s)", sb.Var(spelling)),
fmt.Sprintf("mapContains(attributes, %s)", sb.Var(spelling)),
)
}
return conditions
}

View File

@@ -89,10 +89,12 @@ func (v *TelemetryFieldVisitor) VisitColumnDef(expr *parser.ColumnDef) error {
// Create and store the TelemetryFieldKey
field := &telemetrytypes.TelemetryFieldKey{
Name: fieldName,
FieldContext: fieldContext,
FieldDataType: fieldDataType,
Materialized: true,
Name: fieldName,
FieldContext: fieldContext,
FieldDataType: fieldDataType,
Materialized: true,
MaterializedColumnName: columnName,
MaterializedSemconv: strings.Count(defaultExprStr, "['") > 1,
}
v.Fields = append(v.Fields, field)

View File

@@ -5,8 +5,25 @@ import (
"testing"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestExtractFieldKeyPreservesHistoricalMaterializedColumnName(t *testing.T) {
statement := `CREATE TABLE signoz_traces.signoz_index_v3
(
attributes_string Map(LowCardinality(String), String),
` + "`attribute_string_db$$system`" + ` LowCardinality(String)
DEFAULT if(mapContains(attributes_string, 'db.system.name'), attributes_string['db.system.name'], attributes_string['db.system'])
) ENGINE = MergeTree ORDER BY tuple()`
keys, err := ExtractFieldKeysFromTblStatement(statement)
require.NoError(t, err, "table statement should parse")
require.Len(t, keys, 1, "table statement should contain one materialized key")
assert.Equal(t, "db.system.name", keys[0].Name)
assert.Equal(t, "attribute_string_db$$system", keys[0].MaterializedColumnName)
assert.True(t, keys[0].MaterializedSemconv)
}
func TestExtractFieldKeysFromTblStatement(t *testing.T) {
var statement = `CREATE TABLE signoz_logs.logs_v2

View File

@@ -453,6 +453,9 @@ func (c *conditionBuilder) ConditionFor(
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
if options.ExactSemconv {
matches = querybuilder.MatchingFieldKeysExact(key, fieldKeys)
}
skipResourceFilter := options.SkipResourceFilter
// search() resolves its own (optional) scope; handle it before key resolution.
@@ -499,6 +502,9 @@ func (c *conditionBuilder) ConditionFor(
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
}
if options.ExactSemconv {
keys = querybuilder.ExactSemconvKeys(keys)
}
if skipResourceFilter && !synthesized {
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))

View File

@@ -10,6 +10,7 @@ import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -67,6 +68,20 @@ type fieldMapper struct {
fl flagger.Flagger
}
func logSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
return []string{key.Name}
}
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalLogs,
FieldContext: key.FieldContext,
})
}
func NewFieldMapper(fl flagger.Flagger) qbtypes.FieldMapper {
return &fieldMapper{fl: fl}
}
@@ -141,8 +156,22 @@ func (m *fieldMapper) FieldFor(ctx context.Context, orgID valuer.UUID, tsStart,
case schema.ColumnTypeEnumJSON:
switch key.FieldContext {
case telemetrytypes.FieldContextResource:
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExpr = append(existExpr, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
members := logSemconvMembers(key)
if len(members) > 1 {
values := make([]string, 0, len(members))
guards := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(%s.`%s`::String, '')", columnName, member))
guards = append(guards, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, member))
}
// Missing Dynamic paths are NULL, so this family expression must
// retain the same NULL result as a single JSON-path lookup.
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+")")
existExpr = append(existExpr, "("+strings.Join(guards, " OR ")+")")
} else {
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, members[0]))
existExpr = append(existExpr, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, members[0]))
}
case telemetrytypes.FieldContextBody:
if key.Name == messageSubField {
exprs = append(exprs, messageSubColumn)
@@ -181,13 +210,34 @@ func (m *fieldMapper) FieldFor(ctx context.Context, orgID valuer.UUID, tsStart,
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumBool, schema.ColumnTypeEnumFloat64:
// a key could have been materialized, if so return the materialized column name
if key.Materialized {
members := logSemconvMembers(key)
if key.Materialized && (len(members) == 1 || key.MaterializedSemconv) {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExpr = append(existExpr, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
} else if len(members) > 1 {
guards := make([]string, 0, len(members))
for _, member := range members {
guards = append(guards, fmt.Sprintf("mapContains(%s, '%s')", columnName, member))
}
if valueType.GetType() == schema.ColumnTypeEnumString {
values := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(%s['%s'], '')", columnName, member))
}
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+", '')")
} else {
branches := make([]string, 0, len(members)*2+1)
for i, member := range members {
branches = append(branches, guards[i], fmt.Sprintf("%s['%s']", columnName, member))
}
// Numeric and boolean maps return zero for an absent key. If a
// family of either type is enabled, this tail must become zero too.
exprs = append(exprs, "multiIf("+strings.Join(branches, ", ")+", NULL)")
}
existExpr = append(existExpr, "("+strings.Join(guards, " OR ")+")")
} else {
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, key.Name))
existExpr = append(existExpr, fmt.Sprintf("mapContains(%s, '%s')", columnName, key.Name))
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, members[0]))
existExpr = append(existExpr, fmt.Sprintf("mapContains(%s, '%s')", columnName, members[0]))
}
default:
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "exists operator is not supported for map column type %s", valueType)

View File

@@ -2,6 +2,7 @@ package logstelemetryschema
import (
"context"
"strings"
"testing"
"time"
@@ -14,6 +15,32 @@ import (
"github.com/stretchr/testify/require"
)
func TestFieldForSemconvFamily(t *testing.T) {
ctx := context.Background()
fm := NewFieldMapper(flaggertest.New(t)).(*fieldMapper)
key := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := fm.FieldFor(ctx, valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
assert.Equal(t, "COALESCE(NULLIF(attributes_string['db.system.name'], ''), NULLIF(attributes_string['db.system'], ''), '')", expression)
assert.Less(t, strings.Index(expression, "db.system.name"), strings.Index(expression, "db.system']"), "current spelling must win")
exists, err := fm.existsExpressionFor(ctx, valuer.UUID{}, 0, 0, key, true)
require.NoError(t, err)
assert.Equal(t, "(mapContains(attributes_string, 'db.system.name') OR mapContains(attributes_string, 'db.system'))", exists)
exact := *key
exact.SemconvMembers = []string{"db.system"}
expression, err = fm.FieldFor(ctx, valuer.UUID{}, 0, 0, &exact)
require.NoError(t, err)
assert.Equal(t, "attributes_string['db.system']", expression)
}
func TestGetColumn(t *testing.T) {
ctx := context.Background()

View File

@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/querybuilder"
@@ -135,10 +136,22 @@ func (c *conditionBuilder) conditionFor(
return "true", nil
}
if operator == qbtypes.FilterOperatorExists {
return fmt.Sprintf("has(JSONExtractKeys(labels), '%s')", key.Name), nil
members := metricAttributeMembers(key)
guards := make([]string, 0, len(members))
for _, member := range members {
guards = append(guards, fmt.Sprintf("has(JSONExtractKeys(labels), '%s')", member))
}
return fmt.Sprintf("not has(JSONExtractKeys(labels), '%s')", key.Name), nil
guard := strings.Join(guards, " OR ")
if len(guards) > 1 {
guard = "(" + guard + ")"
}
if operator == qbtypes.FilterOperatorExists {
return guard, nil
}
if len(guards) == 1 {
return "not " + guard, nil
}
return "NOT " + guard, nil
}
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported operator: %v", operator)
}
@@ -151,7 +164,7 @@ func (c *conditionBuilder) ConditionFor(
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey,
_ qbtypes.ConditionBuilderOptions,
options qbtypes.ConditionBuilderOptions,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -162,7 +175,15 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
requestedKey := *key
if requestedKey.Signal == telemetrytypes.SignalUnspecified {
requestedKey.Signal = telemetrytypes.SignalMetrics
}
key = &requestedKey
keys := querybuilder.MatchingFieldKeys(key, fieldKeys)
if options.ExactSemconv {
keys = querybuilder.MatchingFieldKeysExact(key, fieldKeys)
}
var warnings []string
if len(keys) == 0 {
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {
@@ -180,6 +201,9 @@ func (c *conditionBuilder) ConditionFor(
}
}
}
if options.ExactSemconv {
keys = querybuilder.ExactSemconvKeys(keys)
}
conds := make([]string, 0, len(keys))
for _, k := range keys {

View File

@@ -390,3 +390,45 @@ func TestConditionForKeyNotInMetadata(t *testing.T) {
})
}
}
func TestConditionForSemconvMetricLabels(t *testing.T) {
ctx := context.Background()
fm := NewFieldMapper()
conditionBuilder := NewConditionBuilder(fm)
requested := telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
current := requested
legacyNormalized := requested
legacyNormalized.Name = "resource_db_system"
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {&current},
legacyNormalized.Name: {&legacyNormalized},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, warnings, err := conditionBuilder.ConditionFor(
ctx, valuer.UUID{}, 0, 0, &requested, fieldKeys, qbtypes.ConditionBuilderOptions{},
qbtypes.FilterOperatorEqual, "postgresql", sb,
)
require.NoError(t, err)
assert.Empty(t, warnings)
sb.Where(conditions...)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, query, "COALESCE(NULLIF(JSONExtractString(labels, 'db.system.name'), ''), NULLIF(JSONExtractString(labels, 'resource_db_system'), ''), '') = ?")
assert.Equal(t, []any{"postgresql"}, args)
sb = sqlbuilder.NewSelectBuilder()
conditions, _, err = conditionBuilder.ConditionFor(
ctx, valuer.UUID{}, 0, 0, &requested, fieldKeys, qbtypes.ConditionBuilderOptions{ExactSemconv: true},
qbtypes.FilterOperatorExists, nil, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
query, _ = sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, query, "has(JSONExtractKeys(labels), 'db.system.name')")
assert.NotContains(t, query, "resource_db_system")
}

View File

@@ -4,8 +4,10 @@ import (
"context"
"fmt"
"slices"
"strings"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -38,6 +40,23 @@ var (
type fieldMapper struct{}
func metricAttributeMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if key.FieldContext != telemetrytypes.FieldContextResource &&
key.FieldContext != telemetrytypes.FieldContextScope &&
key.FieldContext != telemetrytypes.FieldContextAttribute &&
key.FieldContext != telemetrytypes.FieldContextUnspecified {
return []string{key.Name}
}
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: key.FieldContext,
})
}
// CandidateKeys returns nil: metrics has no attribute-map fallback, so a context-missing
// key stays unresolved and the caller errors.
func (m *fieldMapper) CandidateKeys(_ context.Context, _ valuer.UUID, _ *telemetrytypes.TelemetryFieldKey, _ any, _ map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
@@ -80,19 +99,30 @@ func (m *fieldMapper) FieldFor(ctx context.Context, _ valuer.UUID, startNs, endN
switch key.FieldContext {
case telemetrytypes.FieldContextResource, telemetrytypes.FieldContextScope, telemetrytypes.FieldContextAttribute:
return fmt.Sprintf("JSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
return metricLabelExpression(columns[0].Name, metricAttributeMembers(key)), nil
case telemetrytypes.FieldContextMetric:
return columns[0].Name, nil
case telemetrytypes.FieldContextUnspecified:
if slices.Contains(IntrinsicFields, key.Name) {
return columns[0].Name, nil
}
return fmt.Sprintf("JSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
return metricLabelExpression(columns[0].Name, metricAttributeMembers(key)), nil
}
return columns[0].Name, nil
}
func metricLabelExpression(columnName string, members []string) string {
if len(members) == 1 {
return fmt.Sprintf("JSONExtractString(%s, '%s')", columnName, members[0])
}
values := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(JSONExtractString(%s, '%s'), '')", columnName, member))
}
return "COALESCE(" + strings.Join(values, ", ") + ", '')"
}
func (m *fieldMapper) ColumnFor(ctx context.Context, _ valuer.UUID, tsStart, tsEnd uint64, key *telemetrytypes.TelemetryFieldKey) ([]*schema.Column, error) {
return m.getColumn(ctx, tsStart, tsEnd, key)
}

View File

@@ -2,6 +2,7 @@ package metricstelemetryschema
import (
"context"
"strings"
"testing"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
@@ -12,6 +13,32 @@ import (
"github.com/stretchr/testify/require"
)
func TestMetricLabelSemconvSpellings(t *testing.T) {
fm := NewFieldMapper()
key := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := fm.FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
for _, member := range []string{
"resource_db.system.name", "resource_db_system_name", "db.system.name", "db_system_name",
"resource_db.system", "resource_db_system", "db.system", "db_system",
} {
assert.Contains(t, expression, "'"+member+"'")
}
assert.Less(t, strings.Index(expression, "resource_db.system.name"), strings.Index(expression, "resource_db.system'"))
exact := *key
exact.SemconvMembers = []string{"resource_db_system"}
expression, err = fm.FieldFor(context.Background(), valuer.UUID{}, 0, 0, &exact)
require.NoError(t, err)
assert.Equal(t, "JSONExtractString(labels, 'resource_db_system')", expression)
}
func TestGetColumn(t *testing.T) {
ctx := context.Background()

View File

@@ -18,12 +18,12 @@ import (
)
type conditionBuilder struct {
fm qbtypes.FieldMapper
fm *fieldMapper
}
var _ qbtypes.ConditionBuilder = (*conditionBuilder)(nil)
func NewConditionBuilder(fm qbtypes.FieldMapper) *conditionBuilder {
func NewConditionBuilder(fm *fieldMapper) *conditionBuilder {
return &conditionBuilder{fm: fm}
}
@@ -154,6 +154,17 @@ func (c *conditionBuilder) conditionFor(
// in the query builder, `exists` and `not exists` are used for
// key membership checks, so depending on the column type, the condition changes
case qbtypes.FilterOperatorExists, qbtypes.FilterOperatorNotExists:
// A semantic-convention family is represented by one current-first value
// expression, but presence still has to inspect every physical member. In
// particular, using ExistsExpression below with the requested key would add
// a mapContains check for only that spelling and reject fallback-only rows.
if isTraceSemconvFamily(key) {
pred, err := c.fm.existsExpressionFor(ctx, orgID, startNs, endNs, key, operator == qbtypes.FilterOperatorExists)
if err != nil {
return "", err
}
return sqlbuilder.Escape(pred), nil
}
columns, err := c.fm.ColumnFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err
@@ -211,6 +222,9 @@ func (c *conditionBuilder) ConditionFor(
}
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
if options.ExactSemconv {
matches = querybuilder.MatchingFieldKeysExact(key, fieldKeys)
}
skipResourceFilter := options.SkipResourceFilter
keys, warning := querybuilder.ResolveKeys(key, matches)
@@ -256,6 +270,9 @@ func (c *conditionBuilder) ConditionFor(
synthesized = true
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
if options.ExactSemconv {
keys = querybuilder.ExactSemconvKeys(keys)
}
// When a resource sub-query already covers the term, drop resource keys from the main
// query. Synthesized keys are exempt: the sub-query skips keys absent from metadata.

View File

@@ -308,6 +308,82 @@ func TestConditionFor(t *testing.T) {
}
}
func TestConditionForSemconvFamilyPositiveFilterChecksPresence(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Empty(t, warnings)
assert.Contains(t, sql, "(COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '') = ? AND ((mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))))")
assert.Equal(t, []any{"production"}, args)
}
func TestConditionForSemconvFamilyPreservesMaterializedMemberExistsColumn(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
SemconvMaterializedColumns: map[string]string{
"deployment.environment": "attribute_string_deployment$$environment",
},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Empty(t, warnings)
assert.Contains(t, sql, "`attribute_string_deployment$$environment_exists`")
assert.NotContains(t, sql, "`attribute_string_deployment$environment_exists`")
assert.Equal(t, []any{"production"}, args)
}
func TestConditionForSemconvFamilyNotExistsChecksEveryMember(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotExists, nil, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Empty(t, warnings)
assert.Contains(t, sql, "NOT (((mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))))")
assert.Empty(t, args)
}
func TestConditionForResourceWithEvolution(t *testing.T) {
ctx := context.Background()
releaseTime := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)

View File

@@ -8,6 +8,7 @@ import (
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -167,6 +168,44 @@ func NewFieldMapper() *fieldMapper {
return &fieldMapper{}
}
func traceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
return []string{key.Name}
}
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: key.FieldContext,
})
}
func isTraceSemconvFamily(key *telemetrytypes.TelemetryFieldKey) bool {
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
return false
}
_, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: key.FieldContext,
})
return ok
}
func traceSemconvMapMemberExpressions(columnName string, key *telemetrytypes.TelemetryFieldKey, member string) (string, string) {
if materializedColumn, ok := key.SemconvMaterializedColumns[member]; ok {
return fmt.Sprintf("`%s`", materializedColumn), fmt.Sprintf("`%s_exists`", materializedColumn)
}
if key.Materialized && key.Name == member {
physicalKey := *key
physicalKey.Name = member
return telemetrytypes.FieldKeyToMaterializedColumnName(&physicalKey), telemetrytypes.FieldKeyToMaterializedColumnNameForExists(&physicalKey)
}
return fmt.Sprintf("%s['%s']", columnName, member), fmt.Sprintf("mapContains(%s, '%s')", columnName, member)
}
func (m *fieldMapper) getColumn(
_ context.Context,
_, _ uint64,
@@ -291,10 +330,27 @@ func (m *fieldMapper) resolveColumnExprs(
if key.FieldContext != telemetrytypes.FieldContextResource {
return nil, nil, nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "only resource context fields are supported for json columns, got %s", key.FieldContext.String)
}
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
// once clickHouse dependency is updated, we need to check if we can remove it.
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
members := traceSemconvMembers(key)
if len(members) > 1 {
values := make([]string, 0, len(members))
guards := make([]string, 0, len(members))
for _, member := range members {
// The String cast is required because ClickHouse does not allow
// Variant/Dynamic values in GROUP BY.
value := fmt.Sprintf("%s.`%s`::String", columnName, member)
values = append(values, fmt.Sprintf("NULLIF(%s, '')", value))
guards = append(guards, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, member))
}
// Missing Dynamic paths are NULL, so this family expression must
// retain the same NULL result as a single JSON-path lookup.
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+")")
existExprs = append(existExprs, "("+strings.Join(guards, " OR ")+")")
} else {
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
// once ClickHouse is updated, check whether this cast can be removed.
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, members[0]))
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, members[0]))
}
case schema.ColumnTypeEnumString,
schema.ColumnTypeEnumUInt64,
schema.ColumnTypeEnumUInt32,
@@ -319,13 +375,37 @@ func (m *fieldMapper) resolveColumnExprs(
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumFloat64, schema.ColumnTypeEnumBool:
// a key could have been materialized, if so return the materialized column name
if key.Materialized {
members := traceSemconvMembers(key)
if key.Materialized && (len(members) == 1 || key.MaterializedSemconv) {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
} else if len(members) > 1 {
guards := make([]string, 0, len(members))
memberValues := make([]string, 0, len(members))
for _, member := range members {
valueExpression, existsExpression := traceSemconvMapMemberExpressions(columnName, key, member)
memberValues = append(memberValues, valueExpression)
guards = append(guards, existsExpression)
}
if valueType.GetType() == schema.ColumnTypeEnumString {
values := make([]string, 0, len(members))
for _, memberValue := range memberValues {
values = append(values, fmt.Sprintf("NULLIF(%s, '')", memberValue))
}
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+", '')")
} else {
branches := make([]string, 0, len(members)*2)
for i, memberValue := range memberValues {
branches = append(branches, guards[i], memberValue)
}
// Numeric and boolean maps return zero for an absent key. If a
// family of either type is enabled, this tail must become zero too.
exprs = append(exprs, "multiIf("+strings.Join(branches, ", ")+", NULL)")
}
existExprs = append(existExprs, "("+strings.Join(guards, " OR ")+")")
} else {
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, key.Name))
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, members[0]))
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, members[0]))
}
default:
return nil, nil, nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "value type %s is not supported for map column type %s", valueType, column.Type)
@@ -529,6 +609,25 @@ func (m *fieldMapper) existsExpressionFor(
key *telemetrytypes.TelemetryFieldKey,
exists bool,
) (string, error) {
if isTraceSemconvFamily(key) {
_, existExprs, _, err := m.resolveColumnExprs(ctx, tsStart, tsEnd, key)
if err != nil {
return "", err
}
if len(existExprs) == 0 {
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "no existence expression found for field %s", key.Name)
}
parts := make([]string, 0, len(existExprs))
for _, expression := range existExprs {
parts = append(parts, "("+expression+")")
}
combined := strings.Join(parts, " OR ")
if exists {
return combined, nil
}
return "NOT (" + combined + ")", nil
}
columns, err := m.getColumn(ctx, tsStart, tsEnd, key)
if err != nil {
return "", err

View File

@@ -5,6 +5,8 @@ import (
"testing"
"time"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/querybuilder"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -80,7 +82,7 @@ func TestGetFieldKeyName(t *testing.T) {
Materialized: true,
Evolutions: mockEvolution,
},
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
expectedResult: "multiIf((resource.`deployment.environment.name` IS NOT NULL OR resource.`deployment.environment` IS NOT NULL), COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, '')), (mapContains(resources_string, 'deployment.environment.name') OR `resource_string_deployment$$environment_exists`), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(`resource_string_deployment$$environment`, ''), ''), NULL)",
expectedError: nil,
},
{
@@ -120,6 +122,63 @@ func TestGetFieldKeyName(t *testing.T) {
}
}
func TestFieldForResolvesCurrentTraceSemconvAttributeName(t *testing.T) {
key := telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
require.NoError(t, err)
assert.Equal(t, "COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '')", expression)
}
func TestFieldForResolvesOldTraceSemconvAttributeName(t *testing.T) {
key := telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
require.NoError(t, err)
assert.Equal(t, "COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '')", expression)
}
func TestFieldForPreservesResourceStorageDefaultsForSemconvFamily(t *testing.T) {
key := telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
Evolutions: MockEvolutionData(time.Date(2024, 6, 2, 0, 0, 0, 0, time.UTC)),
}
start := uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano())
end := uint64(time.Date(2024, 6, 5, 0, 0, 0, 0, time.UTC).UnixNano())
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, start, end, &key)
require.NoError(t, err)
assert.Equal(t, "multiIf((resource.`deployment.environment.name` IS NOT NULL OR resource.`deployment.environment` IS NOT NULL), COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, '')), (`resource_string_deployment$$environment$$name_exists` OR mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(`resource_string_deployment$$environment$$name`, ''), NULLIF(resources_string['deployment.environment'], ''), ''), NULL)", expression)
}
func TestFieldForUsesAvailableTraceSemconvMember(t *testing.T) {
key := telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name"},
}
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
require.NoError(t, err)
assert.Equal(t, "attributes_string['deployment.environment.name']", expression)
}
func TestFieldForResourceWithEvolution(t *testing.T) {
ctx := context.Background()
releaseTime := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)
@@ -176,7 +235,7 @@ func TestFieldForResourceWithEvolution(t *testing.T) {
},
tsStart: uint64(time.Date(2025, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
tsEnd: uint64(time.Date(2025, 7, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
expectedResult: "resource.`deployment.environment`::String",
expectedResult: "COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, ''))",
},
{
name: "Window straddles release - materialized resource",
@@ -189,7 +248,7 @@ func TestFieldForResourceWithEvolution(t *testing.T) {
},
tsStart: uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
tsEnd: uint64(time.Date(2025, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
expectedResult: "multiIf((resource.`deployment.environment.name` IS NOT NULL OR resource.`deployment.environment` IS NOT NULL), COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, '')), (mapContains(resources_string, 'deployment.environment.name') OR `resource_string_deployment$$environment_exists`), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(`resource_string_deployment$$environment`, ''), ''), NULL)",
},
}
@@ -303,3 +362,27 @@ func TestColumnExpressionForTimestampAttributeCollision(t *testing.T) {
assert.Contains(t, result, "attributes_number['timestamp']")
})
}
func TestDBSystemFamilyUsesSemconvAwareMaterializedColumn(t *testing.T) {
fm := NewFieldMapper()
key := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
MaterializedColumnName: "attribute_string_db$$system",
MaterializedSemconv: true,
SemconvMembers: []string{"db.system.name", "db.system"},
}
expression, err := fm.FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
assert.Equal(t, "`attribute_string_db$$system`", expression)
exists, err := querybuilder.ExistsExpression(
[]*schema.Column{indexV3Columns["attributes_string"]}, key, 0, 0, expression, true,
)
require.NoError(t, err)
assert.Equal(t, "`attribute_string_db$$system_exists`", exists)
}

View File

@@ -155,6 +155,8 @@ var operatorInverseMapping = map[FilterOperator]FilterOperator{
// doesn't have value "redis"
// Since we don't know the intent, we don't add the exists filter. They are expected
// to add exists filter themselves if exclusion is desired.
// Negative predicates therefore include rows where the key is absent; value
// expressions must preserve the storage column's absent-key default.
//
// For the positive predicates, the key existence is implied.
func (f FilterOperator) AddDefaultExistsFilter() bool {

View File

@@ -45,6 +45,10 @@ type ConditionBuilder interface {
type ConditionBuilderOptions struct {
// SkipResourceFilter drops the resource context from the candidate set.
SkipResourceFilter bool
// ExactSemconv disables semantic-convention family expansion. It is an
// internal escape hatch for diagnostics and migrations that must address one
// physical spelling only; public query APIs continue to resolve families.
ExactSemconv bool
}
type AggExprRewriter interface {
// Rewrite rewrites the aggregation expression to be used in the query.

View File

@@ -25,10 +25,22 @@ type Result struct {
}
type ExecStats struct {
RowsScanned uint64 `json:"rowsScanned"`
BytesScanned uint64 `json:"bytesScanned"`
DurationMS uint64 `json:"durationMs"`
StepIntervals map[string]uint64 `json:"stepIntervals,omitempty"`
RowsScanned uint64 `json:"rowsScanned"`
BytesScanned uint64 `json:"bytesScanned"`
DurationMS uint64 `json:"durationMs"`
StepIntervals map[string]uint64 `json:"stepIntervals,omitempty"`
SemconvResolutions []SemconvResolution `json:"semconvResolutions,omitempty"`
}
// SemconvResolution records a semantic-convention family that the query
// builder resolved. Requested preserves the spelling supplied by the caller so
// agents and editors can update their next query without changing response
// labels in the current response.
type SemconvResolution struct {
Requested string `json:"requested"`
Current string `json:"current"`
Members []string `json:"members"`
Kind string `json:"kind"`
}
var _ jsonschema.Preparer = &ExecStats{}

View File

@@ -141,7 +141,7 @@ func NewSignalFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFi
func NewDefaultQuickFilter(orgID valuer.UUID) ([]*StorableQuickFilter, error) {
tracesFilters := []map[string]interface{}{
{"key": "duration_nano", "dataType": "float64", "type": "tag"},
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "deployment.environment.name", "dataType": "string", "type": "resource"},
{"key": "hasError", "dataType": "bool", "type": "tag"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "name", "dataType": "string", "type": "tag"},
@@ -166,13 +166,13 @@ func NewDefaultQuickFilter(orgID valuer.UUID) ([]*StorableQuickFilter, error) {
}
apiMonitoringFilters := []map[string]interface{}{
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "deployment.environment.name", "dataType": "string", "type": "resource"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "rpc.method", "dataType": "string", "type": "tag"},
}
exceptionsFilters := []map[string]interface{}{
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "deployment.environment.name", "dataType": "string", "type": "resource"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "host.name", "dataType": "string", "type": "resource"},
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},

View File

@@ -0,0 +1,37 @@
package quickfiltertypes
import (
"encoding/json"
"testing"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestDefaultTraceQuickFiltersUseCurrentEnvironmentName(t *testing.T) {
filters, err := NewDefaultQuickFilter(valuer.GenerateUUID())
require.NoError(t, err)
traceSignals := map[string]bool{
SignalTraces.StringValue(): true,
SignalApiMonitoring.StringValue(): true,
SignalExceptions.StringValue(): true,
}
for _, filter := range filters {
if !traceSignals[filter.Signal.StringValue()] {
continue
}
var keys []v3.AttributeKey
require.NoError(t, json.Unmarshal([]byte(filter.Filter), &keys))
found := false
for _, key := range keys {
if key.Key == "deployment.environment.name" {
found = true
}
assert.NotEqual(t, "deployment.environment", key.Key)
}
assert.True(t, found, "missing environment quick filter for %s", filter.Signal.StringValue())
}
}

View File

@@ -46,8 +46,19 @@ type TelemetryFieldKey struct {
JSONPlan JSONAccessPlan `json:"-"`
Indexes []TelemetryFieldKeySkipIndex `json:"-"`
Materialized bool `json:"-"` // refers to promoted in case of body.... fields
// MaterializedColumnName preserves the physical column when its DEFAULT
// expression resolves a newer semantic-convention key than the historical
// column identifier (for example db.system.name in db$$system).
MaterializedColumnName string `json:"-"`
// MaterializedSemconv is true when the column's DEFAULT expression already
// coalesces every enabled family member and is therefore safe for a family query.
MaterializedSemconv bool `json:"-"`
Evolutions []*EvolutionEntry `json:"-"`
Evolutions []*EvolutionEntry `json:"-"`
SemconvMembers []string `json:"-"`
// SemconvMaterializedColumns maps a physical family spelling to its
// materialized column name. It is populated only on resolved query keys.
SemconvMaterializedColumns map[string]string `json:"-"`
}
func (f *TelemetryFieldKey) KeyNameContainsArray() bool {
@@ -126,8 +137,12 @@ func (f *TelemetryFieldKey) OverrideMetadataFrom(src *TelemetryFieldKey) {
f.FieldDataType = src.FieldDataType
f.Indexes = src.Indexes
f.Materialized = src.Materialized
f.MaterializedColumnName = src.MaterializedColumnName
f.MaterializedSemconv = src.MaterializedSemconv
f.JSONPlan = src.JSONPlan
f.Evolutions = src.Evolutions
f.SemconvMembers = src.SemconvMembers
f.SemconvMaterializedColumns = src.SemconvMaterializedColumns
}
func (f *TelemetryFieldKey) Equal(key *TelemetryFieldKey) bool {
@@ -202,6 +217,9 @@ func TelemetryFieldKeyToText(key *TelemetryFieldKey) string {
}
func FieldKeyToMaterializedColumnName(key *TelemetryFieldKey) string {
if key.MaterializedColumnName != "" {
return fmt.Sprintf("`%s`", key.MaterializedColumnName)
}
return fmt.Sprintf("`%s_%s_%s`",
key.FieldContext.String,
fieldDataTypes[key.FieldDataType.StringValue()].StringValue(),
@@ -210,6 +228,9 @@ func FieldKeyToMaterializedColumnName(key *TelemetryFieldKey) string {
}
func FieldKeyToMaterializedColumnNameForExists(key *TelemetryFieldKey) string {
if key.MaterializedColumnName != "" {
return fmt.Sprintf("`%s_exists`", key.MaterializedColumnName)
}
return fmt.Sprintf("`%s_%s_%s_exists`",
key.FieldContext.String,
fieldDataTypes[key.FieldDataType.StringValue()].StringValue(),
@@ -276,6 +297,28 @@ type GettableFieldValues struct {
Complete bool `json:"complete" required:"true"`
}
// PostableSemconvMigrationReportParams selects the metadata window used to
// find services that still emit only historical semantic-convention names.
type PostableSemconvMigrationReportParams struct {
StartUnixMilli int64 `query:"startUnixMilli"`
EndUnixMilli int64 `query:"endUnixMilli"`
}
type SemconvMigrationReportEntry struct {
Current string `json:"current"`
Old string `json:"old"`
Signal string `json:"signal"`
Services []string `json:"services"`
ResourceSets uint64 `json:"resourceSets"`
LastSeenUnixMilli int64 `json:"lastSeenUnixMilli"`
}
type GettableSemconvMigrationReport struct {
StartUnixMilli int64 `json:"startUnixMilli"`
EndUnixMilli int64 `json:"endUnixMilli"`
Entries []*SemconvMigrationReportEntry `json:"entries" required:"true"`
}
type PostableFieldValueParams struct {
PostableFieldKeysParams
Name string `query:"name"`

View File

@@ -26,6 +26,10 @@ type MetadataStore interface {
// GetAllValues returns a list of all values.
GetAllValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *FieldValueSelector) (*TelemetryFieldValues, bool, error)
// GetSemconvMigrationReport returns services whose metadata contains an old
// semantic-convention name but no current member of that family.
GetSemconvMigrationReport(ctx context.Context, orgID valuer.UUID, startUnixMilli, endUnixMilli int64) (*GettableSemconvMigrationReport, error)
// FetchTemporality fetches the temporality for metric
FetchTemporality(ctx context.Context, orgID valuer.UUID, queryTimeRangeStartTs, queryTimeRangeEndTs uint64, metricName string) (metrictypes.Temporality, error)

View File

@@ -23,7 +23,19 @@ type MockMetadataStore struct {
ColumnEvolutionMetadataMap map[string][]*telemetrytypes.EvolutionEntry
LookupKeysMap map[telemetrytypes.MetricMetadataLookupKey]int64
// StaticFields holds signal-specific intrinsic field definitions (e.g. logstelemetryschema.IntrinsicFields).
StaticFields map[string]telemetrytypes.TelemetryFieldKey
StaticFields map[string]telemetrytypes.TelemetryFieldKey
SemconvMigrationReport *telemetrytypes.GettableSemconvMigrationReport
}
func (m *MockMetadataStore) GetSemconvMigrationReport(_ context.Context, _ valuer.UUID, startUnixMilli, endUnixMilli int64) (*telemetrytypes.GettableSemconvMigrationReport, error) {
if m.SemconvMigrationReport != nil {
return m.SemconvMigrationReport, nil
}
return &telemetrytypes.GettableSemconvMigrationReport{
StartUnixMilli: startUnixMilli,
EndUnixMilli: endUnixMilli,
Entries: []*telemetrytypes.SemconvMigrationReportEntry{},
}, nil
}
// NewMockMetadataStore creates a new instance of MockMetadataStore with initialized maps.

View File

@@ -7,5 +7,31 @@ default_enabled: false
families:
deployment.environment.name:
enabled: true
contexts: [resource, attribute]
signals: [traces, logs, metrics]
db.system.name:
enabled: true
contexts: [resource, attribute]
signals: [traces, logs, metrics]
# These metric renames predate the schema history vendored above. Keep them
# in the same generated registry so every v5 metric query uses one source of
# truth instead of the legacy hand-written transition table.
k8s.pod.cpu.usage:
enabled: true
kind: metric
old: [k8s.pod.cpu.utilization]
contexts: [metric]
signals: [metrics]
k8s.node.cpu.usage:
enabled: true
kind: metric
old: [k8s.node.cpu.utilization]
contexts: [metric]
signals: [metrics]
container.cpu.usage:
enabled: true
kind: metric
old: [container.cpu.utilization]
contexts: [metric]
signals: [metrics]

View File

@@ -34,6 +34,7 @@ pytest_plugins = [
"fixtures.role",
"fixtures.savedview",
"fixtures.seed_golden_dataset",
"fixtures.semconv",
]

83
tests/fixtures/semconv.py vendored Normal file
View File

@@ -0,0 +1,83 @@
from collections.abc import Callable, Generator
from datetime import UTC, datetime, timedelta
import pytest
from fixtures import types
from fixtures.metadata import AttributesMetadata
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
SEMCONV_PHASE1_CURRENT = "deployment.environment.name"
SEMCONV_PHASE1_OLD = "deployment.environment"
SEMCONV_PHASE1_PREFIX = "semconv-phase1"
@pytest.fixture(name="semconv_phase1_data")
def semconv_phase1_data(
insert_traces: Callable[[list[Traces]], None],
insert_attributes_metadata: Callable[[list[AttributesMetadata]], None],
clickhouse: types.TestContainerClickhouse,
) -> Generator[datetime]:
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=2)
records = [
(now - timedelta(seconds=5), "old", {SEMCONV_PHASE1_OLD: "production"}),
(now - timedelta(seconds=4), "current", {SEMCONV_PHASE1_CURRENT: "production"}),
(now - timedelta(seconds=3), "both", {SEMCONV_PHASE1_OLD: "production", SEMCONV_PHASE1_CURRENT: "production"}),
(now - timedelta(seconds=2), "conflict", {SEMCONV_PHASE1_OLD: "staging", SEMCONV_PHASE1_CURRENT: "production"}),
(now - timedelta(seconds=1), "staging", {SEMCONV_PHASE1_OLD: "staging"}),
(now, "missing", {}),
]
traces = []
for timestamp, suffix, environment in records:
service = f"{SEMCONV_PHASE1_PREFIX}-{suffix}"
traces.append(
Traces(
timestamp=timestamp,
duration=timedelta(milliseconds=10),
trace_id=TraceIdGenerator.trace_id(),
span_id=TraceIdGenerator.span_id(),
name=service,
kind=TracesKind.SPAN_KIND_SERVER,
status_code=TracesStatusCode.STATUS_CODE_OK,
resources={"service.name": service, **environment},
attributes=dict(environment),
)
)
insert_traces(traces)
# The production collector writes this deduplicated metadata table. Trace
# fixtures insert storage rows directly, so mirror that write explicitly
# to exercise the migration report against the same mixed-generation data.
insert_attributes_metadata(
[
AttributesMetadata(
data_source="traces",
resource_attributes={"service.name": f"{SEMCONV_PHASE1_PREFIX}-{suffix}", **environment},
attributes=environment,
timestamp=timestamp,
)
for timestamp, suffix, environment in records
]
)
# Service-map rows are derived by the collector in production. Seed the
# derived table directly so this test isolates the backend alias allowlist;
# the collector repository owns its write-path integration test.
for environment, suffix in (("production", "production"), ("staging", "staging")):
clickhouse.conn.command(
f"""
INSERT INTO signoz_traces.distributed_dependency_graph_minutes_v2
(src, dest, duration_quantiles_state, error_count, total_count, timestamp,
deployment_environment, k8s_cluster_name, k8s_namespace_name)
SELECT
'{SEMCONV_PHASE1_PREFIX}-map-{suffix}', '{SEMCONV_PHASE1_PREFIX}-map-child',
quantilesState(0.5, 0.75, 0.9, 0.95, 0.99)(toFloat64(1000000)),
toUInt64(0), toUInt64(1), toDateTime({int(now.timestamp())}),
'{environment}', '', ''
"""
)
yield now
cluster = clickhouse.env["SIGNOZ_TELEMETRYSTORE_CLICKHOUSE_CLUSTER"]
clickhouse.conn.command(f"ALTER TABLE signoz_traces.dependency_graph_minutes_v2 ON CLUSTER '{cluster}' DELETE WHERE startsWith(src, '{SEMCONV_PHASE1_PREFIX}-map-') SETTINGS mutations_sync = 1")

View File

@@ -0,0 +1,154 @@
"""Phase 2 semantic-convention checks across logs and metrics."""
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
import requests
from fixtures import querier, types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.logs import Logs
from fixtures.metrics import Metrics
DB_CURRENT = "db.system.name"
DB_OLD = "db.system"
METRIC_CURRENT = "container.cpu.usage"
METRIC_OLD = "container.cpu.utilization"
PREFIX = "semconv-phase2"
def test_logs_resolve_db_system_family(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_logs: Callable[[list[Logs]], None],
) -> None:
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=1)
rows = [
("old", {DB_OLD: "postgresql"}),
("current", {DB_CURRENT: "postgresql"}),
("conflict", {DB_OLD: "mysql", DB_CURRENT: "postgresql"}),
("missing", {}),
]
insert_logs(
[
Logs(
timestamp=now + timedelta(seconds=index),
resources={"service.name": PREFIX, **attributes},
attributes=attributes,
body=f"{PREFIX}-{suffix}",
)
for index, (suffix, attributes) in enumerate(rows)
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms = int((now - timedelta(minutes=2)).timestamp() * 1000)
end_ms = int((now + timedelta(minutes=1)).timestamp() * 1000)
present = {f"{PREFIX}-old", f"{PREFIX}-current", f"{PREFIX}-conflict"}
for context in ("attribute", "resource"):
for requested in (DB_CURRENT, DB_OLD):
field = f"{context}.{requested}"
query_cases = {
"equal": (f'{field} = "postgresql"', present),
"exists": (f"{field} EXISTS", present),
"not_exists": (f"{field} NOT EXISTS", {f"{PREFIX}-missing"}),
}
response = querier.make_query_request(
signoz,
token,
start_ms=start_ms,
end_ms=end_ms,
request_type=querier.RequestType.RAW,
queries=[querier.build_raw_query(name, "logs", limit=100, filter_expression=expression) for name, (expression, _) in query_cases.items()],
)
assert response.status_code == HTTPStatus.OK, response.text
results = response.json()["data"]["data"]["results"]
for name, (_, expected_bodies) in query_cases.items():
result = querier.find_named_result(results, name)
assert result is not None, name
assert {row["data"]["body"] for row in (result.get("rows") or [])} == expected_bodies
values_response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/fields/values"),
timeout=5,
headers={"authorization": f"Bearer {token}"},
params={
"signal": "logs",
"name": requested,
"fieldContext": context,
"fieldDataType": "string",
},
)
assert values_response.status_code == HTTPStatus.OK, values_response.text
assert set(values_response.json()["data"]["values"].get("stringValues") or []) == {"postgresql", "mysql"}
def test_metrics_resolve_label_and_metric_name_families(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_metrics: Callable[[list[Metrics]], None],
) -> None:
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
insert_metrics(
[
Metrics(
metric_name=METRIC_CURRENT,
labels={DB_CURRENT: "postgresql"},
timestamp=now - timedelta(seconds=3),
temporality="Unspecified",
type_="Gauge",
is_monotonic=False,
value=10,
),
Metrics(
metric_name=METRIC_OLD,
labels={"db_system": "mysql"},
timestamp=now - timedelta(seconds=2),
temporality="Unspecified",
type_="Gauge",
is_monotonic=False,
value=20,
),
Metrics(
metric_name=METRIC_OLD,
labels={DB_OLD: "mysql", DB_CURRENT: "postgresql", "series": "conflict"},
timestamp=now - timedelta(seconds=1),
temporality="Unspecified",
type_="Gauge",
is_monotonic=False,
value=30,
),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
for metric_name in (METRIC_CURRENT, METRIC_OLD):
for requested_label in (DB_CURRENT, DB_OLD):
response = querier.make_scalar_query_request(
signoz,
token,
now,
[
querier.build_scalar_query(
name="A",
signal="metrics",
aggregations=[
querier.build_metrics_aggregation(
metric_name,
"latest",
"sum",
"unspecified",
reduce_to="last",
)
],
group_by=[querier.build_group_by_field(requested_label, "string", "attribute")],
filter_expression=f"attribute.{requested_label} EXISTS",
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
data = {row[0]: row[-1] for row in querier.get_scalar_table_data(response.json())}
assert data == {"postgresql": 40.0, "mysql": 20.0}, (metric_name, requested_label, data)

View File

@@ -0,0 +1,204 @@
"""Phase 1 end-to-end checks for semantic-convention name evolution."""
from collections.abc import Callable
from datetime import datetime, timedelta
from http import HTTPStatus
import requests
from fixtures import querier, types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.semconv import SEMCONV_PHASE1_CURRENT as CURRENT
from fixtures.semconv import SEMCONV_PHASE1_OLD as OLD
from fixtures.semconv import SEMCONV_PHASE1_PREFIX as PREFIX
PRODUCTION_SPANS = {
f"{PREFIX}-old",
f"{PREFIX}-current",
f"{PREFIX}-both",
f"{PREFIX}-conflict",
}
STAGING_SPANS = {f"{PREFIX}-staging"}
MISSING_SPANS = {f"{PREFIX}-missing"}
def test_semconv_phase1_mixed_sdk_generations( # pylint: disable=too-many-statements
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
semconv_phase1_data: datetime,
) -> None:
now = semconv_phase1_data
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms = int((now - timedelta(minutes=2)).timestamp() * 1000)
end_ms = int((now + timedelta(minutes=1)).timestamp() * 1000)
# A builder response records the exact spelling that was resolved. Raw SQL
# and PromQL deliberately do not use this resolver.
resolution_response = querier.make_query_request(
signoz,
token,
start_ms=start_ms,
end_ms=end_ms,
request_type=querier.RequestType.RAW,
queries=[
querier.BuilderQuery(
signal="traces",
name="A",
limit=1,
filter_expression=f"resource.{OLD} EXISTS",
).to_dict()
],
)
assert resolution_response.status_code == HTTPStatus.OK, resolution_response.text
assert {
"requested": OLD,
"current": CURRENT,
"members": [CURRENT, OLD],
"kind": "attribute",
} in resolution_response.json()["data"]["meta"]["semconvResolutions"]
# Resource and span-attribute paths share the same matrix. Run every
# operator with both the saved-query (old) and current request spellings.
for context in ("resource", "attribute"):
for requested in (CURRENT, OLD):
field = f"{context}.{requested}"
query_cases = {
"production": (f"{field} = 'production'", PRODUCTION_SPANS),
"staging": (f"{field} = 'staging'", STAGING_SPANS),
# Negative operators intentionally include rows where no family
# member exists; explicit EXISTS is the opt-in presence filter.
"negative": (f"{field} != 'production'", STAGING_SPANS | MISSING_SPANS),
"exists": (f"{field} EXISTS", PRODUCTION_SPANS | STAGING_SPANS),
"not_exists": (f"{field} NOT EXISTS", MISSING_SPANS),
}
matrix_response = querier.make_query_request(
signoz,
token,
start_ms=start_ms,
end_ms=end_ms,
request_type=querier.RequestType.RAW,
queries=[
querier.BuilderQuery(
signal="traces",
name=name,
limit=100,
filter_expression=expression,
select_fields=[querier.TelemetryFieldKey("span.name")],
order=[querier.OrderBy(querier.TelemetryFieldKey("timestamp"), "asc")],
).to_dict()
for name, (expression, _) in query_cases.items()
],
)
assert matrix_response.status_code == HTTPStatus.OK, matrix_response.text
matrix_results = matrix_response.json()["data"]["data"]["results"]
for name, (_, expected_names) in query_cases.items():
result = querier.find_named_result(matrix_results, name)
assert result is not None, name
assert {row["data"]["name"] for row in (result.get("rows") or [])} == expected_names
grouped_response = querier.make_query_request(
signoz,
token,
start_ms=start_ms,
end_ms=end_ms,
request_type=querier.RequestType.SCALAR,
queries=[
querier.BuilderQuery(
signal="traces",
name="A",
filter_expression=f"{field} EXISTS",
aggregations=[querier.Aggregation("count()")],
group_by=[querier.TelemetryFieldKey(requested, "string", context)],
order=[querier.OrderBy(querier.TelemetryFieldKey(requested, "string", context), "asc")],
).to_dict()
],
)
assert grouped_response.status_code == HTTPStatus.OK, grouped_response.text
grouped_results = grouped_response.json()["data"]["data"]["results"]
assert len(grouped_results) == 1
assert grouped_results[0]["columns"][0]["name"] == requested, "response identity must match the request spelling"
assert grouped_results[0]["data"] == [["production", 4], ["staging", 1]]
for requested in (CURRENT, OLD):
values_response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/fields/values"),
timeout=5,
headers={"authorization": f"Bearer {token}"},
params={
"signal": "traces",
"name": requested,
"fieldContext": context,
"fieldDataType": "string",
},
)
assert values_response.status_code == HTTPStatus.OK, values_response.text
assert set(values_response.json()["data"]["values"].get("stringValues") or []) == {"production", "staging"}
keys_response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/fields/keys"),
timeout=5,
headers={"authorization": f"Bearer {token}"},
params={"signal": "traces", "searchText": OLD},
)
assert keys_response.status_code == HTTPStatus.OK, keys_response.text
keys = keys_response.json()["data"]["keys"]
assert CURRENT in keys
assert OLD in keys
start_ns = str(int((now - timedelta(minutes=2)).timestamp() * 1_000_000_000))
end_ns = str(int((now + timedelta(minutes=1)).timestamp() * 1_000_000_000))
for requested in (CURRENT, OLD):
services_response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/services"),
timeout=30,
headers={"authorization": f"Bearer {token}"},
json={
"start": start_ns,
"end": end_ns,
"tags": [
{
"Key": requested,
"Operator": "In",
"StringValues": ["production"],
"TagType": "ResourceAttribute",
}
],
},
)
assert services_response.status_code == HTTPStatus.OK, services_response.text
services = {item["serviceName"] for item in services_response.json()["data"]}
assert services == PRODUCTION_SPANS
map_response = requests.post(
signoz.self.host_configs["8080"].get("/api/v1/dependency_graph"),
timeout=30,
headers={"authorization": f"Bearer {token}"},
json={
"start": start_ns,
"end": end_ns,
"tags": [
{
"key": requested,
"operator": "In",
"stringValues": ["production"],
"tagType": "ResourceAttribute",
}
],
},
)
assert map_response.status_code == HTTPStatus.OK, map_response.text
assert {edge["parent"] for edge in map_response.json()} == {f"{PREFIX}-map-production"}
report_response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/fields/semconv-migration"),
timeout=30,
headers={"authorization": f"Bearer {token}"},
params={
"startUnixMilli": int((now - timedelta(minutes=2)).timestamp() * 1000),
"endUnixMilli": int((now + timedelta(minutes=1)).timestamp() * 1000),
},
)
assert report_response.status_code == HTTPStatus.OK, report_response.text
entry = next(item for item in report_response.json()["data"]["entries"] if item["current"] == CURRENT and item["old"] == OLD and item["signal"] == "traces")
assert set(entry["services"]) == {f"{PREFIX}-old", f"{PREFIX}-staging"}