Compare commits

..

3 Commits

Author SHA1 Message Date
srikanthccv
1e530643ff test: add semantic convention phase one closure gate 2026-08-07 21:20:50 +05:30
srikanthccv
08f4e0ea77 feat: support semconv evolution in services 2026-08-07 21:20:41 +05:30
srikanthccv
7e703723a4 feat: resolve semantic convention names in trace queries 2026-08-07 21:20:19 +05:30
101 changed files with 2995 additions and 968 deletions

View File

@@ -220,6 +220,10 @@ 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-clean
py-clean: ## Clear all pycache and pytest cache from tests directory recursively
@echo ">> cleaning python cache files from tests directory"

View File

@@ -98,6 +98,14 @@ func (ah *APIHandler) getFeatureFlags(w http.ResponseWriter, r *http.Request) {
Route: "",
})
if constants.IsDotMetricsEnabled {
for idx, feature := range featureSet {
if feature.Name == licensetypes.DotMetricsEnabled {
featureSet[idx].Active = true
}
}
}
ah.Respond(w, featureSet)
}

View File

@@ -17,3 +17,15 @@ func GetOrDefaultEnv(key string, fallback string) string {
}
return v
}
// constant functions that override env vars
const DotMetricsEnabled = "DOT_METRICS_ENABLED"
var IsDotMetricsEnabled = false
func init() {
if GetOrDefaultEnv(DotMetricsEnabled, "true") == "true" {
IsDotMetricsEnabled = true
}
}

View File

@@ -24,3 +24,19 @@ export const Logout = async (): Promise<void> => {
window.dispatchEvent(new CustomEvent('LOGOUT'));
history.push(ROUTES.LOGIN);
};
export const UnderscoreToDotMap: Record<string, string> = {
k8s_cluster_name: 'k8s.cluster.name',
k8s_cluster_uid: 'k8s.cluster.uid',
k8s_namespace_name: 'k8s.namespace.name',
k8s_node_name: 'k8s.node.name',
k8s_node_uid: 'k8s.node.uid',
k8s_pod_name: 'k8s.pod.name',
k8s_pod_uid: 'k8s.pod.uid',
k8s_deployment_name: 'k8s.deployment.name',
k8s_daemonset_name: 'k8s.daemonset.name',
k8s_statefulset_name: 'k8s.statefulset.name',
k8s_cronjob_name: 'k8s.cronjob.name',
k8s_job_name: 'k8s.job.name',
k8s_persistentvolumeclaim_name: 'k8s.persistentvolumeclaim.name',
};

View File

@@ -7,6 +7,7 @@ export enum FeatureKeys {
GATEWAY = 'gateway',
PREMIUM_SUPPORT = 'premium_support',
ANOMALY_DETECTION = 'anomaly_detection',
DOT_METRICS_ENABLED = 'dot_metrics_enabled',
USE_JSON_BODY = 'use_json_body',
ENABLE_AI_OBSERVABILITY = 'enable_ai_observability',
ENABLE_METRICS_REDUCTION = 'enable_metrics_reduction',

View File

@@ -37,6 +37,8 @@ import { ErrorResponse, SuccessResponse } from 'types/api';
import { Exception, PayloadProps } from 'types/api/errors/getAll';
import { GlobalReducer } from 'types/reducer/globalTime';
import { FeatureKeys } from '../../constants/features';
import { useAppContext } from '../../providers/App/App';
import { FilterDropdownExtendsProps } from './types';
import {
extractFilterValues,
@@ -416,6 +418,11 @@ function AllErrors(): JSX.Element {
},
];
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const onChangeHandler: TableProps<Exception>['onChange'] = useCallback(
(
paginations: TablePaginationConfig,
@@ -451,7 +458,7 @@ function AllErrors(): JSX.Element {
useEffect(() => {
if (!isUndefined(errorCountResponse.data?.payload)) {
const selectedEnvironments = queries.find(
(val) => val.tagKey === getResourceDeploymentKeys(),
(val) => val.tagKey === getResourceDeploymentKeys(dotMetricsEnabled),
)?.tagValue;
logEvent('Exception: List page visited', {

View File

@@ -35,6 +35,7 @@ import { openInNewTab } from 'utils/navigation';
import triangleRulerUrl from '@/assets/Icons/triangle-ruler.svg';
import { FeatureKeys } from '../../../constants/features';
import { DOCS_LINKS } from '../constants';
import { columns, TIME_PICKER_OPTIONS } from './constants';
@@ -211,13 +212,19 @@ function ServiceMetrics({
const topLevelOperations = useMemo(() => Object.entries(data || {}), [data]);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const queryRangeRequestData = useMemo(
() =>
getQueryRangeRequestData({
topLevelOperations,
globalSelectedInterval,
dotMetricsEnabled,
}),
[globalSelectedInterval, topLevelOperations],
[globalSelectedInterval, topLevelOperations, dotMetricsEnabled],
);
const dataQueries = useGetQueriesRange(

View File

@@ -82,7 +82,7 @@ export function getHostMetricsQueryPayload(
start: number,
end: number,
): ReturnType<typeof getHostQueryPayload> {
return getHostQueryPayload(host.hostName, start, end);
return getHostQueryPayload(host.hostName, start, end, true);
}
export { hostWidgetInfo };

View File

@@ -121,6 +121,12 @@ jest.spyOn(appContextHooks, 'useAppContext').mockReturnValue({
plan_version: 'test-plan-version',
},
},
featureFlags: [
{
name: 'DOT_METRICS_ENABLED',
active: false,
},
],
} as any);
const mockEntity = {

View File

@@ -17,6 +17,8 @@ import { SuccessResponse } from 'types/api';
import { MetricRangePayloadProps } from 'types/api/metrics/getQueryRange';
import uPlot from 'uplot';
import { FeatureKeys } from '../../../constants/features';
import { useAppContext } from '../../../providers/App/App';
import {
getHostQueryPayload,
getNodeQueryPayload,
@@ -51,12 +53,23 @@ function NodeMetrics({
};
}, [timestamp]);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const queryPayloads = useMemo(() => {
if (nodeName) {
return getNodeQueryPayload(clusterName, nodeName, start, end);
return getNodeQueryPayload(
clusterName,
nodeName,
start,
end,
dotMetricsEnabled,
);
}
return getHostQueryPayload(hostName, start, end);
}, [nodeName, hostName, clusterName, start, end]);
return getHostQueryPayload(hostName, start, end, dotMetricsEnabled);
}, [nodeName, hostName, clusterName, start, end, dotMetricsEnabled]);
const widgetInfo = nodeName ? nodeWidgetInfo : hostWidgetInfo;
const queries = useQueries(

View File

@@ -12,11 +12,13 @@ import { useResizeObserver } from 'hooks/useDimensions';
import { GetMetricQueryRange } from 'lib/dashboard/getQueryResults';
import { getUPlotChartOptions } from 'lib/uPlotLib/getUplotChartOptions';
import { getUPlotChartData } from 'lib/uPlotLib/utils/getUplotChartData';
import { useAppContext } from 'providers/App/App';
import { useTimezone } from 'providers/Timezone';
import { SuccessResponse } from 'types/api';
import { MetricRangePayloadProps } from 'types/api/metrics/getQueryRange';
import uPlot from 'uplot';
import { FeatureKeys } from '../../../constants/features';
import { getPodQueryPayload, podWidgetInfo } from './constants';
function PodMetrics({
@@ -52,9 +54,14 @@ function PodMetrics({
scrollLeft: 0,
});
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const queryPayloads = useMemo(
() => getPodQueryPayload(clusterName, podName, start, end),
[clusterName, end, podName, start],
() => getPodQueryPayload(clusterName, podName, start, end, dotMetricsEnabled),
[clusterName, end, podName, start, dotMetricsEnabled],
);
const queries = useQueries(
queryPayloads.map((payload) => ({

View File

@@ -9,21 +9,48 @@ export const getPodQueryPayload = (
podName: string,
start: number,
end: number,
dotMetricsEnabled: boolean,
): GetQueryResultsProps[] => {
const k8sClusterNameKey = 'k8s.cluster.name';
const k8sPodNameKey = 'k8s.pod.name';
const containerCpuUtilKey = 'container.cpu.usage';
const containerMemUsageKey = 'container.memory.usage';
const k8sContainerCpuReqKey = 'k8s.container.cpu_request';
const k8sContainerCpuLimitKey = 'k8s.container.cpu_limit';
const k8sContainerMemReqKey = 'k8s.container.memory_request';
const k8sContainerMemLimitKey = 'k8s.container.memory_limit';
const k8sPodFsAvailKey = 'k8s.pod.filesystem.available';
const k8sPodFsCapKey = 'k8s.pod.filesystem.capacity';
const k8sPodNetIoKey = 'k8s.pod.network.io';
const podLegendTemplate = '{{k8s.pod.name}}';
const podLegendUsage = 'usage - {{k8s.pod.name}}';
const podLegendLimit = 'limit - {{k8s.pod.name}}';
const k8sClusterNameKey = dotMetricsEnabled
? 'k8s.cluster.name'
: 'k8s_cluster_name';
const k8sPodNameKey = dotMetricsEnabled ? 'k8s.pod.name' : 'k8s_pod_name';
const containerCpuUtilKey = dotMetricsEnabled
? 'container.cpu.usage'
: 'container_cpu_usage';
const containerMemUsageKey = dotMetricsEnabled
? 'container.memory.usage'
: 'container_memory_usage';
const k8sContainerCpuReqKey = dotMetricsEnabled
? 'k8s.container.cpu_request'
: 'k8s_container_cpu_request';
const k8sContainerCpuLimitKey = dotMetricsEnabled
? 'k8s.container.cpu_limit'
: 'k8s_container_cpu_limit';
const k8sContainerMemReqKey = dotMetricsEnabled
? 'k8s.container.memory_request'
: 'k8s_container_memory_request';
const k8sContainerMemLimitKey = dotMetricsEnabled
? 'k8s.container.memory_limit'
: 'k8s_container_memory_limit';
const k8sPodFsAvailKey = dotMetricsEnabled
? 'k8s.pod.filesystem.available'
: 'k8s_pod_filesystem_available';
const k8sPodFsCapKey = dotMetricsEnabled
? 'k8s.pod.filesystem.capacity'
: 'k8s_pod_filesystem_capacity';
const k8sPodNetIoKey = dotMetricsEnabled
? 'k8s.pod.network.io'
: 'k8s_pod_network_io';
const podLegendTemplate = dotMetricsEnabled
? '{{k8s.pod.name}}'
: '{{k8s_pod_name}}';
const podLegendUsage = dotMetricsEnabled
? 'usage - {{k8s.pod.name}}'
: 'usage - {{k8s_pod_name}}';
const podLegendLimit = dotMetricsEnabled
? 'limit - {{k8s.pod.name}}'
: 'limit - {{k8s_pod_name}}';
return [
{
@@ -1000,17 +1027,36 @@ export const getNodeQueryPayload = (
nodeName: string,
start: number,
end: number,
dotMetricsEnabled: boolean,
): GetQueryResultsProps[] => {
const k8sClusterNameKey = 'k8s.cluster.name';
const k8sNodeNameKey = 'k8s.node.name';
const k8sNodeCpuTimeKey = 'k8s.node.cpu.time';
const k8sNodeAllocCpuKey = 'k8s.node.allocatable_cpu';
const k8sNodeMemWsKey = 'k8s.node.memory.working_set';
const k8sNodeAllocMemKey = 'k8s.node.allocatable_memory';
const k8sNodeNetIoKey = 'k8s.node.network.io';
const k8sNodeFsAvailKey = 'k8s.node.filesystem.available';
const k8sNodeFsCapKey = 'k8s.node.filesystem.capacity';
const podLegend = '{{k8s.node.name}}';
const k8sClusterNameKey = dotMetricsEnabled
? 'k8s.cluster.name'
: 'k8s_cluster_name';
const k8sNodeNameKey = dotMetricsEnabled ? 'k8s.node.name' : 'k8s_node_name';
const k8sNodeCpuTimeKey = dotMetricsEnabled
? 'k8s.node.cpu.time'
: 'k8s_node_cpu_time';
const k8sNodeAllocCpuKey = dotMetricsEnabled
? 'k8s.node.allocatable_cpu'
: 'k8s_node_allocatable_cpu';
const k8sNodeMemWsKey = dotMetricsEnabled
? 'k8s.node.memory.working_set'
: 'k8s_node_memory_working_set';
const k8sNodeAllocMemKey = dotMetricsEnabled
? 'k8s.node.allocatable_memory'
: 'k8s_node_allocatable_memory';
const k8sNodeNetIoKey = dotMetricsEnabled
? 'k8s.node.network.io'
: 'k8s_node_network_io';
const k8sNodeFsAvailKey = dotMetricsEnabled
? 'k8s.node.filesystem.available'
: 'k8s_node_filesystem_available';
const k8sNodeFsCapKey = dotMetricsEnabled
? 'k8s.node.filesystem.capacity'
: 'k8s_node_filesystem_capacity';
const podLegend = dotMetricsEnabled
? '{{k8s.node.name}}'
: '{{k8s_node_name}}';
return [
{
@@ -1540,23 +1586,48 @@ export const getHostQueryPayload = (
hostName: string,
start: number,
end: number,
dotMetricsEnabled: boolean,
): GetQueryResultsProps[] => {
const hostNameKey = 'host.name';
const cpuTimeKey = 'system.cpu.time';
const memUsageKey = 'system.memory.usage';
const load1mKey = 'system.cpu.load_average.1m';
const load5mKey = 'system.cpu.load_average.5m';
const load15mKey = 'system.cpu.load_average.15m';
const netIoKey = 'system.network.io';
const netPktsKey = 'system.network.packets';
const netErrKey = 'system.network.errors';
const netDropKey = 'system.network.dropped';
const netConnKey = 'system.network.connections';
const diskIoKey = 'system.disk.io';
const diskOpTimeKey = 'system.disk.operation_time';
const diskOpsKey = 'system.disk.operations';
const diskPendingKey = 'system.disk.pending_operations';
const fsUsageKey = 'system.filesystem.usage';
const hostNameKey = dotMetricsEnabled ? 'host.name' : 'host_name';
const cpuTimeKey = dotMetricsEnabled ? 'system.cpu.time' : 'system_cpu_time';
const memUsageKey = dotMetricsEnabled
? 'system.memory.usage'
: 'system_memory_usage';
const load1mKey = dotMetricsEnabled
? 'system.cpu.load_average.1m'
: 'system_cpu_load_average_1m';
const load5mKey = dotMetricsEnabled
? 'system.cpu.load_average.5m'
: 'system_cpu_load_average_5m';
const load15mKey = dotMetricsEnabled
? 'system.cpu.load_average.15m'
: 'system_cpu_load_average_15m';
const netIoKey = dotMetricsEnabled ? 'system.network.io' : 'system_network_io';
const netPktsKey = dotMetricsEnabled
? 'system.network.packets'
: 'system_network_packets';
const netErrKey = dotMetricsEnabled
? 'system.network.errors'
: 'system_network_errors';
const netDropKey = dotMetricsEnabled
? 'system.network.dropped'
: 'system_network_dropped';
const netConnKey = dotMetricsEnabled
? 'system.network.connections'
: 'system_network_connections';
const diskIoKey = dotMetricsEnabled ? 'system.disk.io' : 'system_disk_io';
const diskOpTimeKey = dotMetricsEnabled
? 'system.disk.operation_time'
: 'system_disk_operation_time';
const diskOpsKey = dotMetricsEnabled
? 'system.disk.operations'
: 'system_disk_operations';
const diskPendingKey = dotMetricsEnabled
? 'system.disk.pending_operations'
: 'system_disk_pending_operations';
const fsUsageKey = dotMetricsEnabled
? 'system.filesystem.usage'
: 'system_filesystem_usage';
return [
{

View File

@@ -21,6 +21,7 @@ export const databaseCallsRPS = ({
servicename,
legend,
tagFilterItems,
dotMetricsEnabled,
}: DatabaseCallsRPSProps): QueryBuilderData => {
const autocompleteData: BaseAutocompleteData[] = [
{
@@ -32,7 +33,7 @@ export const databaseCallsRPS = ({
const groupBy: BaseAutocompleteData[] = [
{
dataType: DataTypes.String,
key: WidgetKeys.DbSystem,
key: dotMetricsEnabled ? WidgetKeys.Db_system : WidgetKeys.Db_system_norm,
type: 'tag',
},
];
@@ -41,7 +42,9 @@ export const databaseCallsRPS = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},
@@ -72,6 +75,7 @@ export const databaseCallsRPS = ({
export const databaseCallsAvgDuration = ({
servicename,
tagFilterItems,
dotMetricsEnabled,
}: DatabaseCallProps): QueryBuilderData => {
const autocompleteDataA: BaseAutocompleteData = {
key: WidgetKeys.SignozDbLatencySum,
@@ -88,7 +92,9 @@ export const databaseCallsAvgDuration = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},

View File

@@ -32,6 +32,7 @@ export const externalCallErrorPercent = ({
servicename,
legend,
tagFilterItems,
dotMetricsEnabled,
}: ExternalCallDurationByAddressProps): QueryBuilderData => {
const autocompleteDataA: BaseAutocompleteData = {
key: WidgetKeys.SignozExternalCallLatencyCount,
@@ -48,7 +49,9 @@ export const externalCallErrorPercent = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},
@@ -58,7 +61,7 @@ export const externalCallErrorPercent = ({
{
id: '',
key: {
key: WidgetKeys.StatusCode,
key: dotMetricsEnabled ? WidgetKeys.StatusCode : WidgetKeys.StatusCodeNorm,
dataType: DataTypes.Int64,
type: MetricsType.Tag,
},
@@ -71,7 +74,9 @@ export const externalCallErrorPercent = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},
@@ -115,6 +120,7 @@ export const externalCallErrorPercent = ({
export const externalCallDuration = ({
servicename,
tagFilterItems,
dotMetricsEnabled,
}: ExternalCallProps): QueryBuilderData => {
const autocompleteDataA: BaseAutocompleteData = {
dataType: DataTypes.Float64,
@@ -135,7 +141,9 @@ export const externalCallDuration = ({
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -175,6 +183,7 @@ export const externalCallRpsByAddress = ({
servicename,
legend,
tagFilterItems,
dotMetricsEnabled,
}: ExternalCallDurationByAddressProps): QueryBuilderData => {
const autocompleteData: BaseAutocompleteData[] = [
{
@@ -189,7 +198,9 @@ export const externalCallRpsByAddress = ({
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -220,6 +231,7 @@ export const externalCallDurationByAddress = ({
servicename,
legend,
tagFilterItems,
dotMetricsEnabled,
}: ExternalCallDurationByAddressProps): QueryBuilderData => {
const autocompleteDataA: BaseAutocompleteData = {
dataType: DataTypes.Float64,
@@ -239,7 +251,9 @@ export const externalCallDurationByAddress = ({
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,

View File

@@ -37,10 +37,15 @@ export const latency = ({
tagFilterItems,
isSpanMetricEnable = false,
topLevelOperationsRoute,
dotMetricsEnabled,
}: LatencyProps): QueryBuilderData => {
const signozLatencyBucketMetrics = WidgetKeys.SignozLatencyBucket;
const signozLatencyBucketMetrics = dotMetricsEnabled
? WidgetKeys.Signoz_latency_bucket
: WidgetKeys.Signoz_latency_bucket_norm;
const signozMetricsServiceName = WidgetKeys.OTelServiceName;
const signozMetricsServiceName = dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm;
const newAutoCompleteData: BaseAutocompleteData = {
key: isSpanMetricEnable
? signozLatencyBucketMetrics
@@ -282,21 +287,28 @@ export const apDexMetricsQueryBuilderQueries = ({
threashold,
delta,
metricsBuckets,
dotMetricsEnabled,
}: ApDexMetricsQueryBuilderQueriesProps): QueryBuilderData => {
const autoCompleteDataA: BaseAutocompleteData = {
key: WidgetKeys.SignozLatencyCount,
key: dotMetricsEnabled
? WidgetKeys.SignozLatencyCount
: WidgetKeys.SignozLatencyCountNorm,
dataType: DataTypes.Float64,
type: '',
};
const autoCompleteDataB: BaseAutocompleteData = {
key: WidgetKeys.SignozLatencyBucket,
key: dotMetricsEnabled
? WidgetKeys.Signoz_latency_bucket
: WidgetKeys.Signoz_latency_bucket_norm,
dataType: DataTypes.Float64,
type: '',
};
const autoCompleteDataC: BaseAutocompleteData = {
key: WidgetKeys.SignozLatencyBucket,
key: dotMetricsEnabled
? WidgetKeys.Signoz_latency_bucket
: WidgetKeys.Signoz_latency_bucket_norm,
dataType: DataTypes.Float64,
type: '',
};
@@ -305,7 +317,9 @@ export const apDexMetricsQueryBuilderQueries = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Tag,
},
@@ -329,7 +343,7 @@ export const apDexMetricsQueryBuilderQueries = ({
{
id: '',
key: {
key: WidgetKeys.StatusCode,
key: dotMetricsEnabled ? WidgetKeys.StatusCode : WidgetKeys.StatusCodeNorm,
dataType: DataTypes.String,
type: MetricsType.Tag,
},
@@ -349,7 +363,9 @@ export const apDexMetricsQueryBuilderQueries = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Tag,
},
@@ -383,7 +399,7 @@ export const apDexMetricsQueryBuilderQueries = ({
{
id: '',
key: {
key: WidgetKeys.StatusCode,
key: dotMetricsEnabled ? WidgetKeys.StatusCode : WidgetKeys.StatusCodeNorm,
dataType: DataTypes.String,
type: MetricsType.Tag,
},
@@ -393,7 +409,9 @@ export const apDexMetricsQueryBuilderQueries = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Tag,
},
@@ -456,10 +474,13 @@ export const operationPerSec = ({
servicename,
tagFilterItems,
topLevelOperations,
dotMetricsEnabled,
}: OperationPerSecProps): QueryBuilderData => {
const autocompleteData: BaseAutocompleteData[] = [
{
key: WidgetKeys.SignozLatencyCount,
key: dotMetricsEnabled
? WidgetKeys.SignozLatencyCount
: WidgetKeys.SignozLatencyCountNorm,
dataType: DataTypes.Float64,
type: '',
},
@@ -470,7 +491,9 @@ export const operationPerSec = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},
@@ -511,6 +534,7 @@ export const errorPercentage = ({
servicename,
tagFilterItems,
topLevelOperations,
dotMetricsEnabled,
}: OperationPerSecProps): QueryBuilderData => {
const autocompleteDataA: BaseAutocompleteData = {
key: WidgetKeys.SignozCallsTotal,
@@ -529,7 +553,9 @@ export const errorPercentage = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},
@@ -549,7 +575,7 @@ export const errorPercentage = ({
{
id: '',
key: {
key: WidgetKeys.StatusCode,
key: dotMetricsEnabled ? WidgetKeys.StatusCode : WidgetKeys.StatusCodeNorm,
dataType: DataTypes.Int64,
type: MetricsType.Tag,
},
@@ -563,7 +589,9 @@ export const errorPercentage = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},

View File

@@ -21,9 +21,12 @@ import { getQueryBuilderQuerieswithFormula } from './MetricsPageQueriesFactory';
export const topOperationQueries = ({
servicename,
dotMetricsEnabled,
}: TopOperationQueryFactoryProps): QueryBuilderData => {
const latencyAutoCompleteData: BaseAutocompleteData = {
key: WidgetKeys.SignozLatencyBucket,
key: dotMetricsEnabled
? WidgetKeys.Signoz_latency_bucket
: WidgetKeys.Signoz_latency_bucket_norm,
dataType: DataTypes.Float64,
type: '',
};
@@ -35,7 +38,9 @@ export const topOperationQueries = ({
};
const numOfCallAutoCompleteData: BaseAutocompleteData = {
key: WidgetKeys.SignozLatencyCount,
key: dotMetricsEnabled
? WidgetKeys.SignozLatencyCount
: WidgetKeys.SignozLatencyCountNorm,
dataType: DataTypes.Float64,
type: '',
};
@@ -44,7 +49,9 @@ export const topOperationQueries = ({
{
id: '',
key: {
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
dataType: DataTypes.String,
type: MetricsType.Resource,
},
@@ -58,7 +65,9 @@ export const topOperationQueries = ({
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -68,7 +77,7 @@ export const topOperationQueries = ({
id: '',
key: {
dataType: DataTypes.Int64,
key: WidgetKeys.StatusCode,
key: dotMetricsEnabled ? WidgetKeys.StatusCode : WidgetKeys.StatusCodeNorm,
type: MetricsType.Tag,
},
op: OPERATORS.IN,

View File

@@ -28,6 +28,8 @@ import { TagFilterItem } from 'types/api/queryBuilder/queryBuilderData';
import { EQueryType } from 'types/common/dashboard';
import { v4 as uuid } from 'uuid';
import { FeatureKeys } from '../../../constants/features';
import { useAppContext } from '../../../providers/App/App';
import {
GraphTitle,
MENU_ITEMS,
@@ -87,7 +89,12 @@ function DBCall(): JSX.Element {
[queries],
);
const legend = '{{db.system}}';
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const legend = dotMetricsEnabled ? '{{db.system}}' : '{{db_system}}';
const databaseCallsRPSWidget = useMemo(
() =>
@@ -99,6 +106,7 @@ function DBCall(): JSX.Element {
servicename,
legend,
tagFilterItems,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -109,7 +117,7 @@ function DBCall(): JSX.Element {
id: SERVICE_CHART_ID.dbCallsRPS,
fillSpans: false,
}),
[servicename, tagFilterItems, legend],
[servicename, tagFilterItems, dotMetricsEnabled, legend],
);
const databaseCallsAverageDurationWidget = useMemo(
() =>
@@ -120,6 +128,7 @@ function DBCall(): JSX.Element {
builder: databaseCallsAvgDuration({
servicename,
tagFilterItems,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -130,7 +139,7 @@ function DBCall(): JSX.Element {
id: GraphTitle.DATABASE_CALLS_AVG_DURATION,
fillSpans: true,
}),
[servicename, tagFilterItems],
[servicename, tagFilterItems, dotMetricsEnabled],
);
const stepInterval = useMemo(
@@ -148,7 +157,7 @@ function DBCall(): JSX.Element {
useEffect(() => {
if (!logEventCalledRef.current) {
const selectedEnvironments = queries.find(
(val) => val.tagKey === getResourceDeploymentKeys(),
(val) => val.tagKey === getResourceDeploymentKeys(dotMetricsEnabled),
)?.tagValue;
logEvent('APM: Service detail page visited', {

View File

@@ -30,6 +30,8 @@ import { DataTypes } from 'types/api/queryBuilder/queryAutocompleteResponse';
import { EQueryType } from 'types/common/dashboard';
import { v4 as uuid } from 'uuid';
import { FeatureKeys } from '../../../constants/features';
import { useAppContext } from '../../../providers/App/App';
import {
GraphTitle,
legend,
@@ -82,6 +84,10 @@ function External(): JSX.Element {
handleNonInQueryRange(resourceAttributesToTagFilterItems(queries)) || [],
[queries],
);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const externalCallErrorWidget = useMemo(
() =>
@@ -93,6 +99,7 @@ function External(): JSX.Element {
servicename,
legend: legend.address,
tagFilterItems,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -102,7 +109,7 @@ function External(): JSX.Element {
yAxisUnit: '%',
id: GraphTitle.EXTERNAL_CALL_ERROR_PERCENTAGE,
}),
[servicename, tagFilterItems],
[servicename, tagFilterItems, dotMetricsEnabled],
);
const selectedTraceTags = useMemo(
@@ -119,6 +126,7 @@ function External(): JSX.Element {
builder: externalCallDuration({
servicename,
tagFilterItems,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -129,7 +137,7 @@ function External(): JSX.Element {
id: GraphTitle.EXTERNAL_CALL_DURATION,
fillSpans: true,
}),
[servicename, tagFilterItems],
[servicename, tagFilterItems, dotMetricsEnabled],
);
const errorApmToTraceQuery = useGetAPMToTracesQueries({
@@ -163,7 +171,7 @@ function External(): JSX.Element {
useEffect(() => {
if (!logEventCalledRef.current) {
const selectedEnvironments = queries.find(
(val) => val.tagKey === getResourceDeploymentKeys(),
(val) => val.tagKey === getResourceDeploymentKeys(dotMetricsEnabled),
)?.tagValue;
logEvent('APM: Service detail page visited', {
@@ -186,6 +194,7 @@ function External(): JSX.Element {
servicename,
legend: legend.address,
tagFilterItems,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -196,7 +205,7 @@ function External(): JSX.Element {
id: GraphTitle.EXTERNAL_CALL_RPS_BY_ADDRESS,
fillSpans: true,
}),
[servicename, tagFilterItems],
[servicename, tagFilterItems, dotMetricsEnabled],
);
const externalCallDurationAddressWidget = useMemo(
@@ -209,6 +218,7 @@ function External(): JSX.Element {
servicename,
legend: legend.address,
tagFilterItems,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -219,7 +229,7 @@ function External(): JSX.Element {
id: GraphTitle.EXTERNAL_CALL_DURATION_BY_ADDRESS,
fillSpans: true,
}),
[servicename, tagFilterItems],
[servicename, tagFilterItems, dotMetricsEnabled],
);
const apmToTraceQuery = useGetAPMToTracesQueries({

View File

@@ -93,12 +93,15 @@ function Application(): JSX.Element {
// eslint-disable-next-line react-hooks/exhaustive-deps
[handleSetTimeStamp],
);
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const logEventCalledRef = useRef(false);
useEffect(() => {
if (!logEventCalledRef.current) {
const selectedEnvironments = queries.find(
(val) => val.tagKey === getResourceDeploymentKeys(),
(val) => val.tagKey === getResourceDeploymentKeys(dotMetricsEnabled),
)?.tagValue;
logEvent('APM: Service detail page visited', {
@@ -156,6 +159,7 @@ function Application(): JSX.Element {
servicename,
tagFilterItems,
topLevelOperations: topLevelOperationsRoute,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -165,7 +169,7 @@ function Application(): JSX.Element {
yAxisUnit: 'ops',
id: SERVICE_CHART_ID.rps,
}),
[servicename, tagFilterItems, topLevelOperationsRoute],
[servicename, tagFilterItems, topLevelOperationsRoute, dotMetricsEnabled],
);
const errorPercentageWidget = useMemo(
@@ -178,6 +182,7 @@ function Application(): JSX.Element {
servicename,
tagFilterItems,
topLevelOperations: topLevelOperationsRoute,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -188,7 +193,7 @@ function Application(): JSX.Element {
id: SERVICE_CHART_ID.errorPercentage,
fillSpans: true,
}),
[servicename, tagFilterItems, topLevelOperationsRoute],
[servicename, tagFilterItems, topLevelOperationsRoute, dotMetricsEnabled],
);
const stepInterval = useMemo(

View File

@@ -22,6 +22,8 @@ import { apDexMetricsQueryBuilderQueries } from 'container/MetricsApplication/Me
import { EQueryType } from 'types/common/dashboard';
import { v4 as uuid } from 'uuid';
import { FeatureKeys } from '../../../../../constants/features';
import { useAppContext } from '../../../../../providers/App/App';
import { IServiceName } from '../../types';
import { ApDexMetricsProps } from './types';
@@ -36,6 +38,10 @@ function ApDexMetrics({
}: ApDexMetricsProps): JSX.Element {
const { servicename: encodedServiceName } = useParams<IServiceName>();
const servicename = decodeURIComponent(encodedServiceName);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const apDexMetricsWidget = useMemo(
() =>
getWidgetQueryBuilder({
@@ -49,6 +55,7 @@ function ApDexMetrics({
threashold: thresholdValue || 0,
delta: delta || false,
metricsBuckets: metricsBuckets || [],
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -74,6 +81,7 @@ function ApDexMetrics({
tagFilterItems,
thresholdValue,
topLevelOperationsRoute,
dotMetricsEnabled,
],
);

View File

@@ -3,6 +3,8 @@ import Spinner from 'components/Spinner';
import { useGetMetricMeta } from 'hooks/apDex/useGetMetricMeta';
import useErrorNotification from 'hooks/useErrorNotification';
import { FeatureKeys } from '../../../../../constants/features';
import { useAppContext } from '../../../../../providers/App/App';
import { WidgetKeys } from '../../../constant';
import { IServiceName } from '../../types';
import ApDexMetrics from './ApDexMetrics';
@@ -18,8 +20,17 @@ function ApDexMetricsApplication({
const { servicename: encodedServiceName } = useParams<IServiceName>();
const servicename = decodeURIComponent(encodedServiceName);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const signozLatencyBucketMetrics = dotMetricsEnabled
? WidgetKeys.Signoz_latency_bucket
: WidgetKeys.Signoz_latency_bucket_norm;
const { data, isLoading, error } = useGetMetricMeta(
WidgetKeys.SignozLatencyBucket,
signozLatencyBucketMetrics,
servicename,
);
useErrorNotification(error);

View File

@@ -56,6 +56,10 @@ function ServiceOverview({
[isSpanMetricEnable, queries],
);
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const latencyWidget = useMemo(
() =>
getWidgetQueryBuilder({
@@ -67,6 +71,7 @@ function ServiceOverview({
tagFilterItems,
isSpanMetricEnable,
topLevelOperationsRoute,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
@@ -76,7 +81,13 @@ function ServiceOverview({
yAxisUnit: 'ns',
id: SERVICE_CHART_ID.latency,
}),
[isSpanMetricEnable, servicename, tagFilterItems, topLevelOperationsRoute],
[
isSpanMetricEnable,
servicename,
tagFilterItems,
topLevelOperationsRoute,
dotMetricsEnabled,
],
);
const isQueryEnabled =

View File

@@ -19,6 +19,8 @@ import { EQueryType } from 'types/common/dashboard';
import { GlobalReducer } from 'types/reducer/globalTime';
import { v4 as uuid } from 'uuid';
import { FeatureKeys } from '../../../../constants/features';
import { useAppContext } from '../../../../providers/App/App';
import { IServiceName } from '../types';
import { title } from './config';
import ColumnWithLink from './TableRenderer/ColumnWithLink';
@@ -42,6 +44,11 @@ function TopOperationMetrics(): JSX.Element {
convertRawQueriesToTraceSelectedTags(queries) || [],
);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const keyOperationWidget = useMemo(
() =>
getWidgetQueryBuilder({
@@ -50,13 +57,14 @@ function TopOperationMetrics(): JSX.Element {
promql: [],
builder: topOperationQueries({
servicename,
dotMetricsEnabled,
}),
clickhouse_sql: [],
id: uuid(),
},
panelTypes: PANEL_TYPES.TABLE,
}),
[servicename],
[servicename, dotMetricsEnabled],
);
const updatedQuery = updateStepInterval(keyOperationWidget.query);

View File

@@ -10,6 +10,7 @@ export interface IServiceName {
export interface TopOperationQueryFactoryProps {
servicename: IServiceName['servicename'];
dotMetricsEnabled: boolean;
}
export interface ExternalCallDurationByAddressProps extends ExternalCallProps {
@@ -19,6 +20,7 @@ export interface ExternalCallDurationByAddressProps extends ExternalCallProps {
export interface ExternalCallProps {
servicename: IServiceName['servicename'];
tagFilterItems: TagFilterItem[];
dotMetricsEnabled: boolean;
}
export interface BuilderQueriesProps {
@@ -50,6 +52,7 @@ export interface OperationPerSecProps {
servicename: IServiceName['servicename'];
tagFilterItems: TagFilterItem[];
topLevelOperations: string[];
dotMetricsEnabled: boolean;
}
export interface LatencyProps {
@@ -57,6 +60,7 @@ export interface LatencyProps {
tagFilterItems: TagFilterItem[];
isSpanMetricEnable?: boolean;
topLevelOperationsRoute: string[];
dotMetricsEnabled: boolean;
}
export interface ApDexProps {
@@ -74,4 +78,5 @@ export interface TableRendererProps {
export interface ApDexMetricsQueryBuilderQueriesProps extends ApDexProps {
delta: boolean;
metricsBuckets: number[];
dotMetricsEnabled: boolean;
}

View File

@@ -85,11 +85,14 @@ export enum WidgetKeys {
HasError = 'hasError',
Address = 'address',
DurationNano = 'durationNano',
StatusCodeNorm = 'status_code',
StatusCode = 'status.code',
Operation = 'operation',
OperationName = 'operationName',
OTelServiceName = 'service.name',
Service_name_norm = 'service_name',
Service_name = 'service.name',
ServiceName = 'serviceName',
SignozLatencyCountNorm = 'signoz_latency_count',
SignozLatencyCount = 'signoz_latency.count',
SignozDBLatencyCount = 'signoz_db_latency_count',
DatabaseCallCount = 'signoz_database_call_count',
@@ -98,8 +101,10 @@ export enum WidgetKeys {
SignozCallsTotal = 'signoz_calls_total',
SignozExternalCallLatencyCount = 'signoz_external_call_latency_count',
SignozExternalCallLatencySum = 'signoz_external_call_latency_sum',
SignozLatencyBucket = 'signoz_latency.bucket',
DbSystem = 'db.system',
Signoz_latency_bucket_norm = 'signoz_latency_bucket',
Signoz_latency_bucket = 'signoz_latency.bucket',
Db_system = 'db.system',
Db_system_norm = 'db_system',
}
export const topOperationMetricsDownloadOptions: DownloadOptions = {

View File

@@ -32,4 +32,5 @@ export interface DatabaseCallsRPSProps extends DatabaseCallProps {
export interface DatabaseCallProps {
servicename: IServiceName['servicename'];
tagFilterItems: TagFilterItem[];
dotMetricsEnabled: boolean;
}

View File

@@ -53,6 +53,8 @@ import { getUserOperatingSystem, UserOperatingSystem } from 'utils/getUserOS';
import { useSelectPopupContainer } from 'utils/selectPopupContainer';
import { v4 as uuid } from 'uuid';
import { FeatureKeys } from '../../../../constants/features';
import { useAppContext } from '../../../../providers/App/App';
import { selectStyle } from './config';
import { PLACEHOLDER } from './constant';
import ExampleQueriesRendererForLogs from './ExampleQueriesRendererForLogs';
@@ -102,6 +104,11 @@ function QueryBuilderSearch({
const [isEditingTag, setIsEditingTag] = useState(false);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const {
updateTag,
handleClearTag,
@@ -121,6 +128,7 @@ function QueryBuilderSearch({
exampleQueries,
} = useAutoComplete(
query,
dotMetricsEnabled,
whereClauseConfig,
isLogsExplorerPage,
isInfraMonitoring,
@@ -138,6 +146,7 @@ function QueryBuilderSearch({
const { sourceKeys, handleRemoveSourceKey } = useFetchKeysAndValues(
searchValue,
query,
dotMetricsEnabled,
searchKey,
isLogsExplorerPage,
isInfraMonitoring,

View File

@@ -14,6 +14,8 @@ import { SelectOption } from 'types/common/select';
import { popupContainer } from 'utils/selectPopupContainer';
import { v4 as uuid } from 'uuid';
import { FeatureKeys } from '../../constants/features';
import { useAppContext } from '../../providers/App/App';
import QueryChip from './components/QueryChip';
import { QueryChipItem, SearchContainer } from './styles';
@@ -40,7 +42,12 @@ function ResourceAttributesFilter({
SelectOption<string, string>[]
>([]);
const resourceDeploymentKey = getResourceDeploymentKeys();
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const resourceDeploymentKey = getResourceDeploymentKeys(dotMetricsEnabled);
const [selectedEnvironments, setSelectedEnvironments] = useState<string[]>([]);
@@ -66,14 +73,14 @@ function ResourceAttributesFilter({
}, [queries, resourceDeploymentKey]);
useEffect(() => {
getEnvironmentTagKeys().then((tagKeys) => {
getEnvironmentTagKeys(dotMetricsEnabled).then((tagKeys) => {
if (tagKeys && Array.isArray(tagKeys) && tagKeys.length > 0) {
getEnvironmentTagValues().then((tagValues) => {
getEnvironmentTagValues(dotMetricsEnabled).then((tagValues) => {
setEnvironments(tagValues);
});
}
});
}, []);
}, [dotMetricsEnabled]);
return (
<div className="resourceAttributesFilter-container">

View File

@@ -3,6 +3,8 @@ import {
getResourceDeploymentKeys,
} from 'hooks/useResourceAttribute/utils';
import { FeatureKeys } from '../../../../constants/features';
import { useAppContext } from '../../../../providers/App/App';
import { QueryChipContainer, QueryChipItem } from '../../styles';
import { IQueryChipProps } from './types';
@@ -11,7 +13,13 @@ function QueryChip({ queryData, onClose }: IQueryChipProps): JSX.Element {
onClose(queryData.id);
};
const isClosable = queryData.tagKey !== getResourceDeploymentKeys();
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const isClosable =
queryData.tagKey !== getResourceDeploymentKeys(dotMetricsEnabled);
return (
<QueryChipContainer>

View File

@@ -4,6 +4,8 @@ import { useSelector } from 'react-redux';
import { AppState } from 'store/reducers';
import { GlobalReducer } from 'types/reducer/globalTime';
import { FeatureKeys } from '../../../constants/features';
import { useAppContext } from '../../../providers/App/App';
import { ServiceMetricsProps } from '../types';
import { getQueryRangeRequestData } from '../utils';
import ServiceMetricTable from './ServiceMetricTable';
@@ -16,13 +18,19 @@ function ServiceMetricsApplication({
GlobalReducer
>((state) => state.globalTime);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const queryRangeRequestData = useMemo(
() =>
getQueryRangeRequestData({
topLevelOperations,
globalSelectedInterval,
dotMetricsEnabled,
}),
[globalSelectedInterval, topLevelOperations],
[globalSelectedInterval, topLevelOperations, dotMetricsEnabled],
);
return (
<ServiceMetricTable

View File

@@ -19,10 +19,13 @@ import {
export const serviceMetricsQuery = (
topLevelOperation: [keyof ServiceDataProps, string[]],
dotMetricsEnabled: boolean,
): QueryBuilderData => {
const p99AutoCompleteData: BaseAutocompleteData = {
dataType: DataTypes.Float64,
key: WidgetKeys.SignozLatencyBucket,
key: dotMetricsEnabled
? WidgetKeys.Signoz_latency_bucket
: WidgetKeys.Signoz_latency_bucket_norm,
type: '',
};
@@ -50,7 +53,9 @@ export const serviceMetricsQuery = (
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -73,7 +78,9 @@ export const serviceMetricsQuery = (
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -83,7 +90,7 @@ export const serviceMetricsQuery = (
id: '',
key: {
dataType: DataTypes.Int64,
key: WidgetKeys.StatusCode,
key: dotMetricsEnabled ? WidgetKeys.StatusCode : WidgetKeys.StatusCodeNorm,
type: MetricsType.Tag,
},
op: OPERATORS.IN,
@@ -106,7 +113,9 @@ export const serviceMetricsQuery = (
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -129,7 +138,9 @@ export const serviceMetricsQuery = (
id: '',
key: {
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Resource,
},
op: OPERATORS.IN,
@@ -182,7 +193,9 @@ export const serviceMetricsQuery = (
const groupBy: BaseAutocompleteData[] = [
{
dataType: DataTypes.String,
key: WidgetKeys.OTelServiceName,
key: dotMetricsEnabled
? WidgetKeys.Service_name
: WidgetKeys.Service_name_norm,
type: MetricsType.Tag,
},
];

View File

@@ -17,6 +17,8 @@ import { AppState } from 'store/reducers';
import { GlobalReducer } from 'types/reducer/globalTime';
import { Tags } from 'types/reducer/trace';
import { FeatureKeys } from '../../../constants/features';
import { useAppContext } from '../../../providers/App/App';
import SkipOnBoardingModal from '../SkipOnBoardModal';
import ServiceTraceTable from './ServiceTracesTable';
@@ -38,6 +40,11 @@ function ServiceTraces(): JSX.Element {
selectedTags,
});
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
useErrorNotification(error);
const services = data || [];
@@ -55,7 +62,7 @@ function ServiceTraces(): JSX.Element {
useEffect(() => {
if (!logEventCalledRef.current && !isUndefined(data)) {
const selectedEnvironments = queries.find(
(val) => val.tagKey === getResourceDeploymentKeys(),
(val) => val.tagKey === getResourceDeploymentKeys(dotMetricsEnabled),
)?.tagValue;
const rps = data.reduce((total, service) => total + service.callRate, 0);

View File

@@ -26,6 +26,7 @@ export interface ServiceMetricsTableProps {
export interface GetQueryRangeRequestDataProps {
topLevelOperations: [keyof ServiceDataProps, string[]][];
globalSelectedInterval: Time | CustomTimeType;
dotMetricsEnabled: boolean;
}
export interface GetServiceListFromQueryProps {

View File

@@ -26,6 +26,7 @@ export function getSeriesValue(
export const getQueryRangeRequestData = ({
topLevelOperations,
globalSelectedInterval,
dotMetricsEnabled,
}: GetQueryRangeRequestDataProps): GetQueryResultsProps[] => {
const requestData: GetQueryResultsProps[] = [];
topLevelOperations.forEach((operation) => {
@@ -33,7 +34,7 @@ export const getQueryRangeRequestData = ({
query: {
queryType: EQueryType.QUERY_BUILDER,
promql: [],
builder: serviceMetricsQuery(operation),
builder: serviceMetricsQuery(operation, dotMetricsEnabled),
clickhouse_sql: [],
id: uuid(),
},

View File

@@ -27,6 +27,7 @@ export type WhereClauseConfig = {
export const useAutoComplete = (
query: IBuilderQuery,
dotMetricsEnabled: boolean,
whereClauseConfig?: WhereClauseConfig,
shouldUseSuggestions?: boolean,
isInfraMonitoring?: boolean,
@@ -39,6 +40,7 @@ export const useAutoComplete = (
const { keys, results, isFetching, exampleQueries } = useFetchKeysAndValues(
searchValue,
query,
dotMetricsEnabled,
searchKey,
shouldUseSuggestions,
isInfraMonitoring,

View File

@@ -48,6 +48,7 @@ type IuseFetchKeysAndValues = {
export const useFetchKeysAndValues = (
searchValue: string,
query: IBuilderQuery,
dotMetricsEnabled: boolean,
searchKey: string,
shouldUseSuggestions?: boolean,
isInfraMonitoring?: boolean,

View File

@@ -6,6 +6,8 @@ import { useSafeNavigate } from 'hooks/useSafeNavigate';
import useUrlQuery from 'hooks/useUrlQuery';
import { encode } from 'js-base64';
import { FeatureKeys } from '../../constants/features';
import { useAppContext } from '../../providers/App/App';
import { whilelistedKeys } from './config';
import { ResourceContext } from './context';
import {
@@ -56,6 +58,11 @@ function ResourceProvider({ children }: Props): JSX.Element {
}
};
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const dispatchQueries = useCallback(
(queries: IResourceAttribute[]): void => {
urlQuery.set(
@@ -71,7 +78,7 @@ function ResourceProvider({ children }: Props): JSX.Element {
const loadTagKeys = (): void => {
handleLoading(true);
GetTagKeys()
GetTagKeys(dotMetricsEnabled)
.then((tagKeys) => {
const options = mappingWithRoutesAndKeys(pathname, tagKeys);
setOptionsData({ options, mode: undefined });
@@ -154,15 +161,15 @@ function ResourceProvider({ children }: Props): JSX.Element {
setSelectedQueries([...value]);
},
[optionsData.mode, step, staging, pathname],
[optionsData.mode, step, staging, dotMetricsEnabled, pathname],
);
const handleEnvironmentChange = useCallback(
(environments: string[]): void => {
const staging = [getResourceDeploymentKeys(), 'IN'];
const staging = [getResourceDeploymentKeys(dotMetricsEnabled), 'IN'];
const queriesCopy = queries.filter(
(query) => query.tagKey !== getResourceDeploymentKeys(),
(query) => query.tagKey !== getResourceDeploymentKeys(dotMetricsEnabled),
);
if (environments && Array.isArray(environments) && environments.length > 0) {
@@ -177,7 +184,7 @@ function ResourceProvider({ children }: Props): JSX.Element {
setStep('Idle');
},
[dispatchQueries, queries],
[dispatchQueries, dotMetricsEnabled, queries],
);
const handleClose = useCallback(

View File

@@ -2,9 +2,13 @@ import { ReactNode } from 'react';
import { QueryClient, QueryClientProvider } from 'react-query';
import { Router } from 'react-router-dom';
import { act, renderHook, waitFor } from '@testing-library/react';
import { FeatureKeys } from 'constants/features';
import ROUTES from 'constants/routes';
import { createMemoryHistory, MemoryHistory } from 'history';
import { encode } from 'js-base64';
import { AppContext } from 'providers/App/App';
import { IAppContext } from 'providers/App/types';
import { getAppContextMock } from 'tests/test-utils';
import ResourceProvider from '../ResourceProvider';
import useResourceAttribute from '../useResourceAttribute';
@@ -51,8 +55,10 @@ const mockTagValues = getResourceAttributesTagValues as jest.MockedFunction<
function createWrapper({
routerHistory,
appContextOverrides,
}: {
routerHistory: MemoryHistory;
appContextOverrides?: Partial<IAppContext>;
}): ({ children }: { children: ReactNode }) => JSX.Element {
const queryClient = new QueryClient({
defaultOptions: { queries: { retry: false } },
@@ -60,9 +66,13 @@ function createWrapper({
return function Wrapper({ children }: { children: ReactNode }): JSX.Element {
return (
<QueryClientProvider client={queryClient}>
<Router history={routerHistory}>
<ResourceProvider>{children}</ResourceProvider>
</Router>
<AppContext.Provider
value={getAppContextMock('ADMIN', appContextOverrides)}
>
<Router history={routerHistory}>
<ResourceProvider>{children}</ResourceProvider>
</Router>
</AppContext.Provider>
</QueryClientProvider>
);
};
@@ -401,7 +411,7 @@ describe('ResourceProvider', () => {
});
describe('handleEnvironmentChange', () => {
it('adds a dotted environment query when envs are provided', async () => {
it('adds an environment query when envs are provided', async () => {
const routerHistory = createMemoryHistory({ initialEntries: ['/'] });
const { result } = renderHook(() => useResourceAttribute(), {
wrapper: createWrapper({ routerHistory }),
@@ -414,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',
operator: 'IN',
tagValue: ['production'],
});
@@ -425,7 +435,7 @@ describe('ResourceProvider', () => {
const seeded = [
{
id: 'env',
tagKey: 'resource_deployment.environment',
tagKey: 'resource_deployment_environment',
operator: 'IN',
tagValue: ['production'],
},
@@ -449,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');
expect(tagKeys).toContain('resource_service_name');
});
});
@@ -458,7 +468,7 @@ describe('ResourceProvider', () => {
const seeded = [
{
id: 'env',
tagKey: 'resource_deployment.environment',
tagKey: 'resource_deployment_environment',
operator: 'IN',
tagValue: ['production'],
},
@@ -476,13 +486,43 @@ describe('ResourceProvider', () => {
await waitFor(() => {
const envQueries = result.current.queries.filter(
(q) => q.tagKey === 'resource_deployment.environment',
(q) => q.tagKey === 'resource_deployment_environment',
);
expect(envQueries).toHaveLength(1);
expect(envQueries[0].tagValue).toStrictEqual(['staging']);
});
});
it('uses the dotted deployment env key when DOT_METRICS_ENABLED is active', async () => {
const routerHistory = createMemoryHistory({ initialEntries: ['/'] });
const { result } = renderHook(() => useResourceAttribute(), {
wrapper: createWrapper({
routerHistory,
appContextOverrides: {
featureFlags: [
{
name: FeatureKeys.DOT_METRICS_ENABLED,
active: true,
usage: 0,
usage_limit: -1,
route: '',
},
],
},
}),
});
act(() => {
result.current.handleEnvironmentChange(['production']);
});
await waitFor(() => {
expect(result.current.queries[0].tagKey).toBe(
'resource_deployment.environment',
);
});
});
it('preserves unrelated query params when dispatching', async () => {
const routerHistory = createMemoryHistory({
initialEntries: ['/?tab=overview'],

View File

@@ -5,13 +5,13 @@ import { mappingWithRoutesAndKeys } from '../utils';
describe('useResourceAttribute config', () => {
describe('whilelistedKeys', () => {
it('should include underscore-notation keys', () => {
it('should include underscore-notation keys (DOT_METRICS_ENABLED=false)', () => {
expect(whilelistedKeys).toContain('resource_deployment_environment');
expect(whilelistedKeys).toContain('resource_k8s_cluster_name');
expect(whilelistedKeys).toContain('resource_k8s_cluster_namespace');
});
it('should include dot-notation keys', () => {
it('should include dot-notation keys (DOT_METRICS_ENABLED=true)', () => {
expect(whilelistedKeys).toContain('resource_deployment.environment');
expect(whilelistedKeys).toContain('resource_k8s.cluster.name');
expect(whilelistedKeys).toContain('resource_k8s.cluster.namespace');

View File

@@ -144,11 +144,19 @@ export const OperatorSchema: IOption[] = OperatorConversions.map(
}),
);
export const getResourceDeploymentKeys = (): string =>
'resource_deployment.environment';
export const getResourceDeploymentKeys = (
dotMetricsEnabled: boolean,
): string => {
if (dotMetricsEnabled) {
return 'resource_deployment.environment';
}
return 'resource_deployment_environment';
};
export const GetTagKeys = async (): Promise<IOption[]> => {
const resourceDeploymentKey = getResourceDeploymentKeys();
export const GetTagKeys = async (
dotMetricsEnabled: boolean,
): Promise<IOption[]> => {
const resourceDeploymentKey = getResourceDeploymentKeys(dotMetricsEnabled);
const { payload } = await getResourceAttributesTagKeys({
metricName: 'signoz_calls_total',
match: 'resource_',
@@ -168,10 +176,12 @@ export const GetTagKeys = async (): Promise<IOption[]> => {
}));
};
export const getEnvironmentTagKeys = async (): Promise<IOption[]> => {
export const getEnvironmentTagKeys = async (
dotMetricsEnabled: boolean,
): Promise<IOption[]> => {
const { payload } = await getResourceAttributesTagKeys({
metricName: 'signoz_calls_total',
match: getResourceDeploymentKeys(),
match: getResourceDeploymentKeys(dotMetricsEnabled),
});
if (!payload || !payload?.data) {
return [];
@@ -184,9 +194,11 @@ export const getEnvironmentTagKeys = async (): Promise<IOption[]> => {
}));
};
export const getEnvironmentTagValues = async (): Promise<IOption[]> => {
export const getEnvironmentTagValues = async (
dotMetricsEnabled: boolean,
): Promise<IOption[]> => {
const { payload } = await getResourceAttributesTagValues({
tagKey: getResourceDeploymentKeys(),
tagKey: getResourceDeploymentKeys(dotMetricsEnabled),
metricName: 'signoz_calls_total',
});

View File

@@ -4,6 +4,8 @@ import { CardContainer } from 'container/GridCardLayout/styles';
import { useIsDarkMode } from 'hooks/useDarkMode';
import { Widgets } from 'types/api/dashboard/getAll';
import { FeatureKeys } from '../../../../constants/features';
import { useAppContext } from '../../../../providers/App/App';
import MetricPageGridGraph from './MetricPageGraph';
import {
getAverageRequestLatencyWidgetData,
@@ -71,15 +73,20 @@ function MetricColumnGraphs({
}): JSX.Element {
const { t } = useTranslation('messagingQueues');
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const metricsData = [
{
title: t('metricGraphCategory.brokerMetrics.title'),
description: t('metricGraphCategory.brokerMetrics.description'),
graphCount: [
getBrokerCountWidgetData(),
getRequestTimesWidgetData(),
getProducerFetchRequestPurgatoryWidgetData(),
getBrokerNetworkThroughputWidgetData(),
getBrokerCountWidgetData(dotMetricsEnabled),
getRequestTimesWidgetData(dotMetricsEnabled),
getProducerFetchRequestPurgatoryWidgetData(dotMetricsEnabled),
getBrokerNetworkThroughputWidgetData(dotMetricsEnabled),
],
id: 'broker-metrics',
},
@@ -87,11 +94,11 @@ function MetricColumnGraphs({
title: t('metricGraphCategory.producerMetrics.title'),
description: t('metricGraphCategory.producerMetrics.description'),
graphCount: [
getIoWaitTimeWidgetData(),
getRequestResponseWidgetData(),
getAverageRequestLatencyWidgetData(),
getKafkaProducerByteRateWidgetData(),
getBytesConsumedWidgetData(),
getIoWaitTimeWidgetData(dotMetricsEnabled),
getRequestResponseWidgetData(dotMetricsEnabled),
getAverageRequestLatencyWidgetData(dotMetricsEnabled),
getKafkaProducerByteRateWidgetData(dotMetricsEnabled),
getBytesConsumedWidgetData(dotMetricsEnabled),
],
id: 'producer-metrics',
},
@@ -99,11 +106,11 @@ function MetricColumnGraphs({
title: t('metricGraphCategory.consumerMetrics.title'),
description: t('metricGraphCategory.consumerMetrics.description'),
graphCount: [
getConsumerOffsetWidgetData(),
getConsumerGroupMemberWidgetData(),
getConsumerLagByGroupWidgetData(),
getConsumerFetchRateWidgetData(),
getMessagesConsumedWidgetData(),
getConsumerOffsetWidgetData(dotMetricsEnabled),
getConsumerGroupMemberWidgetData(dotMetricsEnabled),
getConsumerLagByGroupWidgetData(dotMetricsEnabled),
getConsumerFetchRateWidgetData(dotMetricsEnabled),
getMessagesConsumedWidgetData(dotMetricsEnabled),
],
id: 'consumer-metrics',
},

View File

@@ -8,6 +8,8 @@ import { useIsDarkMode } from 'hooks/useDarkMode';
import { ChevronDown, ChevronUp } from '@signozhq/icons';
import { Widgets } from 'types/api/dashboard/getAll';
import { FeatureKeys } from '../../../../constants/features';
import { useAppContext } from '../../../../providers/App/App';
import MetricColumnGraphs from './MetricColumnGraphs';
import MetricPageGridGraph from './MetricPageGraph';
import {
@@ -95,6 +97,11 @@ function MetricPage(): JSX.Element {
}));
};
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const { t } = useTranslation('messagingQueues');
const metricSections = [
@@ -103,10 +110,10 @@ function MetricPage(): JSX.Element {
title: t('metricGraphCategory.brokerJVMMetrics.title'),
description: t('metricGraphCategory.brokerJVMMetrics.description'),
graphCount: [
getJvmGCCountWidgetData(),
getJvmGcCollectionsElapsedWidgetData(),
getCpuRecentUtilizationWidgetData(),
getJvmMemoryHeapWidgetData(),
getJvmGCCountWidgetData(dotMetricsEnabled),
getJvmGcCollectionsElapsedWidgetData(dotMetricsEnabled),
getCpuRecentUtilizationWidgetData(dotMetricsEnabled),
getJvmMemoryHeapWidgetData(dotMetricsEnabled),
],
},
{
@@ -114,10 +121,10 @@ function MetricPage(): JSX.Element {
title: t('metricGraphCategory.partitionMetrics.title'),
description: t('metricGraphCategory.partitionMetrics.description'),
graphCount: [
getPartitionCountPerTopicWidgetData(),
getCurrentOffsetPartitionWidgetData(),
getOldestOffsetWidgetData(),
getInsyncReplicasWidgetData(),
getPartitionCountPerTopicWidgetData(dotMetricsEnabled),
getCurrentOffsetPartitionWidgetData(dotMetricsEnabled),
getOldestOffsetWidgetData(dotMetricsEnabled),
getInsyncReplicasWidgetData(dotMetricsEnabled),
],
},
];
@@ -131,7 +138,7 @@ function MetricPage(): JSX.Element {
// Only log when first graph has rendered and we haven't logged yet
if (renderedGraphCountRef.current === 1 && !hasLoggedRef.current) {
void logEvent('MQ Kafka: Metric view', {
logEvent('MQ Kafka: Metric view', {
graphRendered: true,
});
hasLoggedRef.current = true;

View File

@@ -78,15 +78,21 @@ export function getWidgetQuery(
};
}
export const getRequestTimesWidgetData = (): Widgets =>
export const getRequestTimesWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.request.time.avg',
id: 'kafka.request.time.avg--float64--Gauge--true',
// choose key based on flag
key: dotMetricsEnabled
? 'kafka.request.time.avg'
: 'kafka_request_time_avg',
// mirror into the id as well
id: 'kafka_request_time_avg--float64--Gauge--true',
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -116,15 +122,15 @@ export const getRequestTimesWidgetData = (): Widgets =>
}),
);
export const getBrokerCountWidgetData = (): Widgets =>
export const getBrokerCountWidgetData = (dotMetricsEnabled: boolean): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.brokers',
id: 'kafka.brokers--float64--Gauge--true',
key: dotMetricsEnabled ? 'kafka.brokers' : 'kafka_brokers',
id: 'kafka_brokers--float64--Gauge--true',
type: 'Gauge',
},
aggregateOperator: 'sum',
@@ -150,15 +156,20 @@ export const getBrokerCountWidgetData = (): Widgets =>
}),
);
export const getProducerFetchRequestPurgatoryWidgetData = (): Widgets =>
export const getProducerFetchRequestPurgatoryWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.purgatory.size',
id: 'kafka.purgatory.size--float64--Gauge--true',
// inline ternary based on dotMetricsEnabled
key: dotMetricsEnabled ? 'kafka.purgatory.size' : 'kafka_purgatory_size',
id: `${
dotMetricsEnabled ? 'kafka.purgatory.size' : 'kafka_purgatory_size'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -185,15 +196,24 @@ export const getProducerFetchRequestPurgatoryWidgetData = (): Widgets =>
}),
);
export const getBrokerNetworkThroughputWidgetData = (): Widgets =>
export const getBrokerNetworkThroughputWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka_server_brokertopicmetrics_total_replicationbytesinpersec_oneminuterate',
id: 'kafka_server_brokertopicmetrics_total_replicationbytesinpersec_oneminuterate--float64--Gauge--true',
// inline ternary based on dotMetricsEnabled
key: dotMetricsEnabled
? 'kafka_server_brokertopicmetrics_total_replicationbytesinpersec_oneminuterate'
: 'kafka_server_brokertopicmetrics_bytesoutpersec_oneminuterate',
id: `${
dotMetricsEnabled
? 'kafka_server_brokertopicmetrics_total_replicationbytesinpersec_oneminuterate'
: 'kafka_server_brokertopicmetrics_bytesoutpersec_oneminuterate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -220,15 +240,22 @@ export const getBrokerNetworkThroughputWidgetData = (): Widgets =>
}),
);
export const getIoWaitTimeWidgetData = (): Widgets =>
export const getIoWaitTimeWidgetData = (dotMetricsEnabled: boolean): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.producer.io_waittime_total',
id: 'kafka.producer.io_waittime_total--float64--Sum--true',
// inline ternary based on dotMetricsEnabled
key: dotMetricsEnabled
? 'kafka.producer.io_waittime_total'
: 'kafka_producer_io_waittime_total',
id: `${
dotMetricsEnabled
? 'kafka.producer.io_waittime_total'
: 'kafka_producer_io_waittime_total'
}--float64--Sum--true`,
type: 'Sum',
},
aggregateOperator: 'rate',
@@ -255,15 +282,23 @@ export const getIoWaitTimeWidgetData = (): Widgets =>
}),
);
export const getRequestResponseWidgetData = (): Widgets =>
export const getRequestResponseWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.producer.request_rate',
id: 'kafka.producer.request_rate--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.producer.request_rate'
: 'kafka_producer_request_rate',
id: `${
dotMetricsEnabled
? 'kafka.producer.request_rate'
: 'kafka_producer_request_rate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -286,8 +321,14 @@ export const getRequestResponseWidgetData = (): Widgets =>
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.producer.response_rate',
id: 'kafka.producer.response_rate--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.producer.response_rate'
: 'kafka_producer_response_rate',
id: `${
dotMetricsEnabled
? 'kafka.producer.response_rate'
: 'kafka_producer_response_rate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -314,15 +355,23 @@ export const getRequestResponseWidgetData = (): Widgets =>
}),
);
export const getAverageRequestLatencyWidgetData = (): Widgets =>
export const getAverageRequestLatencyWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.producer.request_latency_avg',
id: 'kafka.producer.request_latency_avg--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.producer.request_latency_avg'
: 'kafka_producer_request_latency_avg',
id: `${
dotMetricsEnabled
? 'kafka.producer.request_latency_avg'
: 'kafka_producer_request_latency_avg'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -349,15 +398,23 @@ export const getAverageRequestLatencyWidgetData = (): Widgets =>
}),
);
export const getKafkaProducerByteRateWidgetData = (): Widgets =>
export const getKafkaProducerByteRateWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.producer.byte_rate',
id: 'kafka.producer.byte_rate--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.producer.byte_rate'
: 'kafka_producer_byte_rate',
id: `${
dotMetricsEnabled
? 'kafka.producer.byte_rate'
: 'kafka_producer_byte_rate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -385,21 +442,31 @@ export const getKafkaProducerByteRateWidgetData = (): Widgets =>
timeAggregation: 'avg',
},
],
title: 'kafka.producer.byte_rate',
title: dotMetricsEnabled
? 'kafka.producer.byte_rate'
: 'kafka_producer_byte_rate',
description:
'Helps measure the data output rate from the producer, indicating the load a producer is placing on Kafka brokers.',
}),
);
export const getBytesConsumedWidgetData = (): Widgets =>
export const getBytesConsumedWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.consumer.bytes_consumed_rate',
id: 'kafka.consumer.bytes_consumed_rate--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.consumer.bytes_consumed_rate'
: 'kafka_consumer_bytes_consumed_rate',
id: `${
dotMetricsEnabled
? 'kafka.consumer.bytes_consumed_rate'
: 'kafka_consumer_bytes_consumed_rate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -427,15 +494,23 @@ export const getBytesConsumedWidgetData = (): Widgets =>
}),
);
export const getConsumerOffsetWidgetData = (): Widgets =>
export const getConsumerOffsetWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.consumer_group.offset',
id: 'kafka.consumer_group.offset--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.consumer_group.offset'
: 'kafka_consumer_group_offset',
id: `${
dotMetricsEnabled
? 'kafka.consumer_group.offset'
: 'kafka_consumer_group_offset'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -481,15 +556,23 @@ export const getConsumerOffsetWidgetData = (): Widgets =>
}),
);
export const getConsumerGroupMemberWidgetData = (): Widgets =>
export const getConsumerGroupMemberWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.consumer_group.members',
id: 'kafka.consumer_group.members--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.consumer_group.members'
: 'kafka_consumer_group_members',
id: `${
dotMetricsEnabled
? 'kafka.consumer_group.members'
: 'kafka_consumer_group_members'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'sum',
@@ -522,15 +605,23 @@ export const getConsumerGroupMemberWidgetData = (): Widgets =>
}),
);
export const getConsumerLagByGroupWidgetData = (): Widgets =>
export const getConsumerLagByGroupWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.consumer_group.lag',
id: 'kafka.consumer_group.lag--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.consumer_group.lag'
: 'kafka_consumer_group_lag',
id: `${
dotMetricsEnabled
? 'kafka.consumer_group.lag'
: 'kafka_consumer_group_lag'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -576,15 +667,23 @@ export const getConsumerLagByGroupWidgetData = (): Widgets =>
}),
);
export const getConsumerFetchRateWidgetData = (): Widgets =>
export const getConsumerFetchRateWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.consumer.fetch_rate',
id: 'kafka.consumer.fetch_rate--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.consumer.fetch_rate'
: 'kafka_consumer_fetch_rate',
id: `${
dotMetricsEnabled
? 'kafka.consumer.fetch_rate'
: 'kafka_consumer_fetch_rate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -597,7 +696,7 @@ export const getConsumerFetchRateWidgetData = (): Widgets =>
{
dataType: DataTypes.String,
id: 'service_name--string--tag--false',
key: 'service.name',
key: dotMetricsEnabled ? 'service.name' : 'service_name',
type: 'tag',
},
],
@@ -618,15 +717,23 @@ export const getConsumerFetchRateWidgetData = (): Widgets =>
}),
);
export const getMessagesConsumedWidgetData = (): Widgets =>
export const getMessagesConsumedWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.consumer.records_consumed_rate',
id: 'kafka.consumer.records_consumed_rate--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.consumer.records_consumed_rate'
: 'kafka_consumer_records_consumed_rate',
id: `${
dotMetricsEnabled
? 'kafka.consumer.records_consumed_rate'
: 'kafka_consumer_records_consumed_rate'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -653,15 +760,21 @@ export const getMessagesConsumedWidgetData = (): Widgets =>
}),
);
export const getJvmGCCountWidgetData = (): Widgets =>
export const getJvmGCCountWidgetData = (dotMetricsEnabled: boolean): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'jvm.gc.collections.count',
id: 'jvm.gc.collections.count--float64--Sum--true',
key: dotMetricsEnabled
? 'jvm.gc.collections.count'
: 'jvm_gc_collections_count',
id: `${
dotMetricsEnabled
? 'jvm.gc.collections.count'
: 'jvm_gc_collections_count'
}--float64--Sum--true`,
type: 'Sum',
},
aggregateOperator: 'rate',
@@ -688,15 +801,23 @@ export const getJvmGCCountWidgetData = (): Widgets =>
}),
);
export const getJvmGcCollectionsElapsedWidgetData = (): Widgets =>
export const getJvmGcCollectionsElapsedWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'jvm.gc.collections.elapsed',
id: 'jvm.gc.collections.elapsed--float64--Sum--true',
key: dotMetricsEnabled
? 'jvm.gc.collections.elapsed'
: 'jvm_gc_collections_elapsed',
id: `${
dotMetricsEnabled
? 'jvm.gc.collections.elapsed'
: 'jvm_gc_collections_elapsed'
}--float64--Sum--true`,
type: 'Sum',
},
aggregateOperator: 'rate',
@@ -717,21 +838,31 @@ export const getJvmGcCollectionsElapsedWidgetData = (): Widgets =>
timeAggregation: 'rate',
},
],
title: 'jvm.gc.collections.elapsed',
title: dotMetricsEnabled
? 'jvm.gc.collections.elapsed'
: 'jvm_gc_collections_elapsed',
description:
'Measures the total time (usually in milliseconds) spent on garbage collection (GC) events in the Java Virtual Machine (JVM).',
}),
);
export const getCpuRecentUtilizationWidgetData = (): Widgets =>
export const getCpuRecentUtilizationWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'jvm.cpu.recent_utilization',
id: 'jvm.cpu.recent_utilization--float64--Gauge--true',
key: dotMetricsEnabled
? 'jvm.cpu.recent_utilization'
: 'jvm_cpu_recent_utilization',
id: `${
dotMetricsEnabled
? 'jvm.cpu.recent_utilization'
: 'jvm_cpu_recent_utilization'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -758,15 +889,19 @@ export const getCpuRecentUtilizationWidgetData = (): Widgets =>
}),
);
export const getJvmMemoryHeapWidgetData = (): Widgets =>
export const getJvmMemoryHeapWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'jvm.memory.heap.max',
id: 'jvm.memory.heap.max--float64--Gauge--true',
key: dotMetricsEnabled ? 'jvm.memory.heap.max' : 'jvm_memory_heap_max',
id: `${
dotMetricsEnabled ? 'jvm.memory.heap.max' : 'jvm_memory_heap_max'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -793,15 +928,21 @@ export const getJvmMemoryHeapWidgetData = (): Widgets =>
}),
);
export const getPartitionCountPerTopicWidgetData = (): Widgets =>
export const getPartitionCountPerTopicWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.topic.partitions',
id: 'kafka.topic.partitions--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.topic.partitions'
: 'kafka_topic_partitions',
id: `${
dotMetricsEnabled ? 'kafka.topic.partitions' : 'kafka_topic_partitions'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'sum',
@@ -834,15 +975,23 @@ export const getPartitionCountPerTopicWidgetData = (): Widgets =>
}),
);
export const getCurrentOffsetPartitionWidgetData = (): Widgets =>
export const getCurrentOffsetPartitionWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.partition.current_offset',
id: 'kafka.partition.current_offset--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.partition.current_offset'
: 'kafka_partition_current_offset',
id: `${
dotMetricsEnabled
? 'kafka.partition.current_offset'
: 'kafka_partition_current_offset'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -882,15 +1031,23 @@ export const getCurrentOffsetPartitionWidgetData = (): Widgets =>
}),
);
export const getOldestOffsetWidgetData = (): Widgets =>
export const getOldestOffsetWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.partition.oldest_offset',
id: 'kafka.partition.oldest_offset--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.partition.oldest_offset'
: 'kafka_partition_oldest_offset',
id: `${
dotMetricsEnabled
? 'kafka.partition.oldest_offset'
: 'kafka_partition_oldest_offset'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',
@@ -930,15 +1087,23 @@ export const getOldestOffsetWidgetData = (): Widgets =>
}),
);
export const getInsyncReplicasWidgetData = (): Widgets =>
export const getInsyncReplicasWidgetData = (
dotMetricsEnabled: boolean,
): Widgets =>
getWidgetQueryBuilder(
getWidgetQuery({
queryData: [
{
aggregateAttribute: {
dataType: DataTypes.Float64,
key: 'kafka.partition.replicas_in_sync',
id: 'kafka.partition.replicas_in_sync--float64--Gauge--true',
key: dotMetricsEnabled
? 'kafka.partition.replicas_in_sync'
: 'kafka_partition_replicas_in_sync',
id: `${
dotMetricsEnabled
? 'kafka.partition.replicas_in_sync'
: 'kafka_partition_replicas_in_sync'
}--float64--Gauge--true`,
type: 'Gauge',
},
aggregateOperator: 'avg',

View File

@@ -11,6 +11,8 @@ import useDebouncedFn from 'hooks/useDebouncedFunction';
import useUrlQuery from 'hooks/useUrlQuery';
import { Check, Share2 } from '@signozhq/icons';
import { FeatureKeys } from '../../../constants/features';
import { useAppContext } from '../../../providers/App/App';
import { useGetAllConfigOptions } from './useGetAllConfigOptions';
import './MQConfigOptions.styles.scss';
@@ -38,11 +40,19 @@ const useConfigOptions = (
isFetching: boolean;
options: DefaultOptionType[];
} => {
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const [searchText, setSearchText] = useState<string>('');
const { isFetching, options } = useGetAllConfigOptions({
attributeKey: type,
searchText,
});
const { isFetching, options } = useGetAllConfigOptions(
{
attributeKey: type,
searchText,
},
dotMetricsEnabled,
);
const handleDebouncedSearch = useDebouncedFn((searchText): void => {
setSearchText(searchText as string);
}, 500);

View File

@@ -3,6 +3,7 @@ import { useCallback, useMemo, useRef } from 'react';
import { useDispatch } from 'react-redux';
import { useHistory, useLocation } from 'react-router-dom';
import logEvent from 'api/common/logEvent';
import { FeatureKeys } from 'constants/features';
import { QueryParams } from 'constants/query';
import { PANEL_TYPES } from 'constants/queryBuilder';
import { ViewMenuAction } from 'container/GridCardLayout/config';
@@ -11,6 +12,7 @@ import { Card } from 'container/GridCardLayout/styles';
import { getWidgetQueryBuilder } from 'container/MetricsApplication/MetricsApplication.factory';
import { useIsDarkMode } from 'hooks/useDarkMode';
import useUrlQuery from 'hooks/useUrlQuery';
import { useAppContext } from 'providers/App/App';
import { UpdateTimeInterval } from 'store/actions';
import {
@@ -32,9 +34,15 @@ function MessagingQueuesGraph(): JSX.Element {
[consumerGrp, topic, partition],
);
const { featureFlags } = useAppContext();
const dotMetricsEnabled =
featureFlags?.find((flag) => flag.name === FeatureKeys.DOT_METRICS_ENABLED)
?.active || false;
const widgetData = useMemo(
() => getWidgetQueryBuilder(getWidgetQuery({ filterItems })),
[filterItems],
() =>
getWidgetQueryBuilder(getWidgetQuery({ filterItems, dotMetricsEnabled })),
[filterItems, dotMetricsEnabled],
);
const history = useHistory();
@@ -73,7 +81,7 @@ function MessagingQueuesGraph(): JSX.Element {
const checkIfDataExists = (isDataAvailable: boolean): void => {
if (!isLogEventCalled.current) {
isLogEventCalled.current = true;
void logEvent('Messaging Queues: Graph data fetched', {
logEvent('Messaging Queues: Graph data fetched', {
isDataAvailable,
});
}

View File

@@ -16,6 +16,7 @@ export interface GetAllConfigOptionsResponse {
export function useGetAllConfigOptions(
props: ConfigOptions,
dotMetricsEnabled: boolean,
): GetAllConfigOptionsResponse {
const { attributeKey, searchText } = props;
@@ -25,7 +26,9 @@ export function useGetAllConfigOptions(
const { payload } = await getAttributesValues({
aggregateOperator: 'avg',
dataSource: DataSource.METRICS,
aggregateAttribute: 'kafka.consumer_group.lag',
aggregateAttribute: dotMetricsEnabled
? 'kafka.consumer_group.lag'
: 'kafka_consumer_group_lag',
attributeKey,
searchText: searchText ?? '',
filterAttributeKeyDataType: DataTypes.String,

View File

@@ -94,8 +94,10 @@ export function getFiltersFromConfigOptions(
export function getWidgetQuery({
filterItems,
dotMetricsEnabled,
}: {
filterItems: TagFilterItem[];
dotMetricsEnabled: boolean;
}): GetWidgetQueryBuilderProps {
return {
title: 'Consumer Lag',
@@ -110,8 +112,14 @@ export function getWidgetQuery({
{
aggregateAttribute: {
dataType: DataTypes.Float64,
id: 'kafka.consumer_group.lag--float64--Gauge--true',
key: 'kafka.consumer_group.lag',
id: `${
dotMetricsEnabled
? 'kafka.consumer_group.lag'
: 'kafka_consumer_group_lag'
}--float64--Gauge--true`,
key: dotMetricsEnabled
? 'kafka.consumer_group.lag'
: 'kafka_consumer_group_lag',
type: 'Gauge',
},
aggregateOperator: 'max',

View File

@@ -3190,6 +3190,11 @@ 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 {
@@ -3201,7 +3206,7 @@ func (r *ClickHouseReader) GetMetricAttributeValues(ctx context.Context, orgID v
query = query + fmt.Sprintf(" LIMIT %d;", req.Limit)
}
names := []string{req.AggregateAttribute}
names = append(names, metrics.GetTransitionedMetric(req.AggregateAttribute))
names = append(names, metrics.GetTransitionedMetric(req.AggregateAttribute, normalized))
rows, err = r.db.Query(ctx, query, req.FilterAttributeKey, names, req.FilterAttributeKey, fmt.Sprintf("%%%s%%", req.SearchText), common.PastDayRoundOff())
@@ -5443,3 +5448,112 @@ func (r *ClickHouseReader) SearchTraces(ctx context.Context, params *model.Searc
return &searchSpansResult, nil
}
func (r *ClickHouseReader) GetNormalizedStatus(
ctx context.Context,
orgID valuer.UUID,
metricNames []string,
) (map[string]bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalMetrics.StringValue(),
instrumentationtypes.CodeNamespace: "clickhouse-reader",
instrumentationtypes.CodeFunctionName: "GetNormalizedStatus",
})
if len(metricNames) == 0 {
return map[string]bool{}, nil
}
result := make(map[string]bool, len(metricNames))
buildKey := func(name string) string {
return constants.NormalizedMetricsMapCacheKey + ":" + name
}
uncached := make([]string, 0, len(metricNames))
for _, m := range metricNames {
var status model.MetricsNormalizedMap
if err := r.cache.Get(ctx, orgID, buildKey(m), &status); err == nil {
result[m] = status.IsUnNormalized
} else {
uncached = append(uncached, m)
}
}
if len(uncached) == 0 {
return result, nil
}
placeholders := "'" + strings.Join(uncached, "', '") + "'"
reductionEnabled := r.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableMetricsReduction, featuretypes.NewFlaggerEvaluationContext(orgID))
var q string
if reductionEnabled {
q = fmt.Sprintf(
`SELECT metric_name, toUInt8(__normalized)
FROM (
SELECT metric_name, __normalized FROM %s.%s WHERE metric_name IN (%s)
UNION ALL
SELECT metric_name, __normalized FROM %s.%s WHERE metric_name IN (%s)
)
GROUP BY metric_name, __normalized`,
signozMetricDBName, signozTSTableNameV41Day, placeholders,
signozMetricDBName, signozTSTableNameV4Reduced, placeholders,
)
} else {
q = fmt.Sprintf(
`SELECT metric_name, toUInt8(__normalized)
FROM %s.%s
WHERE metric_name IN (%s)
GROUP BY metric_name, __normalized`,
signozMetricDBName, signozTSTableNameV41Day, placeholders,
)
}
rows, err := r.db.Query(ctx, q)
if err != nil {
return nil, err
}
defer rows.Close()
// tmp[m] collects the set {0,1} for a metric name, truth table
tmp := make(map[string]map[uint8]struct{}, len(uncached))
for rows.Next() {
var (
name string
normalized uint8
)
if err := rows.Scan(&name, &normalized); err != nil {
return nil, err
}
if _, ok := tmp[name]; !ok {
tmp[name] = make(map[uint8]struct{}, 2)
}
tmp[name][normalized] = struct{}{}
}
if err := rows.Err(); err != nil {
return nil, err
}
for _, m := range uncached {
set := tmp[m]
switch {
case len(set) == 0:
return nil, fmt.Errorf("metric %q not found in ClickHouse", m)
case len(set) == 2:
result[m] = true
default:
_, hasUnnorm := set[0]
result[m] = hasUnnorm
}
status := model.MetricsNormalizedMap{
MetricName: m,
IsUnNormalized: result[m],
}
_ = r.cache.Set(ctx, orgID, buildKey(m), &status, 0)
}
return result, nil
}

View File

@@ -56,6 +56,7 @@ import (
"github.com/SigNoz/signoz/pkg/query-service/app/queryBuilder"
tracesV3 "github.com/SigNoz/signoz/pkg/query-service/app/traces/v3"
tracesV4 "github.com/SigNoz/signoz/pkg/query-service/app/traces/v4"
"github.com/SigNoz/signoz/pkg/query-service/constants"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/query-service/postprocess"
"github.com/SigNoz/signoz/pkg/types"
@@ -1051,6 +1052,10 @@ func prepareQuery(r *http.Request) (string, error) {
return "", tmplErr
}
if !constants.IsDotMetricsEnabled {
return queryBuf.String(), nil
}
query = queryBuf.String()
// Now handle $var replacements (simple string replace)
@@ -1603,6 +1608,13 @@ func (aH *APIHandler) getFeatureFlags(w http.ResponseWriter, r *http.Request) {
Route: "",
})
if constants.IsDotMetricsEnabled {
for idx, feature := range featureSet {
if feature.Name == licensetypes.DotMetricsEnabled {
featureSet[idx].Active = true
}
}
}
aH.Respond(w, featureSet)
}
@@ -2043,8 +2055,12 @@ func (aH *APIHandler) onboardKafka(w http.ResponseWriter, r *http.Request) {
}
}
}
var kafkaConsumerFetchLatencyAvg string = "kafka.consumer.fetch_latency_avg"
var kafkaConsumerLag string = "kafka.consumer_group.lag"
var kafkaConsumerFetchLatencyAvg string = "kafka_consumer_fetch_latency_avg"
var kafkaConsumerLag string = "kafka_consumer_group_lag"
if constants.IsDotMetricsEnabled {
kafkaConsumerLag = "kafka.consumer_group.lag"
kafkaConsumerFetchLatencyAvg = "kafka.consumer.fetch_latency_avg"
}
if !fetchLatencyState && !consumerLagState {
entries = append(entries, kafka.OnboardingResponse{

View File

@@ -18,12 +18,12 @@ import (
)
var (
metricToUseForClusters = "k8s.node.cpu.usage"
metricToUseForClusters = GetDotMetrics("k8s_node_cpu_usage")
clusterAttrsToEnrich = []string{"k8s.cluster.name"}
clusterAttrsToEnrich = []string{GetDotMetrics("k8s_cluster_name")}
// TODO(srikanthccv): change this to k8s_cluster_uid after showing the missing data banner
k8sClusterUIDAttrKey = "k8s.cluster.name"
k8sClusterUIDAttrKey = GetDotMetrics("k8s_cluster_name")
queryNamesForClusters = map[string][]string{
"cpu": {"A"},

View File

@@ -9,6 +9,250 @@ import (
"github.com/SigNoz/signoz/pkg/query-service/model"
)
var dotMetricMap = map[string]string{
"system_uptime": "system.uptime",
"system_cpu_physical_count": "system.cpu.physical.count",
"system_cpu_logical_count": "system.cpu.logical.count",
"system_cpu_time": "system.cpu.time",
"system_cpu_frequency": "system.cpu.frequency",
"system_cpu_utilization": "system.cpu.utilization",
"system_cpu_load_average_15m": "system.cpu.load_average.15m",
"system_memory_usage": "system.memory.usage",
"system_memory_limit": "system.memory.limit",
"system_memory_utilization": "system.memory.utilization",
"system_memory_linux_available": "system.memory.linux.available",
"system_memory_linux_shared": "system.memory.linux.shared",
"system_memory_linux_slab_usage": "system.memory.linux.slab.usage",
"system_paging_usage": "system.paging.usage",
"system_paging_utilization": "system.paging.utilization",
"system_paging_faults": "system.paging.faults",
"system_paging_operations": "system.paging.operations",
"system_disk_io": "system.disk.io",
"system_disk_operations": "system.disk.operations",
"system_disk_io_time": "system.disk.io_time",
"system_disk_operation_time": "system.disk.operation_time",
"system_disk_merged": "system.disk.merged",
"system_disk_limit": "system.disk.limit",
"system_filesystem_usage": "system.filesystem.usage",
"system_filesystem_utilization": "system.filesystem.utilization",
"system_filesystem_limit": "system.filesystem.limit",
"system_network_errors": "system.network.errors",
"system_network_io": "system.network.io",
"system_network_connections": "system.network.connections",
"system_network_dropped": "system.network.dropped",
"system_network_packets": "system.network.packets",
"system_processes_count": "system.processes.count",
"system_processes_created": "system.processes.created",
"system_disk_pending_operations": "system.disk.pending_operations",
"system_disk_weighted_io_time": "system.disk.weighted_io_time",
"system_filesystem_inodes_usage": "system.filesystem.inodes.usage",
"system_network_conntrack_count": "system.network.conntrack.count",
"system_network_conntrack_max": "system.network.conntrack.max",
"system_cpu_load_average_1m": "system.cpu.load_average.1m",
"system_cpu_load_average_5m": "system.cpu.load_average.5m",
"host_name": "host.name",
"k8s_cluster_name": "k8s.cluster.name",
"k8s_node_name": "k8s.node.name",
"k8s_pod_memory_usage": "k8s.pod.memory.usage",
"k8s_pod_cpu_request_utilization": "k8s.pod.cpu_request_utilization",
"k8s_pod_memory_request_utilization": "k8s.pod.memory_request_utilization",
"k8s_pod_cpu_limit_utilization": "k8s.pod.cpu_limit_utilization",
"k8s_pod_memory_limit_utilization": "k8s.pod.memory_limit_utilization",
"k8s_container_restarts": "k8s.container.restarts",
"k8s_pod_phase": "k8s.pod.phase",
"k8s_node_allocatable_cpu": "k8s.node.allocatable_cpu",
"k8s_node_allocatable_memory": "k8s.node.allocatable_memory",
"k8s_node_memory_usage": "k8s.node.memory.usage",
"k8s_node_condition_ready": "k8s.node.condition_ready",
"k8s_daemonset_desired_scheduled_nodes": "k8s.daemonset.desired_scheduled_nodes",
"k8s_daemonset_current_scheduled_nodes": "k8s.daemonset.current_scheduled_nodes",
"k8s_deployment_desired": "k8s.deployment.desired",
"k8s_deployment_available": "k8s.deployment.available",
"k8s_job_desired_successful_pods": "k8s.job.desired_successful_pods",
"k8s_job_active_pods": "k8s.job.active_pods",
"k8s_job_failed_pods": "k8s.job.failed_pods",
"k8s_job_successful_pods": "k8s.job.successful_pods",
"k8s_statefulset_desired_pods": "k8s.statefulset.desired_pods",
"k8s_statefulset_current_pods": "k8s.statefulset.current_pods",
"k8s_namespace_name": "k8s.namespace.name",
"k8s_deployment_name": "k8s.deployment.name",
"k8s_cronjob_name": "k8s.cronjob.name",
"k8s_job_name": "k8s.job.name",
"k8s_daemonset_name": "k8s.daemonset.name",
"os_type": "os.type",
"process_cgroup": "process.cgroup",
"process_pid": "process.pid",
"process_parent_pid": "process.parent_pid",
"process_owner": "process.owner",
"process_executable_path": "process.executable.path",
"process_executable_name": "process.executable.name",
"process_command_line": "process.command_line",
"process_command": "process.command",
"process_memory_usage": "process.memory.usage",
"process_memory_virtual": "process.memory.virtual",
"process_cpu_time": "process.cpu.time",
"process_disk_io": "process.disk.io",
"nfs_client_net_count": "nfs.client.net.count",
"nfs_client_net_tcp_connection_accepted": "nfs.client.net.tcp.connection.accepted",
"nfs_client_operation_count": "nfs.client.operation.count",
"nfs_client_procedure_count": "nfs.client.procedure.count",
"nfs_client_rpc_authrefresh_count": "nfs.client.rpc.authrefresh.count",
"nfs_client_rpc_count": "nfs.client.rpc.count",
"nfs_client_rpc_retransmit_count": "nfs.client.rpc.retransmit.count",
"nfs_server_fh_stale_count": "nfs.server.fh.stale.count",
"nfs_server_io": "nfs.server.io",
"nfs_server_net_count": "nfs.server.net.count",
"nfs_server_net_tcp_connection_accepted": "nfs.server.net.tcp.connection.accepted",
"nfs_server_operation_count": "nfs.server.operation.count",
"nfs_server_procedure_count": "nfs.server.procedure.count",
"nfs_server_repcache_requests": "nfs.server.repcache.requests",
"nfs_server_rpc_count": "nfs.server.rpc.count",
"nfs_server_thread_count": "nfs.server.thread.count",
"k8s_persistentvolumeclaim_name": "k8s.persistentvolumeclaim.name",
"k8s_volume_available": "k8s.volume.available",
"k8s_volume_capacity": "k8s.volume.capacity",
"k8s_volume_inodes": "k8s.volume.inodes",
"k8s_volume_inodes_free": "k8s.volume.inodes.free",
"k8s_pod_uid": "k8s.pod.uid",
"k8s_pod_name": "k8s.pod.name",
"k8s_container_name": "k8s.container.name",
"container_id": "container.id",
"k8s_volume_name": "k8s.volume.name",
"k8s_volume_type": "k8s.volume.type",
"aws_volume_id": "aws.volume.id",
"fs_type": "fs.type",
"partition": "partition",
"gce_pd_name": "gce.pd.name",
"glusterfs_endpoints_name": "glusterfs.endpoints.name",
"glusterfs_path": "glusterfs.path",
"interface": "interface",
"direction": "direction",
"k8s_node_cpu_usage": "k8s.node.cpu.usage",
"k8s_node_cpu_time": "k8s.node.cpu.time",
"k8s_node_memory_available": "k8s.node.memory.available",
"k8s_node_memory_rss": "k8s.node.memory.rss",
"k8s_node_memory_working_set": "k8s.node.memory.working_set",
"k8s_node_memory_page_faults": "k8s.node.memory.page_faults",
"k8s_node_memory_major_page_faults": "k8s.node.memory.major_page_faults",
"k8s_node_filesystem_available": "k8s.node.filesystem.available",
"k8s_node_filesystem_capacity": "k8s.node.filesystem.capacity",
"k8s_node_filesystem_usage": "k8s.node.filesystem.usage",
"k8s_node_network_io": "k8s.node.network.io",
"k8s_node_network_errors": "k8s.node.network.errors",
"k8s_node_uptime": "k8s.node.uptime",
"k8s_pod_cpu_usage": "k8s.pod.cpu.usage",
"k8s_pod_cpu_time": "k8s.pod.cpu.time",
"k8s_pod_memory_available": "k8s.pod.memory.available",
"k8s_pod_cpu_node_utilization": "k8s.pod.cpu.node.utilization",
"k8s_pod_memory_node_utilization": "k8s.pod.memory.node.utilization",
"k8s_pod_memory_rss": "k8s.pod.memory.rss",
"k8s_pod_memory_working_set": "k8s.pod.memory.working_set",
"k8s_pod_memory_page_faults": "k8s.pod.memory.page_faults",
"k8s_pod_memory_major_page_faults": "k8s.pod.memory.major_page_faults",
"k8s_pod_filesystem_available": "k8s.pod.filesystem.available",
"k8s_pod_filesystem_capacity": "k8s.pod.filesystem.capacity",
"k8s_pod_filesystem_usage": "k8s.pod.filesystem.usage",
"k8s_pod_network_io": "k8s.pod.network.io",
"k8s_pod_network_errors": "k8s.pod.network.errors",
"k8s_pod_uptime": "k8s.pod.uptime",
"container_cpu_usage": "container.cpu.usage",
"container_cpu_time": "container.cpu.time",
"container_memory_available": "container.memory.available",
"container_memory_usage": "container.memory.usage",
"k8s_container_cpu_node_utilization": "k8s.container.cpu.node.utilization",
"k8s_container_cpu_limit_utilization": "k8s.container.cpu_limit_utilization",
"k8s_container_cpu_request_utilization": "k8s.container.cpu_request_utilization",
"k8s_container_memory_node_utilization": "k8s.container.memory.node.utilization",
"k8s_container_memory_limit_utilization": "k8s.container.memory_limit_utilization",
"k8s_container_memory_request_utilization": "k8s.container.memory_request_utilization",
"container_memory_rss": "container.memory.rss",
"container_memory_working_set": "container.memory.working_set",
"container_memory_page_faults": "container.memory.page_faults",
"container_memory_major_page_faults": "container.memory.major_page_faults",
"container_filesystem_available": "container.filesystem.available",
"container_filesystem_capacity": "container.filesystem.capacity",
"container_filesystem_usage": "container.filesystem.usage",
"container_uptime": "container.uptime",
"k8s_volume_inodes_used": "k8s.volume.inodes.used",
"k8s_namespace_uid": "k8s.namespace.uid",
"container_image_name": "container.image.name",
"container_image_tag": "container.image.tag",
"k8s_pod_qos_class": "k8s.pod.qos_class",
"k8s_replicaset_name": "k8s.replicaset.name",
"k8s_replicaset_uid": "k8s.replicaset.uid",
"k8s_replicationcontroller_name": "k8s.replicationcontroller.name",
"k8s_replicationcontroller_uid": "k8s.replicationcontroller.uid",
"k8s_resourcequota_uid": "k8s.resourcequota.uid",
"k8s_resourcequota_name": "k8s.resourcequota.name",
"k8s_statefulset_uid": "k8s.statefulset.uid",
"k8s_statefulset_name": "k8s.statefulset.name",
"k8s_deployment_uid": "k8s.deployment.uid",
"k8s_cronjob_uid": "k8s.cronjob.uid",
"k8s_daemonset_uid": "k8s.daemonset.uid",
"k8s_hpa_uid": "k8s.hpa.uid",
"k8s_hpa_name": "k8s.hpa.name",
"k8s_hpa_scaletargetref_kind": "k8s.hpa.scaletargetref.kind",
"k8s_hpa_scaletargetref_name": "k8s.hpa.scaletargetref.name",
"k8s_hpa_scaletargetref_apiversion": "k8s.hpa.scaletargetref.apiversion",
"k8s_job_uid": "k8s.job.uid",
"k8s_kubelet_version": "k8s.kubelet.version",
"container_runtime": "container.runtime",
"container_runtime_version": "container.runtime.version",
"os_description": "os.description",
"openshift_clusterquota_uid": "openshift.clusterquota.uid",
"openshift_clusterquota_name": "openshift.clusterquota.name",
"k8s_container_status_last_terminated_reason": "k8s.container.status.last_terminated_reason",
"resource": "resource",
"condition": "condition",
"k8s_container_cpu_request": "k8s.container.cpu_request",
"k8s_container_cpu_limit": "k8s.container.cpu_limit",
"k8s_container_memory_request": "k8s.container.memory_request",
"k8s_container_memory_limit": "k8s.container.memory_limit",
"k8s_container_storage_request": "k8s.container.storage_request",
"k8s_container_storage_limit": "k8s.container.storage_limit",
"k8s_container_ephemeralstorage_request": "k8s.container.ephemeralstorage_request",
"k8s_container_ephemeralstorage_limit": "k8s.container.ephemeralstorage_limit",
"k8s_container_ready": "k8s.container.ready",
"k8s_pod_status_reason": "k8s.pod.status_reason",
"k8s_cronjob_active_jobs": "k8s.cronjob.active_jobs",
"k8s_daemonset_misscheduled_nodes": "k8s.daemonset.misscheduled_nodes",
"k8s_daemonset_ready_nodes": "k8s.daemonset.ready_nodes",
"k8s_hpa_max_replicas": "k8s.hpa.max_replicas",
"k8s_hpa_min_replicas": "k8s.hpa.min_replicas",
"k8s_hpa_current_replicas": "k8s.hpa.current_replicas",
"k8s_hpa_desired_replicas": "k8s.hpa.desired_replicas",
"k8s_job_max_parallel_pods": "k8s.job.max_parallel_pods",
"k8s_namespace_phase": "k8s.namespace.phase",
"k8s_replicaset_desired": "k8s.replicaset.desired",
"k8s_replicaset_available": "k8s.replicaset.available",
"k8s_replication_controller_desired": "k8s.replication_controller.desired",
"k8s_replication_controller_available": "k8s.replication_controller.available",
"k8s_resource_quota_hard_limit": "k8s.resource_quota.hard_limit",
"k8s_resource_quota_used": "k8s.resource_quota.used",
"k8s_statefulset_updated_pods": "k8s.statefulset.updated_pods",
"k8s_node_condition": "k8s.node.condition",
}
const fromWhereQuery = `
FROM %s.%s
WHERE metric_name IN (%s)
@@ -18,39 +262,39 @@ WHERE metric_name IN (%s)
var (
// TODO(srikanthccv): import metadata yaml from receivers and use generated files to check the metrics
podMetricNamesToCheck = []string{
"k8s.pod.cpu.usage",
"k8s.pod.memory.working_set",
"k8s.pod.cpu_request_utilization",
"k8s.pod.memory_request_utilization",
"k8s.pod.cpu_limit_utilization",
"k8s.pod.memory_limit_utilization",
"k8s.container.restarts",
"k8s.pod.phase",
GetDotMetrics("k8s_pod_cpu_usage"),
GetDotMetrics("k8s_pod_memory_working_set"),
GetDotMetrics("k8s_pod_cpu_request_utilization"),
GetDotMetrics("k8s_pod_memory_request_utilization"),
GetDotMetrics("k8s_pod_cpu_limit_utilization"),
GetDotMetrics("k8s_pod_memory_limit_utilization"),
GetDotMetrics("k8s_container_restarts"),
GetDotMetrics("k8s_pod_phase"),
}
nodeMetricNamesToCheck = []string{
"k8s.node.cpu.usage",
"k8s.node.allocatable_cpu",
"k8s.node.memory.working_set",
"k8s.node.allocatable_memory",
"k8s.node.condition_ready",
GetDotMetrics("k8s_node_cpu_usage"),
GetDotMetrics("k8s_node_allocatable_cpu"),
GetDotMetrics("k8s_node_memory_working_set"),
GetDotMetrics("k8s_node_allocatable_memory"),
GetDotMetrics("k8s_node_condition_ready"),
}
clusterMetricNamesToCheck = []string{
"k8s.daemonset.desired_scheduled_nodes",
"k8s.daemonset.current_scheduled_nodes",
"k8s.deployment.desired",
"k8s.deployment.available",
"k8s.job.desired_successful_pods",
"k8s.job.active_pods",
"k8s.job.failed_pods",
"k8s.job.successful_pods",
"k8s.statefulset.desired_pods",
"k8s.statefulset.current_pods",
GetDotMetrics("k8s_daemonset_desired_scheduled_nodes"),
GetDotMetrics("k8s_daemonset_current_scheduled_nodes"),
GetDotMetrics("k8s_deployment_desired"),
GetDotMetrics("k8s_deployment_available"),
GetDotMetrics("k8s_job_desired_successful_pods"),
GetDotMetrics("k8s_job_active_pods"),
GetDotMetrics("k8s_job_failed_pods"),
GetDotMetrics("k8s_job_successful_pods"),
GetDotMetrics("k8s_statefulset_desired_pods"),
GetDotMetrics("k8s_statefulset_current_pods"),
}
optionalPodMetricNamesToCheck = []string{
"k8s.pod.cpu_request_utilization",
"k8s.pod.memory_request_utilization",
"k8s.pod.cpu_limit_utilization",
"k8s.pod.memory_limit_utilization",
GetDotMetrics("k8s_pod_cpu_request_utilization"),
GetDotMetrics("k8s_pod_memory_request_utilization"),
GetDotMetrics("k8s_pod_cpu_limit_utilization"),
GetDotMetrics("k8s_pod_memory_limit_utilization"),
}
// did they ever send _any_ pod metrics?
@@ -88,15 +332,15 @@ SELECT
any(JSONExtractString(labels, '%s')) as k8s_job_name,
JSONExtractString(labels, '%s') as k8s_pod_name
`,
"k8s.cluster.name",
"k8s.node.name",
"k8s.namespace.name",
"k8s.deployment.name",
"k8s.statefulset.name",
"k8s.daemonset.name",
"k8s.cronjob.name",
"k8s.job.name",
"k8s.pod.name",
GetDotMetrics("k8s_cluster_name"),
GetDotMetrics("k8s_node_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_deployment_name"),
GetDotMetrics("k8s_statefulset_name"),
GetDotMetrics("k8s_daemonset_name"),
GetDotMetrics("k8s_cronjob_name"),
GetDotMetrics("k8s_job_name"),
GetDotMetrics("k8s_pod_name"),
)
filterGroupQuery = fmt.Sprintf(`
@@ -105,7 +349,7 @@ AND JSONExtractString(labels, '%s')
GROUP BY k8s_pod_name
LIMIT 1 BY k8s_cluster_name, k8s_node_name, k8s_namespace_name
`,
"k8s.namespace.name",
GetDotMetrics("k8s_namespace_name"),
)
isSendingRequiredMetadataQuery = selectQuery + fromWhereQuery + filterGroupQuery
@@ -206,3 +450,12 @@ func getParamsForTopVolumes(req model.VolumeListRequest) (int64, string, string)
func localQueryToDistributedQuery(query string) string {
return strings.Replace(query, ".time_series_v4", ".distributed_time_series_v4", 1)
}
func GetDotMetrics(key string) string {
if constants.IsDotMetricsEnabled {
if _, ok := dotMetricMap[key]; ok {
return dotMetricMap[key]
}
}
return key
}

View File

@@ -18,18 +18,18 @@ import (
)
var (
metricToUseForDaemonSets = "k8s.pod.cpu.usage"
k8sDaemonSetNameAttrKey = "k8s.daemonset.name"
metricToUseForDaemonSets = GetDotMetrics("k8s_pod_cpu_usage")
k8sDaemonSetNameAttrKey = GetDotMetrics("k8s_daemonset_name")
metricNamesForDaemonSets = map[string]string{
"desired_nodes": "k8s.daemonset.desired_scheduled_nodes",
"available_nodes": "k8s.daemonset.current_scheduled_nodes",
"desired_nodes": GetDotMetrics("k8s_daemonset_desired_scheduled_nodes"),
"available_nodes": GetDotMetrics("k8s_daemonset_current_scheduled_nodes"),
}
daemonSetAttrsToEnrich = []string{
"k8s.daemonset.name",
"k8s.namespace.name",
"k8s.cluster.name",
GetDotMetrics("k8s_daemonset_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_cluster_name"),
}
queryNamesForDaemonSets = map[string][]string{

View File

@@ -18,18 +18,18 @@ import (
)
var (
metricToUseForDeployments = "k8s.pod.cpu.usage"
k8sDeploymentNameAttrKey = "k8s.deployment.name"
metricToUseForDeployments = GetDotMetrics("k8s_pod_cpu_usage")
k8sDeploymentNameAttrKey = GetDotMetrics("k8s_deployment_name")
metricNamesForDeployments = map[string]string{
"desired_pods": "k8s.deployment.desired",
"available_pods": "k8s.deployment.available",
"desired_pods": GetDotMetrics("k8s_deployment_desired"),
"available_pods": GetDotMetrics("k8s_deployment_available"),
}
deploymentAttrsToEnrich = []string{
"k8s.deployment.name",
"k8s.namespace.name",
"k8s.cluster.name",
GetDotMetrics("k8s_deployment_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_cluster_name"),
}
queryNamesForDeployments = map[string][]string{

View File

@@ -45,15 +45,15 @@ var (
"mode",
"mountpoint",
"type",
"os.type",
"process.cgroup",
"process.command",
"process.command_line",
"process.executable.name",
"process.executable.path",
"process.owner",
"process.parent_pid",
"process.pid",
GetDotMetrics("os_type"),
GetDotMetrics("process_cgroup"),
GetDotMetrics("process_command"),
GetDotMetrics("process_command_line"),
GetDotMetrics("process_executable_name"),
GetDotMetrics("process_executable_path"),
GetDotMetrics("process_owner"),
GetDotMetrics("process_parent_pid"),
GetDotMetrics("process_pid"),
}
queryNamesForTopHosts = map[string][]string{
@@ -64,65 +64,65 @@ var (
}
// TODO(srikanthccv): remove hardcoded metric name and support keys from any system metric
metricToUseForHostAttributes = "system.cpu.load_average.15m"
hostNameAttrKey = "host.name"
metricToUseForHostAttributes = GetDotMetrics("system_cpu_load_average_15m")
hostNameAttrKey = GetDotMetrics("host_name")
agentNameToIgnore = "k8s-infra-otel-agent"
hostAttrsToEnrich = []string{
"os.type",
GetDotMetrics("os_type"),
}
metricNamesForHosts = map[string]string{
"filesystem": "system.filesystem.usage",
"cpu": "system.cpu.time",
"memory": "system.memory.usage",
"load15": "system.cpu.load_average.15m",
"wait": "system.cpu.time",
"filesystem": GetDotMetrics("system_filesystem_usage"),
"cpu": GetDotMetrics("system_cpu_time"),
"memory": GetDotMetrics("system_memory_usage"),
"load15": GetDotMetrics("system_cpu_load_average_15m"),
"wait": GetDotMetrics("system_cpu_time"),
}
uniqueMetricNamesForHosts = []string{
"system.uptime",
"system.cpu.time",
"system.cpu.load_average.1m",
"system.cpu.load_average.5m",
"system.cpu.load_average.15m",
"system.memory.usage",
"system.paging.usage",
"system.paging.faults",
"system.paging.operations",
"system.disk.io",
"system.disk.operations",
"system.disk.io_time",
"system.disk.operation_time",
"system.disk.merged",
"system.disk.pending_operations",
"system.disk.weighted_io_time",
"system.filesystem.usage",
"system.filesystem.inodes.usage",
"system.network.io",
"system.network.errors",
"system.network.connections",
"system.network.dropped",
"system.network.packets",
"system.processes.count",
"system.processes.created",
"process.cpu.time",
"process.disk.io",
"process.memory.usage",
"process.memory.virtual",
"nfs.client.net.count",
"nfs.client.net.tcp.connection.accepted",
"nfs.client.operation.count",
"nfs.client.procedure.count",
"nfs.client.rpc.authrefresh.count",
"nfs.client.rpc.count",
"nfs.client.rpc.retransmit.count",
"nfs.server.fh.stale.count",
"nfs.server.io",
"nfs.server.net.count",
"nfs.server.net.tcp.connection.accepted",
"nfs.server.operation.count",
"nfs.server.procedure.count",
"nfs.server.repcache.requests",
"nfs.server.rpc.count",
"nfs.server.thread.count",
GetDotMetrics("system_uptime"),
GetDotMetrics("system_cpu_time"),
GetDotMetrics("system_cpu_load_average_1m"),
GetDotMetrics("system_cpu_load_average_5m"),
GetDotMetrics("system_cpu_load_average_15m"),
GetDotMetrics("system_memory_usage"),
GetDotMetrics("system_paging_usage"),
GetDotMetrics("system_paging_faults"),
GetDotMetrics("system_paging_operations"),
GetDotMetrics("system_disk_io"),
GetDotMetrics("system_disk_operations"),
GetDotMetrics("system_disk_io_time"),
GetDotMetrics("system_disk_operation_time"),
GetDotMetrics("system_disk_merged"),
GetDotMetrics("system_disk_pending_operations"),
GetDotMetrics("system_disk_weighted_io_time"),
GetDotMetrics("system_filesystem_usage"),
GetDotMetrics("system_filesystem_inodes_usage"),
GetDotMetrics("system_network_io"),
GetDotMetrics("system_network_errors"),
GetDotMetrics("system_network_connections"),
GetDotMetrics("system_network_dropped"),
GetDotMetrics("system_network_packets"),
GetDotMetrics("system_processes_count"),
GetDotMetrics("system_processes_created"),
GetDotMetrics("process_cpu_time"),
GetDotMetrics("process_disk_io"),
GetDotMetrics("process_memory_usage"),
GetDotMetrics("process_memory_virtual"),
GetDotMetrics("nfs_client_net_count"),
GetDotMetrics("nfs_client_net_tcp_connection_accepted"),
GetDotMetrics("nfs_client_operation_count"),
GetDotMetrics("nfs_client_procedure_count"),
GetDotMetrics("nfs_client_rpc_authrefresh_count"),
GetDotMetrics("nfs_client_rpc_count"),
GetDotMetrics("nfs_client_rpc_retransmit_count"),
GetDotMetrics("nfs_server_fh_stale_count"),
GetDotMetrics("nfs_server_io"),
GetDotMetrics("nfs_server_net_count"),
GetDotMetrics("nfs_server_net_tcp_connection_accepted"),
GetDotMetrics("nfs_server_operation_count"),
GetDotMetrics("nfs_server_procedure_count"),
GetDotMetrics("nfs_server_repcache_requests"),
GetDotMetrics("nfs_server_rpc_count"),
GetDotMetrics("nfs_server_thread_count"),
}
)
@@ -351,8 +351,8 @@ func (h *HostsRepo) IsSendingK8SAgentMetrics(ctx context.Context, req model.Host
AND unix_milli >= toUnixTimestamp(now() - INTERVAL 60 MINUTE) * 1000
AND JSONExtractString(labels, '%s') LIKE '%%-otel-agent%%'
AND fingerprint GLOBAL IN (%s)`,
"k8s.cluster.name", "k8s.node.name",
constants.SIGNOZ_METRIC_DBNAME, constants.SIGNOZ_TIMESERIES_V4_TABLENAME, namesStr, "host.name", queryForRecentFingerprints)
GetDotMetrics("k8s_cluster_name"), GetDotMetrics("k8s_node_name"),
constants.SIGNOZ_METRIC_DBNAME, constants.SIGNOZ_TIMESERIES_V4_TABLENAME, namesStr, GetDotMetrics("host_name"), queryForRecentFingerprints)
result, err := h.reader.GetListResultV3(ctx, query)
if err != nil {
@@ -363,13 +363,13 @@ func (h *HostsRepo) IsSendingK8SAgentMetrics(ctx context.Context, req model.Host
nodeNames := make(map[string]struct{})
for _, row := range result {
switch v := row.Data["k8s.cluster.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_cluster_name")].(type) {
case string:
clusterNames[v] = struct{}{}
case *string:
clusterNames[*v] = struct{}{}
}
switch v := row.Data["k8s.node.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_node_name")].(type) {
case string:
nodeNames[v] = struct{}{}
case *string:
@@ -535,7 +535,7 @@ func (h *HostsRepo) GetHostList(ctx context.Context, orgID valuer.UUID, req mode
if _, ok := hostAttrs[record.HostName]; ok {
record.Meta = hostAttrs[record.HostName]
}
if osType, ok := record.Meta["os.type"]; ok {
if osType, ok := record.Meta[GetDotMetrics("os_type")]; ok {
record.OS = osType
}
record.Active = activeHosts[record.HostName]

View File

@@ -18,20 +18,20 @@ import (
)
var (
metricToUseForJobs = "k8s.job.desired_successful_pods"
k8sJobNameAttrKey = "k8s.job.name"
metricToUseForJobs = GetDotMetrics("k8s_job_desired_successful_pods")
k8sJobNameAttrKey = GetDotMetrics("k8s_job_name")
metricNamesForJobs = map[string]string{
"desired_successful_pods": "k8s.job.desired_successful_pods",
"active_pods": "k8s.job.active_pods",
"failed_pods": "k8s.job.failed_pods",
"successful_pods": "k8s.job.successful_pods",
"desired_successful_pods": GetDotMetrics("k8s_job_desired_successful_pods"),
"active_pods": GetDotMetrics("k8s_job_active_pods"),
"failed_pods": GetDotMetrics("k8s_job_failed_pods"),
"successful_pods": GetDotMetrics("k8s_job_successful_pods"),
}
jobAttrsToEnrich = []string{
"k8s.job.name",
"k8s.namespace.name",
"k8s.cluster.name",
GetDotMetrics("k8s_job_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_cluster_name"),
}
queryNamesForJobs = map[string][]string{
@@ -54,7 +54,7 @@ var (
QueryName: "H",
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: metricNamesForJobs["desired_successful_pods"],
Key: GetDotMetrics(metricNamesForJobs["desired_successful_pods"]),
DataType: v3.AttributeKeyDataTypeFloat64,
},
Temporality: v3.Unspecified,
@@ -74,7 +74,7 @@ var (
QueryName: "I",
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: metricNamesForJobs["active_pods"],
Key: GetDotMetrics(metricNamesForJobs["active_pods"]),
DataType: v3.AttributeKeyDataTypeFloat64,
},
Temporality: v3.Unspecified,
@@ -94,7 +94,7 @@ var (
QueryName: "J",
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: metricNamesForJobs["failed_pods"],
Key: GetDotMetrics(metricNamesForJobs["failed_pods"]),
DataType: v3.AttributeKeyDataTypeFloat64,
},
Temporality: v3.Unspecified,
@@ -114,7 +114,7 @@ var (
QueryName: "K",
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: metricNamesForJobs["successful_pods"],
Key: GetDotMetrics(metricNamesForJobs["successful_pods"]),
DataType: v3.AttributeKeyDataTypeFloat64,
},
Temporality: v3.Unspecified,
@@ -327,7 +327,7 @@ func (d *JobsRepo) GetJobList(ctx context.Context, orgID valuer.UUID, req model.
}
if req.OrderBy == nil {
req.OrderBy = &v3.OrderBy{ColumnName: "desired_pods", Order: v3.DirectionDesc}
req.OrderBy = &v3.OrderBy{ColumnName: GetDotMetrics("desired_pods"), Order: v3.DirectionDesc}
}
if req.GroupBy == nil {

View File

@@ -18,11 +18,11 @@ import (
)
var (
metricToUseForNamespaces = "k8s.pod.cpu.usage"
metricToUseForNamespaces = GetDotMetrics("k8s_pod_cpu_usage")
namespaceAttrsToEnrich = []string{
"k8s.namespace.name",
"k8s.cluster.name",
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_cluster_name"),
}
queryNamesForNamespaces = map[string][]string{
@@ -33,11 +33,11 @@ var (
namespaceQueryNames = []string{"A", "D", "H", "I", "J", "K"}
attributesKeysForNamespaces = []v3.AttributeKey{
{Key: "k8s.namespace.name"},
{Key: "k8s.cluster.name"},
{Key: GetDotMetrics("k8s_namespace_name")},
{Key: GetDotMetrics("k8s_cluster_name")},
}
k8sNamespaceNameAttrKey = "k8s.namespace.name"
k8sNamespaceNameAttrKey = GetDotMetrics("k8s_namespace_name")
)
type NamespacesRepo struct {

View File

@@ -21,11 +21,11 @@ import (
)
var (
metricToUseForNodes = "k8s.node.cpu.usage"
metricToUseForNodes = GetDotMetrics("k8s_node_cpu_usage")
nodeAttrsToEnrich = []string{"k8s.node.name", "k8s.node.uid", "k8s.cluster.name"}
nodeAttrsToEnrich = []string{GetDotMetrics("k8s_node_name"), GetDotMetrics("k8s_node_uid"), GetDotMetrics("k8s_cluster_name")}
k8sNodeGroupAttrKey = "k8s.node.name"
k8sNodeGroupAttrKey = GetDotMetrics("k8s_node_name")
queryNamesForNodes = map[string][]string{
"cpu": {"A"},
@@ -36,11 +36,11 @@ var (
nodeQueryNames = []string{"A", "B", "C", "D", "E", "F"}
metricNamesForNodes = map[string]string{
"cpu": "k8s.node.cpu.usage",
"cpu_allocatable": "k8s.node.allocatable_cpu",
"memory": "k8s.node.memory.working_set",
"memory_allocatable": "k8s.node.allocatable_memory",
"node_condition": "k8s.node.condition_ready",
"cpu": GetDotMetrics("k8s_node_cpu_usage"),
"cpu_allocatable": GetDotMetrics("k8s_node_allocatable_cpu"),
"memory": GetDotMetrics("k8s_node_memory_working_set"),
"memory_allocatable": GetDotMetrics("k8s_node_allocatable_memory"),
"node_condition": GetDotMetrics("k8s_node_condition_ready"),
}
)

View File

@@ -21,22 +21,22 @@ import (
)
var (
metricToUseForPods = "k8s.pod.cpu.usage"
metricToUseForPods = GetDotMetrics("k8s_pod_cpu_usage")
podAttrsToEnrich = []string{
"k8s.pod.uid",
"k8s.pod.name",
"k8s.namespace.name",
"k8s.node.name",
"k8s.deployment.name",
"k8s.statefulset.name",
"k8s.daemonset.name",
"k8s.job.name",
"k8s.cronjob.name",
"k8s.cluster.name",
GetDotMetrics("k8s_pod_uid"),
GetDotMetrics("k8s_pod_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_node_name"),
GetDotMetrics("k8s_deployment_name"),
GetDotMetrics("k8s_statefulset_name"),
GetDotMetrics("k8s_daemonset_name"),
GetDotMetrics("k8s_job_name"),
GetDotMetrics("k8s_cronjob_name"),
GetDotMetrics("k8s_cluster_name"),
}
k8sPodUIDAttrKey = "k8s.pod.uid"
k8sPodUIDAttrKey = GetDotMetrics("k8s_pod_uid")
queryNamesForPods = map[string][]string{
"cpu": {"A"},
@@ -51,14 +51,14 @@ var (
podQueryNames = []string{"A", "B", "C", "D", "E", "F", "G", "H", "I", "J", "K"}
metricNamesForPods = map[string]string{
"cpu": "k8s.pod.cpu.usage",
"cpu_request": "k8s.pod.cpu_request_utilization",
"cpu_limit": "k8s.pod.cpu_limit_utilization",
"memory": "k8s.pod.memory.working_set",
"memory_request": "k8s.pod.memory_request_utilization",
"memory_limit": "k8s.pod.memory_limit_utilization",
"restarts": "k8s.container.restarts",
"pod_phase": "k8s.pod.phase",
"cpu": GetDotMetrics("k8s_pod_cpu_usage"),
"cpu_request": GetDotMetrics("k8s_pod_cpu_request_utilization"),
"cpu_limit": GetDotMetrics("k8s_pod_cpu_limit_utilization"),
"memory": GetDotMetrics("k8s_pod_memory_working_set"),
"memory_request": GetDotMetrics("k8s_pod_memory_request_utilization"),
"memory_limit": GetDotMetrics("k8s_pod_memory_limit_utilization"),
"restarts": GetDotMetrics("k8s_container_restarts"),
"pod_phase": GetDotMetrics("k8s_pod_phase"),
}
)
@@ -169,7 +169,7 @@ func (p *PodsRepo) SendingRequiredMetadata(ctx context.Context) ([]model.PodOnbo
// for each pod, check if we have all the required metadata
for _, row := range result {
status := model.PodOnboardingStatus{}
switch v := row.Data["k8s.cluster.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_cluster_name")].(type) {
case string:
status.HasClusterName = true
status.ClusterName = v
@@ -177,7 +177,7 @@ func (p *PodsRepo) SendingRequiredMetadata(ctx context.Context) ([]model.PodOnbo
status.HasClusterName = *v != ""
status.ClusterName = *v
}
switch v := row.Data["k8s.node.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_node_name")].(type) {
case string:
status.HasNodeName = true
status.NodeName = v
@@ -185,7 +185,7 @@ func (p *PodsRepo) SendingRequiredMetadata(ctx context.Context) ([]model.PodOnbo
status.HasNodeName = *v != ""
status.NodeName = *v
}
switch v := row.Data["k8s.namespace.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_namespace_name")].(type) {
case string:
status.HasNamespaceName = true
status.NamespaceName = v
@@ -193,38 +193,38 @@ func (p *PodsRepo) SendingRequiredMetadata(ctx context.Context) ([]model.PodOnbo
status.HasNamespaceName = *v != ""
status.NamespaceName = *v
}
switch v := row.Data["k8s.deployment.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_deployment_name")].(type) {
case string:
status.HasDeploymentName = true
case *string:
status.HasDeploymentName = *v != ""
}
switch v := row.Data["k8s.statefulset.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_statefulset_name")].(type) {
case string:
status.HasStatefulsetName = true
case *string:
status.HasStatefulsetName = *v != ""
}
switch v := row.Data["k8s.daemonset.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_daemonset_name")].(type) {
case string:
status.HasDaemonsetName = true
case *string:
status.HasDaemonsetName = *v != ""
}
switch v := row.Data["k8s.cronjob.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_cronjob_name")].(type) {
case string:
status.HasCronjobName = true
case *string:
status.HasCronjobName = *v != ""
}
switch v := row.Data["k8s.job.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_job_name")].(type) {
case string:
status.HasJobName = true
case *string:
status.HasJobName = *v != ""
}
switch v := row.Data["k8s.pod.name"].(type) {
switch v := row.Data[GetDotMetrics("k8s_pod_name")].(type) {
case string:
status.PodName = v
case *string:

View File

@@ -23,15 +23,15 @@ var (
"memory": {"C"},
}
processPIDAttrKey = "process.pid"
processPIDAttrKey = GetDotMetrics("process_pid")
metricNamesForProcesses = map[string]string{
"cpu": "process.cpu.time",
"memory": "process.memory.usage",
"cpu": GetDotMetrics("process_cpu_time"),
"memory": GetDotMetrics("process_memory_usage"),
}
metricToUseForProcessAttributes = "process.memory.usage"
processNameAttrKey = "process.executable.name"
processCMDAttrKey = "process.command"
processCMDLineAttrKey = "process.command_line"
metricToUseForProcessAttributes = GetDotMetrics("process_memory_usage")
processNameAttrKey = GetDotMetrics("process_executable_name")
processCMDAttrKey = GetDotMetrics("process_command")
processCMDLineAttrKey = GetDotMetrics("process_command_line")
)
type ProcessesRepo struct {
@@ -46,7 +46,7 @@ func NewProcessesRepo(reader interfaces.Reader, querierV2 interfaces.Querier) *P
func (p *ProcessesRepo) GetProcessAttributeKeys(ctx context.Context, orgID valuer.UUID, req v3.FilterAttributeKeyRequest) (*v3.FilterAttributeKeyResponse, error) {
// TODO(srikanthccv): remove hardcoded metric name and support keys from any system metric
req.DataSource = v3.DataSourceMetrics
req.AggregateAttribute = "process.memory.usage"
req.AggregateAttribute = GetDotMetrics("process_memory_usage")
if req.Limit == 0 {
req.Limit = 50
}
@@ -71,7 +71,7 @@ func (p *ProcessesRepo) GetProcessAttributeKeys(ctx context.Context, orgID value
func (p *ProcessesRepo) GetProcessAttributeValues(ctx context.Context, orgID valuer.UUID, req v3.FilterAttributeValueRequest) (*v3.FilterAttributeValueResponse, error) {
req.DataSource = v3.DataSourceMetrics
req.AggregateAttribute = "process.memory.usage"
req.AggregateAttribute = GetDotMetrics("process_memory_usage")
if req.Limit == 0 {
req.Limit = 50
}
@@ -87,7 +87,7 @@ func (p *ProcessesRepo) getMetadataAttributes(ctx context.Context,
req model.ProcessListRequest) (map[string]map[string]string, error) {
processAttrs := map[string]map[string]string{}
keysToAdd := []string{"process.pid", "process.executable.name", "process.command", "process.command_line"}
keysToAdd := []string{GetDotMetrics("process_pid"), GetDotMetrics("process_executable_name"), GetDotMetrics("process_command"), GetDotMetrics("process_command_line")}
for _, key := range keysToAdd {
hasKey := false
for _, groupByKey := range req.GroupBy {

View File

@@ -18,19 +18,19 @@ import (
)
var (
metricToUseForVolumes = "k8s.volume.available"
metricToUseForVolumes = GetDotMetrics("k8s_volume_available")
volumeAttrsToEnrich = []string{
"k8s.pod.uid",
"k8s.pod.name",
"k8s.namespace.name",
"k8s.node.name",
"k8s.statefulset.name",
"k8s.cluster.name",
"k8s.persistentvolumeclaim.name",
GetDotMetrics("k8s_pod_uid"),
GetDotMetrics("k8s_pod_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_node_name"),
GetDotMetrics("k8s_statefulset_name"),
GetDotMetrics("k8s_cluster_name"),
GetDotMetrics("k8s_persistentvolumeclaim_name"),
}
k8sPersistentVolumeClaimNameAttrKey = "k8s.persistentvolumeclaim.name"
k8sPersistentVolumeClaimNameAttrKey = GetDotMetrics("k8s_persistentvolumeclaim_name")
queryNamesForVolumes = map[string][]string{
"available": {"A"},
@@ -44,11 +44,11 @@ var (
volumeQueryNames = []string{"A", "B", "C", "D", "E", "F1"}
metricNamesForVolumes = map[string]string{
"available": "k8s.volume.available",
"capacity": "k8s.volume.capacity",
"inodes": "k8s.volume.inodes",
"inodes_free": "k8s.volume.inodes.free",
"inodes_used": "k8s.volume.inodes.used",
"available": GetDotMetrics("k8s_volume_available"),
"capacity": GetDotMetrics("k8s_volume_capacity"),
"inodes": GetDotMetrics("k8s_volume_inodes"),
"inodes_free": GetDotMetrics("k8s_volume_inodes_free"),
"inodes_used": GetDotMetrics("k8s_volume_inodes_used"),
}
)

View File

@@ -18,18 +18,18 @@ import (
)
var (
metricToUseForStatefulSets = "k8s.pod.cpu.usage"
k8sStatefulSetNameAttrKey = "k8s.statefulset.name"
metricToUseForStatefulSets = GetDotMetrics("k8s_pod_cpu_usage")
k8sStatefulSetNameAttrKey = GetDotMetrics("k8s_statefulset_name")
metricNamesForStatefulSets = map[string]string{
"desired_pods": "k8s.statefulset.desired_pods",
"available_pods": "k8s.statefulset.current_pods",
"desired_pods": GetDotMetrics("k8s_statefulset_desired_pods"),
"available_pods": GetDotMetrics("k8s_statefulset_current_pods"),
}
statefulSetAttrsToEnrich = []string{
"k8s.statefulset.name",
"k8s.namespace.name",
"k8s.cluster.name",
GetDotMetrics("k8s_statefulset_name"),
GetDotMetrics("k8s_namespace_name"),
GetDotMetrics("k8s_cluster_name"),
}
queryNamesForStatefulSets = map[string][]string{

View File

@@ -4,13 +4,13 @@ import v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
var (
metricNamesForWorkloads = map[string]string{
"cpu": "k8s.pod.cpu.usage",
"cpu_request": "k8s.pod.cpu_request_utilization",
"cpu_limit": "k8s.pod.cpu_limit_utilization",
"memory": "k8s.pod.memory.working_set",
"memory_request": "k8s.pod.memory_request_utilization",
"memory_limit": "k8s.pod.memory_limit_utilization",
"restarts": "k8s.container.restarts",
"cpu": GetDotMetrics("k8s_pod_cpu_usage"),
"cpu_request": GetDotMetrics("k8s_pod_cpu_request_utilization"),
"cpu_limit": GetDotMetrics("k8s_pod_cpu_limit_utilization"),
"memory": GetDotMetrics("k8s_pod_memory_working_set"),
"memory_request": GetDotMetrics("k8s_pod_memory_request_utilization"),
"memory_limit": GetDotMetrics("k8s_pod_memory_limit_utilization"),
"restarts": GetDotMetrics("k8s_container_restarts"),
}
)

View File

@@ -4,6 +4,7 @@ import (
"fmt"
"github.com/SigNoz/signoz/pkg/query-service/common"
"github.com/SigNoz/signoz/pkg/query-service/constants"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
)
@@ -66,6 +67,11 @@ func buildBuilderQueriesProducerBytes(
attributeCache *Clients,
) (map[string]*v3.BuilderQuery, error) {
normalized := true
if constants.IsDotMetricsEnabled {
normalized = false
}
bq := make(map[string]*v3.BuilderQuery)
queryName := "byte_rate"
@@ -74,7 +80,7 @@ func buildBuilderQueriesProducerBytes(
StepInterval: common.MinAllowedStepInterval(unixMilliStart, unixMilliEnd),
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: "kafka.producer.byte-rate",
Key: getDotMetrics("kafka_producer_byte_rate", normalized),
DataType: v3.AttributeKeyDataTypeFloat64,
Type: v3.AttributeKeyType("Gauge"),
IsColumn: true,
@@ -88,7 +94,7 @@ func buildBuilderQueriesProducerBytes(
Items: []v3.FilterItem{
{
Key: v3.AttributeKey{
Key: "service.name",
Key: getDotMetrics("service_name", normalized),
Type: v3.AttributeKeyTypeTag,
DataType: v3.AttributeKeyDataTypeString,
},
@@ -110,7 +116,7 @@ func buildBuilderQueriesProducerBytes(
ReduceTo: v3.ReduceToOperatorAvg,
GroupBy: []v3.AttributeKey{
{
Key: "service.name",
Key: getDotMetrics("service_name", normalized),
DataType: v3.AttributeKeyDataTypeString,
Type: v3.AttributeKeyTypeTag,
},
@@ -133,12 +139,17 @@ func buildBuilderQueriesNetwork(
bq := make(map[string]*v3.BuilderQuery)
queryName := "latency"
normalized := true
if constants.IsDotMetricsEnabled {
normalized = false
}
chq := &v3.BuilderQuery{
QueryName: queryName,
StepInterval: common.MinAllowedStepInterval(unixMilliStart, unixMilliEnd),
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: "kafka.consumer.fetch_latency_avg",
Key: getDotMetrics("kafka_consumer_fetch_latency_avg", normalized),
},
AggregateOperator: v3.AggregateOperatorAvg,
Temporality: v3.Unspecified,
@@ -149,7 +160,7 @@ func buildBuilderQueriesNetwork(
Items: []v3.FilterItem{
{
Key: v3.AttributeKey{
Key: "service.name",
Key: getDotMetrics("service_name", normalized),
Type: v3.AttributeKeyTypeTag,
DataType: v3.AttributeKeyDataTypeString,
},
@@ -158,7 +169,7 @@ func buildBuilderQueriesNetwork(
},
{
Key: v3.AttributeKey{
Key: "client-id",
Key: getDotMetrics("client_id", normalized),
Type: v3.AttributeKeyTypeTag,
DataType: v3.AttributeKeyDataTypeString,
},
@@ -167,7 +178,7 @@ func buildBuilderQueriesNetwork(
},
{
Key: v3.AttributeKey{
Key: "service.instance.id",
Key: getDotMetrics("service_instance_id", normalized),
Type: v3.AttributeKeyTypeTag,
DataType: v3.AttributeKeyDataTypeString,
},
@@ -180,17 +191,17 @@ func buildBuilderQueriesNetwork(
ReduceTo: v3.ReduceToOperatorAvg,
GroupBy: []v3.AttributeKey{
{
Key: "service.name",
Key: getDotMetrics("service_name", normalized),
DataType: v3.AttributeKeyDataTypeString,
Type: v3.AttributeKeyTypeTag,
},
{
Key: "client-id",
Key: getDotMetrics("client_id", normalized),
DataType: v3.AttributeKeyDataTypeString,
Type: v3.AttributeKeyTypeTag,
},
{
Key: "service.instance.id",
Key: getDotMetrics("service_instance_id", normalized),
DataType: v3.AttributeKeyDataTypeString,
Type: v3.AttributeKeyTypeTag,
},
@@ -207,12 +218,17 @@ func BuildBuilderQueriesKafkaOnboarding(messagingQueue *MessagingQueue) (*v3.Que
unixMilliStart := messagingQueue.Start / 1000000
unixMilliEnd := messagingQueue.End / 1000000
normalized := true
if constants.IsDotMetricsEnabled {
normalized = false
}
buiderQuery := &v3.BuilderQuery{
QueryName: "fetch_latency",
StepInterval: common.MinAllowedStepInterval(unixMilliStart, unixMilliEnd),
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: "kafka.consumer.fetch_latency_avg",
Key: getDotMetrics("kafka_consumer_fetch_latency_avg", normalized),
},
AggregateOperator: v3.AggregateOperatorCount,
Temporality: v3.Unspecified,
@@ -227,7 +243,7 @@ func BuildBuilderQueriesKafkaOnboarding(messagingQueue *MessagingQueue) (*v3.Que
StepInterval: common.MinAllowedStepInterval(unixMilliStart, unixMilliEnd),
DataSource: v3.DataSourceMetrics,
AggregateAttribute: v3.AttributeKey{
Key: "kafka.consumer_group.lag",
Key: getDotMetrics("kafka_consumer_group_lag", normalized),
},
AggregateOperator: v3.AggregateOperatorCount,
Temporality: v3.Unspecified,
@@ -411,3 +427,19 @@ func buildCompositeQuery(chq *v3.ClickHouseQuery, queryContext string) (*v3.Comp
PanelType: v3.PanelTypeTable,
}, nil
}
func getDotMetrics(metricName string, normalized bool) string {
dotMetricsMap := map[string]string{
"kafka_producer_byte_rate": "kafka.producer.byte-rate",
"service_name": "service.name",
"kafka_consumer_fetch_latency_avg": "kafka.consumer.fetch_latency_avg",
"service_instance_id": "service.instance.id",
"client_id": "client-id",
"kafka_consumer_group_lag": "kafka.consumer_group.lag",
}
if _, ok := dotMetricsMap[metricName]; ok && !normalized {
return dotMetricsMap[metricName]
} else {
return metricName
}
}

View File

@@ -258,7 +258,11 @@ func PrepareTimeseriesFilterQuery(start, end int64, mq *v3.BuilderQuery) (string
conditions = append(conditions, fmt.Sprintf("metric_name IN %s", utils.ClickHouseFormattedMetricNames(mq.AggregateAttribute.Key)))
conditions = append(conditions, fmt.Sprintf("temporality = '%s'", mq.Temporality))
conditions = append(conditions, "__normalized = false")
if constants.IsDotMetricsEnabled {
conditions = append(conditions, "__normalized = false")
} else {
conditions = append(conditions, "__normalized = true")
}
start, end, tableName := whichTSTableToUse(start, end, mq)
@@ -350,7 +354,11 @@ func PrepareTimeseriesFilterQueryV3(start, end int64, mq *v3.BuilderQuery) (stri
conditions = append(conditions, fmt.Sprintf("metric_name IN %s", utils.ClickHouseFormattedMetricNames(mq.AggregateAttribute.Key)))
conditions = append(conditions, fmt.Sprintf("temporality = '%s'", mq.Temporality))
conditions = append(conditions, "__normalized = false")
if constants.IsDotMetricsEnabled {
conditions = append(conditions, "__normalized = false")
} else {
conditions = append(conditions, "__normalized = true")
}
start, end, tableName := whichTSTableToUse(start, end, mq)

View File

@@ -6,6 +6,9 @@ import (
"strings"
"sync"
"github.com/prometheus/prometheus/promql/parser"
"github.com/SigNoz/signoz/pkg/errors"
logsV4 "github.com/SigNoz/signoz/pkg/query-service/app/logs/v4"
metricsV3 "github.com/SigNoz/signoz/pkg/query-service/app/metrics/v3"
metricsV4 "github.com/SigNoz/signoz/pkg/query-service/app/metrics/v4"
@@ -275,3 +278,59 @@ func (q *querier) runBuilderQuery(
Series: resultSeries,
}
}
// ValidateMetricNames function is used to print all those queries who are still using old normalized metrics and not new metrics.
func (q *querier) ValidateMetricNames(ctx context.Context, query *v3.CompositeQuery, orgID valuer.UUID) {
var metricNames []string
switch query.QueryType {
case v3.QueryTypePromQL:
for _, query := range query.PromQueries {
expr, err := q.parser.ParseExpr(query.Query)
if err != nil {
q.logger.DebugContext(ctx, "error parsing promql expression", "query", query.Query, errors.Attr(err))
continue
}
parser.Inspect(expr, func(node parser.Node, path []parser.Node) error {
if vs, ok := node.(*parser.VectorSelector); ok {
for _, m := range vs.LabelMatchers {
if m.Name == "__name__" {
metricNames = append(metricNames, m.Value)
}
}
}
return nil
})
}
metrics, err := q.reader.GetNormalizedStatus(ctx, orgID, metricNames)
if err != nil {
q.logger.DebugContext(ctx, "error getting corresponding normalized metrics", errors.Attr(err))
return
}
for metricName, metricPresent := range metrics {
if metricPresent {
continue
} else {
q.logger.WarnContext(ctx, "using normalized metric name", "metrics", metricName)
continue
}
}
case v3.QueryTypeBuilder:
for _, query := range query.BuilderQueries {
metricName := query.AggregateAttribute.Key
metricNames = append(metricNames, metricName)
}
metrics, err := q.reader.GetNormalizedStatus(ctx, orgID, metricNames)
if err != nil {
q.logger.DebugContext(ctx, "error getting corresponding normalized metrics", errors.Attr(err))
return
}
for metricName, metricPresent := range metrics {
if metricPresent {
continue
} else {
q.logger.WarnContext(ctx, "using normalized metric name", "metrics", metricName)
continue
}
}
}
}

View File

@@ -515,6 +515,9 @@ func (q *querier) QueryRange(ctx context.Context, orgID valuer.UUID, params *v3.
var results []*v3.Result
var err error
var errQueriesByName map[string]error
if !q.testingMode && q.reader != nil {
q.ValidateMetricNames(ctx, params.CompositeQuery, orgID)
}
if params.CompositeQuery != nil {
switch params.CompositeQuery.QueryType {
case v3.QueryTypeBuilder:

View File

@@ -212,7 +212,7 @@ func TestBuildQueryWithThreeOrMoreQueriesRefAndFormula(t *testing.T) {
// So(queries["F5"], ShouldContainSubstring, "SELECT A.ts as ts, ((A.value - B.value) / B.value) * 100")
// So(strings.Count(queries["F5"], " ON "), ShouldEqual, 1)
})
t.Run("TestBuildQueryWithMetricNameAndAttribute", func(t *testing.T) {
t.Run("TestBuildQueryWithDotMetricNameAndAttribute", func(t *testing.T) {
q := &v3.QueryRangeParamsV3{
Start: 1735036101000,
End: 1735637901000,

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

@@ -3,6 +3,7 @@ package constants
import (
"maps"
"os"
"regexp"
"strconv"
"github.com/SigNoz/signoz/pkg/query-service/model"
@@ -25,6 +26,12 @@ const OrderBySpanCount = "span_count"
var MetricsExplorerClickhouseThreads = GetOrDefaultEnvInt("METRICS_EXPLORER_CLICKHOUSE_THREADS", 8)
var UpdatedMetricsMetadataCachePrefix = GetOrDefaultEnv("METRICS_UPDATED_METADATA_CACHE_KEY", "UPDATED_METRICS_METADATA")
const NormalizedMetricsMapCacheKey = "NORMALIZED_METRICS_MAP_CACHE_KEY"
const NormalizedMetricsMapQueryThreads = 10
var NormalizedMetricsMapRegex = regexp.MustCompile(`[^a-zA-Z0-9]`)
var NormalizedMetricsMapQuantileRegex = regexp.MustCompile(`(?i)([._-]?quantile.*)$`)
func GetEvalDelay() valuer.TextDuration {
evalDelayStr := GetOrDefaultEnv("RULES_EVAL_DELAY", "2m")
evalDelayDuration, err := valuer.ParseTextDuration(evalDelayStr)
@@ -664,11 +671,16 @@ var OldToNewTraceFieldsMap = map[string]string{
var StaticFieldsTraces = map[string]v3.AttributeKey{}
var IsDotMetricsEnabled = false
var MaxJSONFlatteningDepth = 1
func init() {
StaticFieldsTraces = maps.Clone(NewStaticFieldsTraces)
maps.Copy(StaticFieldsTraces, DeprecatedStaticFieldsTraces)
if GetOrDefaultEnv(DotMetricsEnabled, "true") == "true" {
IsDotMetricsEnabled = true
}
// set max flattening depth
depth, err := strconv.Atoi(GetOrDefaultEnv(maxJSONFlatteningDepth, "1"))
if err == nil {
@@ -696,4 +708,5 @@ var MaterializedDataTypeMap = map[string]string{
const InspectMetricsMaxTimeDiff = 1800000
const DotMetricsEnabled = "DOT_METRICS_ENABLED"
const maxJSONFlatteningDepth = "MAX_JSON_FLATTENING_DEPTH"

View File

@@ -108,6 +108,7 @@ type Reader interface {
GetUpdatedMetricsMetadata(ctx context.Context, orgID valuer.UUID, metricNames ...string) (map[string]*model.UpdateMetricsMetadata, *model.ApiError)
CheckForLabelsInMetric(ctx context.Context, orgID valuer.UUID, metricName string, labels []string) (bool, *model.ApiError)
GetNormalizedStatus(ctx context.Context, orgID valuer.UUID, metricNames []string) (map[string]bool, error)
}
type Querier interface {

View File

@@ -1,14 +1,27 @@
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) string {
if transitionedMetric, ok := MetricsUnderTransition[metric]; ok {
return transitionedMetric
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
}
return metric
}

View File

@@ -0,0 +1,15 @@
package model
import "encoding/json"
type MetricsNormalizedMap struct {
MetricName string `json:"metricName"`
IsUnNormalized bool `json:"isUnNormalized"`
}
func (c *MetricsNormalizedMap) MarshalBinary() (data []byte, err error) {
return json.Marshal(c)
}
func (c *MetricsNormalizedMap) UnmarshalBinary(data []byte) error {
return json.Unmarshal(data, c)
}

View File

@@ -234,7 +234,7 @@ func ClickHouseFormattedValue(v interface{}) string {
func ClickHouseFormattedMetricNames(v interface{}) string {
if name, ok := v.(string); ok {
transitionedMetrics := metrics.GetTransitionedMetric(name)
transitionedMetrics := metrics.GetTransitionedMetric(name, !constants.IsDotMetricsEnabled)
if transitionedMetrics != name {
return ClickHouseFormattedValue([]interface{}{transitionedMetrics})
} else {

View File

@@ -10,6 +10,7 @@ import (
"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"
@@ -982,25 +983,78 @@ 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{}
selector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: field.Signal,
FieldContext: field.FieldContext,
}
members := semconv.Members(semconv.KindAttribute, selector)
isFamily := len(members) > 1
fieldKeysForName := make([]*telemetrytypes.TelemetryFieldKey, 0)
indexByIdentity := make(map[string]int)
// 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)
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
}
// A wildcard lookup may have found a same-named field in a scope where
// this family does not apply. Keep exact names, but reject cross-member
// matches outside the generated family scope.
if memberName != field.Name {
itemSelector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
}
if !slices.Contains(semconv.Members(semconv.KindAttribute, itemSelector), memberName) {
continue
}
}
physicalMembers := item.SemconvMembers
if len(physicalMembers) == 0 {
physicalMembers = []string{memberName}
}
identity := item.Signal.StringValue() + ";" + item.FieldContext.StringValue() + ";" + item.FieldDataType.StringValue()
if isFamily {
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)
}
}
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 isFamily {
resolved.Name = field.Name
resolved.SemconvMembers = slices.Clone(physicalMembers)
}
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)
}
}

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,53 @@ func TestVisitKey(t *testing.T) {
}
}
func TestMatchingFieldKeysResolvesSemconvFamily(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,
}
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
}
for _, requestedName := range []string{current.Name, old.Name} {
requested := telemetrytypes.NewTelemetryFieldKey(
requestedName,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeys(requested, fieldKeys)
require.Len(t, matches, 1, "family lookup must return one field before its metadata is inspected")
assert.Equal(t, requestedName, matches[0].Name)
assert.Equal(t, "current metadata", matches[0].Description)
assert.Equal(t, []string{current.Name, old.Name}, matches[0].SemconvMembers)
}
// A current-name query still resolves when metadata has seen only the old
// spelling. The returned name remains the request identity.
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{old.Name: {old}})
require.Len(t, matches, 1, "family lookup must return one field before its metadata is inspected")
assert.Equal(t, current.Name, matches[0].Name)
assert.Equal(t, "old metadata", matches[0].Description)
assert.Equal(t, []string{old.Name}, matches[0].SemconvMembers)
}
// ---------------------------------------------------------------------------
// TestVisitComparison
// ---------------------------------------------------------------------------

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

@@ -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

@@ -198,6 +198,75 @@ func TestConditionBuilder(t *testing.T) {
expected: "match(simpleJSONExtractString(labels, 'k8s.namespace.name'), ?) AND labels LIKE ?",
expectedArgs: []any{"ban.*", "%k8s.namespace.name%"},
},
{
name: "semantic convention family equality uses current-first fallback",
key: &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
},
op: qbtypes.FilterOperatorEqual,
value: "production",
expected: "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), '')) = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)",
expectedArgs: []any{
"production",
"%deployment.environment.name%",
"%deployment.environment%",
`%deployment.environment.name":"production%`,
`%deployment.environment":"production%`,
},
},
{
name: "old semantic convention request uses only current metadata member",
key: &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name"},
},
op: qbtypes.FilterOperatorEqual,
value: "production",
expected: "simpleJSONExtractString(labels, 'deployment.environment.name') = ? AND labels LIKE ? AND labels LIKE ?",
expectedArgs: []any{"production", "%deployment.environment.name%", `%deployment.environment.name":"production%`},
},
{
name: "semantic convention family negative filter does not reject fallback rows",
key: &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
},
op: qbtypes.FilterOperatorNotEqual,
value: "staging",
expected: "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), '')) <> ?",
expectedArgs: []any{"staging"},
},
{
name: "semantic convention family exists checks every member",
key: &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
},
op: qbtypes.FilterOperatorExists,
expected: "(simpleJSONHas(labels, 'deployment.environment.name') = ? OR simpleJSONHas(labels, 'deployment.environment') = ?) AND (labels LIKE ? OR labels LIKE ?)",
expectedArgs: []any{true, true, "%deployment.environment.name%", "%deployment.environment%"},
},
{
name: "semantic convention family not exists checks every member",
key: &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
},
op: qbtypes.FilterOperatorNotExists,
expected: "(simpleJSONHas(labels, 'deployment.environment.name') <> ? AND simpleJSONHas(labels, 'deployment.environment') <> ?)",
expectedArgs: []any{true, true},
},
}
fm := NewFieldMapper()

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,20 @@ func NewFieldMapper() *defaultFieldMapper {
return &defaultFieldMapper{}
}
func resourceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if 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 +82,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

@@ -1,35 +0,0 @@
package telemetrymetadata
import "github.com/SigNoz/signoz/pkg/types/telemetrytypes"
type BackwardCompatibleKeyMap map[string]string
var (
TracesBackwardCompatKeys = BackwardCompatibleKeyMap{
"net.peer.name": "server.address",
"server.address": "net.peer.name",
"http.url": "url.full",
"url.full": "http.url",
}
// LogsBackwardCompatKeys contains bidirectional mappings for logs.
// Currently empty, can be extended in the future.
LogsBackwardCompatKeys = BackwardCompatibleKeyMap{}
// MetricsBackwardCompatKeys contains bidirectional mappings for metrics.
// Currently empty, can be extended in the future.
MetricsBackwardCompatKeys = BackwardCompatibleKeyMap{}
)
func GetBackwardCompatKeysForSignal(signal telemetrytypes.Signal) BackwardCompatibleKeyMap {
switch signal {
case telemetrytypes.SignalTraces:
return TracesBackwardCompatKeys
case telemetrytypes.SignalLogs:
return LogsBackwardCompatKeys
case telemetrytypes.SignalMetrics:
return MetricsBackwardCompatKeys
default:
return BackwardCompatibleKeyMap{}
}
}

View File

@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"log/slog"
"slices"
"strings"
"time"
@@ -14,6 +15,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 +153,102 @@ func (t *telemetryMetaStore) tracesTblStatementToFieldKeys(ctx context.Context)
return materialisedKeys, nil
}
func traceSemconvMembers(name string, fieldContext telemetrytypes.FieldContext) []string {
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: fieldContext,
})
}
func traceSemconvDuplicateFactor() int {
factor := 1
for _, family := range semconv.All() {
if family.Kind != semconv.KindAttribute {
continue
}
if _, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: family.Current,
Signal: telemetrytypes.SignalTraces,
}); ok {
factor = max(factor, len(family.Old)+1)
}
}
return factor
}
// canonicalizeTraceSemconvKeys presents one current-name key for each family.
// If metadata contains both spellings, metadata attached to the current name
// wins; otherwise the old entry is copied under the current response name.
func canonicalizeTraceSemconvKeys(keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
result := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
indexByIdentity := make(map[string]int)
currentSourceByIdentity := make(map[string]bool)
for _, key := range keys {
if key.Signal != telemetrytypes.SignalTraces {
result = append(result, key)
continue
}
family, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: key.FieldContext,
})
if !ok {
result = append(result, key)
continue
}
resolved := *key
resolved.Name = family.Current
resolved.SemconvMembers = []string{key.Name}
identity := resolved.Name + ";" + resolved.Signal.StringValue() + ";" + resolved.FieldContext.StringValue() + ";" + resolved.FieldDataType.StringValue()
fromCurrent := key.Name == family.Current
if index, found := indexByIdentity[identity]; found {
physicalMembers := result[index].SemconvMembers
if fromCurrent && !currentSourceByIdentity[identity] {
result[index] = &resolved
currentSourceByIdentity[identity] = true
}
for _, member := range physicalMembers {
if !slices.Contains(result[index].SemconvMembers, member) {
result[index].SemconvMembers = append(result[index].SemconvMembers, member)
}
}
if !slices.Contains(result[index].SemconvMembers, key.Name) {
result[index].SemconvMembers = append(result[index].SemconvMembers, key.Name)
}
continue
}
indexByIdentity[identity] = len(result)
currentSourceByIdentity[identity] = fromCurrent
result = append(result, &resolved)
}
for _, key := range result {
if len(key.SemconvMembers) < 2 {
continue
}
present := make(map[string]bool, len(key.SemconvMembers))
for _, member := range key.SemconvMembers {
present[member] = true
}
ordered := make([]string, 0, len(key.SemconvMembers))
for _, member := range traceSemconvMembers(key.Name, key.FieldContext) {
if present[member] {
ordered = append(ordered, member)
}
}
key.SemconvMembers = ordered
}
return result
}
// 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{
@@ -202,10 +300,23 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
// key part of the selector
fieldKeyConds := []string{}
members := traceSemconvMembers(fieldKeySelector.Name, fieldKeySelector.FieldContext)
if fieldKeySelector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeExact {
fieldKeyConds = append(fieldKeyConds, sb.E("tagKey", fieldKeySelector.Name))
if len(members) == 1 {
fieldKeyConds = append(fieldKeyConds, sb.E("tagKey", members[0]))
} else {
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
fieldKeyConds = append(fieldKeyConds, sb.In("tagKey", memberValues...))
}
} else {
fieldKeyConds = append(fieldKeyConds, sb.ILike("tagKey", "%"+escapeForLike(fieldKeySelector.Name)+"%"))
memberConditions := make([]string, 0, len(members))
for _, member := range members {
memberConditions = append(memberConditions, sb.ILike("tagKey", "%"+escapeForLike(member)+"%"))
}
fieldKeyConds = append(fieldKeyConds, sb.Or(memberConditions...))
}
searchTexts = append(searchTexts, fieldKeySelector.Name)
@@ -238,8 +349,10 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
mainSb.From(mainSb.BuilderAs(sb, "sub_query"))
mainSb.GroupBy("tag_key", "tag_type", "tag_data_type")
mainSb.OrderBy("priority")
// query one extra to check if we hit the limit
mainSb.Limit(limit + 1)
// Family members collapse after the database query. In the worst case each
// logical key occupies one row per family member, so fetch enough physical
// rows to return the requested logical page, plus one to detect truncation.
mainSb.Limit(limit*traceSemconvDuplicateFactor() + 1)
query, args := mainSb.BuildWithFlavor(sqlbuilder.ClickHouse)
@@ -249,14 +362,7 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
}
defer rows.Close()
keys := []*telemetrytypes.TelemetryFieldKey{}
rowCount := 0
for rows.Next() {
rowCount++
// reached the limit, we know there are more results
if rowCount > limit {
break
}
var name string
var fieldContext telemetrytypes.FieldContext
var fieldDataType telemetrytypes.FieldDataType
@@ -285,8 +391,11 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
return nil, false, errors.Wrap(rows.Err(), errors.TypeInternal, errors.CodeInternal, ErrFailedToGetTracesKeys.Error())
}
// hit the limit? (only counting DB results)
complete := rowCount <= limit
keys = canonicalizeTraceSemconvKeys(keys)
complete := len(keys) <= limit
if !complete {
keys = keys[:limit]
}
staticKeys := []string{"isRoot", "isEntryPoint"}
staticKeys = append(staticKeys, maps.Keys(tracestelemetryschema.IntrinsicFields)...)
@@ -1108,40 +1217,6 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
}
// applyBackwardCompatibleKeys adds backward compatible key aliases to the map.
func applyBackwardCompatibleKeys(mapOfKeys map[string][]*telemetrytypes.TelemetryFieldKey) {
// Get backward compatible keys for all signals
backwardCompatKeysBySignal := map[telemetrytypes.Signal]BackwardCompatibleKeyMap{
telemetrytypes.SignalTraces: GetBackwardCompatKeysForSignal(telemetrytypes.SignalTraces),
telemetrytypes.SignalLogs: GetBackwardCompatKeysForSignal(telemetrytypes.SignalLogs),
telemetrytypes.SignalMetrics: GetBackwardCompatKeysForSignal(telemetrytypes.SignalMetrics),
}
// Iterate over existing keys and add aliases if they exist in backward compat mapping
for srcKey, srcKeys := range mapOfKeys {
for _, srcKeyEntry := range srcKeys {
backwardCompatKeys := backwardCompatKeysBySignal[srcKeyEntry.Signal]
if backwardCompatKeys == nil {
continue
}
if aliasKey, ok := backwardCompatKeys[srcKey]; ok {
if _, aliasExists := mapOfKeys[aliasKey]; !aliasExists {
aliasKeyEntry := &telemetrytypes.TelemetryFieldKey{
Name: aliasKey,
Signal: srcKeyEntry.Signal,
FieldContext: srcKeyEntry.FieldContext,
FieldDataType: srcKeyEntry.FieldDataType,
}
mapOfKeys[aliasKey] = []*telemetrytypes.TelemetryFieldKey{aliasKeyEntry}
}
// Found the alias for this signal, no need to check other entries
break
}
}
}
}
func enrichWithIntrinsicMetricKeys(keys map[string][]*telemetrytypes.TelemetryFieldKey, selectors []*telemetrytypes.FieldKeySelector) map[string][]*telemetrytypes.TelemetryFieldKey {
if len(selectors) == 0 {
return keys
@@ -1272,7 +1347,6 @@ func (t *telemetryMetaStore) GetKeys(ctx context.Context, orgID valuer.UUID, fie
mapOfKeys[key.Name] = append(mapOfKeys[key.Name], key)
}
applyBackwardCompatibleKeys(mapOfKeys)
mapOfKeys = enrichWithIntrinsicMetricKeys(mapOfKeys, selectors)
if t.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableAIObservability, featuretypes.NewFlaggerEvaluationContext(orgID)) {
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, selectors)
@@ -1353,7 +1427,6 @@ func (t *telemetryMetaStore) GetKeysMulti(ctx context.Context, orgID valuer.UUID
mapOfKeys[key.Name] = append(mapOfKeys[key.Name], key)
}
applyBackwardCompatibleKeys(mapOfKeys)
mapOfKeys = enrichWithIntrinsicMetricKeys(mapOfKeys, fieldKeySelectors)
if t.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableAIObservability, featuretypes.NewFlaggerEvaluationContext(orgID)) {
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, fieldKeySelectors)
@@ -1367,7 +1440,18 @@ func (t *telemetryMetaStore) GetKey(ctx context.Context, orgID valuer.UUID, fiel
if err != nil {
return nil, err
}
return keys[fieldKeySelector.Name], nil
members := semconv.Members(semconv.KindAttribute, *fieldKeySelector)
resolved := make([]*telemetrytypes.TelemetryFieldKey, 0)
seen := make(map[*telemetrytypes.TelemetryFieldKey]bool)
for _, member := range members {
for _, key := range keys[member] {
if !seen[key] {
resolved = append(resolved, key)
seen[key] = true
}
}
}
return resolved, nil
}
func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *telemetrytypes.FieldValueSelector) ([]string, bool, error) {
@@ -1542,7 +1626,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

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,83 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
t.Errorf("there were unfulfilled expectations: %s", err)
}
}
func TestCanonicalizeTraceSemconvKeys(t *testing.T) {
oldResource := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Description: "old resource metadata",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
currentResource := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Description: "current resource metadata",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
oldAttribute := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
oldLogResource := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
result := canonicalizeTraceSemconvKeys([]*telemetrytypes.TelemetryFieldKey{
oldResource,
currentResource,
oldAttribute,
oldLogResource,
})
require.Len(t, result, 3)
assert.Equal(t, "deployment.environment.name", result[0].Name)
assert.Equal(t, "current resource metadata", result[0].Description)
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, result[0].SemconvMembers)
assert.Equal(t, "deployment.environment.name", result[1].Name)
assert.Equal(t, telemetrytypes.FieldContextAttribute, result[1].FieldContext)
assert.Equal(t, []string{"deployment.environment"}, result[1].SemconvMembers)
assert.Equal(t, "deployment.environment", result[2].Name, "phase 1 must not rewrite raw log metadata")
}
func TestGetSpanFieldValuesMergesSemconvFamily(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

@@ -154,6 +154,15 @@ 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) {
if fm, ok := c.fm.(*fieldMapper); ok {
return fm.existsExpressionFor(ctx, orgID, startNs, endNs, key, operator == qbtypes.FilterOperatorExists)
}
}
columns, err := c.fm.ColumnFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err

View File

@@ -210,6 +210,32 @@ func TestConditionFor(t *testing.T) {
expectedSQL: "NOT mapContains(attributes_string, 'user.id')",
expectedError: nil,
},
{
name: "Equal operator - semantic convention family",
key: telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
},
operator: qbtypes.FilterOperatorEqual,
value: "production",
expectedSQL: "(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'))))",
expectedArgs: []any{"production"},
expectedError: nil,
},
{
name: "Not Exists operator - semantic convention family",
key: telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
},
operator: qbtypes.FilterOperatorNotExists,
expectedSQL: "NOT (((mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))))",
expectedError: nil,
},
{
name: "Exists operator - json field",
key: telemetrytypes.TelemetryFieldKey{

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,32 @@ 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 (m *fieldMapper) getColumn(
_ context.Context,
_, _ uint64,
@@ -291,10 +318,25 @@ 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))
}
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 +361,35 @@ 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 {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
members := traceSemconvMembers(key)
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)
for i, member := range members {
branches = append(branches, guards[i], fmt.Sprintf("%s['%s']", columnName, member))
}
exprs = append(exprs, "multiIf("+strings.Join(branches, ", ")+", NULL)")
}
existExprs = append(existExprs, "("+strings.Join(guards, " OR ")+")")
} else if key.Materialized {
// a key could have been materialized, if so return the materialized column name
physicalKey := *key
physicalKey.Name = members[0]
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(&physicalKey))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(&physicalKey))
} 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 +593,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

@@ -80,7 +80,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 mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(resources_string['deployment.environment'], '')), NULL)",
expectedError: nil,
},
{
@@ -120,6 +120,51 @@ func TestGetFieldKeyName(t *testing.T) {
}
}
func TestFieldForSemconvFamily(t *testing.T) {
ctx := context.Background()
fm := NewFieldMapper()
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())
for _, requestedName := range []string{"deployment.environment.name", "deployment.environment"} {
attributeKey := telemetrytypes.TelemetryFieldKey{
Name: requestedName,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
attributeExpression, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, &attributeKey)
require.NoError(t, err)
assert.Equal(t,
"COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''))",
attributeExpression,
)
resourceKey := telemetrytypes.TelemetryFieldKey{
Name: requestedName,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
Evolutions: MockEvolutionData(time.Date(2024, 6, 2, 0, 0, 0, 0, time.UTC)),
}
resourceExpression, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, &resourceKey)
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, '')), (mapContains(resources_string, 'deployment.environment.name') OR mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(resources_string['deployment.environment'], '')), NULL)",
resourceExpression,
)
}
oldRequestWithCurrentOnlyMetadata := telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
SemconvMembers: []string{"deployment.environment.name"},
}
expression, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, &oldRequestWithCurrentOnlyMetadata)
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 +221,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 +234,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 mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(resources_string['deployment.environment'], '')), NULL)",
},
}

View File

@@ -12,7 +12,8 @@ var (
Gateway = valuer.NewString("gateway")
PremiumSupport = valuer.NewString("premium_support")
AnomalyDetection = valuer.NewString("anomaly_detection")
AnomalyDetection = valuer.NewString("anomaly_detection")
DotMetricsEnabled = valuer.NewString("dot_metrics_enabled")
// License State.
LicenseStatusInvalid = valuer.NewString("invalid")
@@ -59,6 +60,13 @@ var BasicPlan = []*Feature{
UsageLimit: -1,
Route: "",
},
{
Name: DotMetricsEnabled,
Active: false,
Usage: 0,
UsageLimit: -1,
Route: "",
},
}
var EnterprisePlan = []*Feature{
@@ -104,6 +112,21 @@ var EnterprisePlan = []*Feature{
UsageLimit: -1,
Route: "",
},
{
Name: DotMetricsEnabled,
Active: false,
Usage: 0,
UsageLimit: -1,
Route: "",
},
}
var DefaultFeatureSet = []*Feature{}
var DefaultFeatureSet = []*Feature{
{
Name: DotMetricsEnabled,
Active: false,
Usage: 0,
UsageLimit: -1,
Route: "",
},
}

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

@@ -47,7 +47,8 @@ type TelemetryFieldKey struct {
Indexes []TelemetryFieldKeySkipIndex `json:"-"`
Materialized bool `json:"-"` // refers to promoted in case of body.... fields
Evolutions []*EvolutionEntry `json:"-"`
Evolutions []*EvolutionEntry `json:"-"`
SemconvMembers []string `json:"-"`
}
func (f *TelemetryFieldKey) KeyNameContainsArray() bool {
@@ -128,6 +129,7 @@ func (f *TelemetryFieldKey) OverrideMetadataFrom(src *TelemetryFieldKey) {
f.Materialized = src.Materialized
f.JSONPlan = src.JSONPlan
f.Evolutions = src.Evolutions
f.SemconvMembers = src.SemconvMembers
}
func (f *TelemetryFieldKey) Equal(key *TelemetryFieldKey) bool {

View File

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

View File

@@ -704,7 +704,6 @@ def build_raw_query(
order: list[dict] | None = None,
limit: int | None = None,
filter_expression: str | None = None,
select_fields: list[dict] | None = None,
step_interval: int = DEFAULT_STEP_INTERVAL,
disabled: bool = False,
) -> dict:
@@ -724,9 +723,6 @@ def build_raw_query(
if filter_expression:
spec["filter"] = {"expression": filter_expression}
if select_fields:
spec["selectFields"] = select_fields
return {"type": "builder_query", "spec": spec}

View File

@@ -1,124 +0,0 @@
"""Seed data for the queriercommon keyless-semantics tests.
Three identities exist in every signal. GOLD and SILVER carry the test keys.
NONE carries no key at all. The tests assert which identities a filter
returns, so the membership of NONE is the point of every case.
The attribute names are outside every semantic-convention family, so the
seeded data pins base behavior with any semconv overlay state.
"""
from collections.abc import Callable, Generator
from datetime import UTC, datetime, timedelta
import pytest
from fixtures.logs import Logs
from fixtures.metrics import Metrics
from fixtures.querier import aligned_epoch
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
PREFIX = "keyless-sem"
STRING_KEY = "tenant.tier"
NUMBER_KEY = "retry.count"
METRIC_NAME = "keyless_semantics_gauge"
METRIC_LABEL = "tenant_tier"
# Row identities, keyed by the value of the string key that each row carries.
GOLD = f"{PREFIX}-gold"
SILVER = f"{PREFIX}-silver"
NONE = f"{PREFIX}-none" # carries no string key and no number key
# (identity, string-key value, number-key value, insert offset)
_ROWS = [
(GOLD, "gold", 0, timedelta(seconds=3)),
(SILVER, "silver", 5, timedelta(seconds=2)),
(NONE, None, None, timedelta(seconds=1)),
]
def _resources(identity: str, tier: str | None) -> dict:
base = {"service.name": identity}
if tier is not None:
base[STRING_KEY] = tier
return base
def _attributes(tier: str | None, retries: int | None) -> dict:
attrs: dict = {}
if tier is not None:
attrs[STRING_KEY] = tier
if retries is not None:
attrs[NUMBER_KEY] = retries
return attrs
@pytest.fixture(name="keyless_rows", scope="function")
def keyless_rows(
insert_logs: Callable[[list[Logs]], None],
insert_traces: Callable[[list[Traces]], None],
) -> Generator[datetime]:
"""Inserts one span and one log per identity: GOLD (string "gold",
number 0), SILVER (string "silver", number 5), and NONE (no keys).
Yields the base timestamp. Span name and log body are the identity."""
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=1)
insert_traces(
[
Traces(
timestamp=now - offset,
duration=timedelta(milliseconds=10),
trace_id=TraceIdGenerator.trace_id(),
span_id=TraceIdGenerator.span_id(),
name=identity,
kind=TracesKind.SPAN_KIND_SERVER,
status_code=TracesStatusCode.STATUS_CODE_OK,
resources=_resources(identity, tier),
attributes=_attributes(tier, retries),
)
for identity, tier, retries, offset in _ROWS
]
)
insert_logs(
[
Logs(
timestamp=now - offset,
body=identity,
resources=_resources(identity, tier),
attributes=_attributes(tier, retries),
)
for identity, tier, retries, offset in _ROWS
]
)
yield now
@pytest.fixture(name="keyless_series", scope="function")
def keyless_series(insert_metrics: Callable[[list[Metrics]], None]) -> Generator[tuple[int, int]]:
"""Inserts three gauge series: GOLD and SILVER carry the metric label,
NONE does not. The `service` label is the identity. Yields the
(start, end) epoch-second window that covers the points."""
start = aligned_epoch(timedelta(minutes=30))
points = 5
def labels(identity: str, tier: str | None) -> dict:
base = {"service": identity}
if tier is not None:
base[METRIC_LABEL] = tier
return base
insert_metrics(
[
Metrics(
metric_name=METRIC_NAME,
labels=labels(identity, tier),
timestamp=datetime.fromtimestamp(start + minute * 60, tz=UTC),
value=10.0,
type_="Gauge",
is_monotonic=False,
)
for identity, tier in ((GOLD, "gold"), (SILVER, "silver"), (NONE, None))
for minute in range(points)
]
)
yield start, start + points * 60

View File

@@ -1,202 +0,0 @@
"""Pins the keyless-row contract for filter operators, per signal.
The contract (deliberate product semantics, enforced by
`FilterOperator.AddDefaultExistsFilter` in
pkg/types/querybuildertypes/querybuildertypesv5/builder_elements.go):
- Negative operators (!=, NOT IN, NOT LIKE, NOT CONTAINS, ...) are a set
complement over ALL rows: a row that does not carry the key at all MUST
match. Users opt into presence explicitly with `AND key EXISTS`.
- Positive operators carry an implicit existence guard: a keyless row must
NOT match `key = ''`-style comparisons against sentinel defaults.
- EXISTS / NOT EXISTS partition rows exactly by key presence.
- Numeric attributes inherit the map-default sentinel: a missing key reads
as 0, so `num != 0` excludes keyless rows while `num != 5` includes them.
This conflation is deliberate and pinned here as the reference for any
value-expression change (for example coalesce tails in semconv families).
Any implementation change that makes these assertions fail is a behavior
break, not a cleanup. Family-field behavior must mirror this matrix; see
queriertraces/13_semconv_evolution.py.
Seed data lives in fixtures/queriercommon.py: GOLD and SILVER carry the
keys, NONE carries none. Every case asserts which identities a filter
returns, so the membership of NONE is the point of each case.
"""
from collections.abc import Callable
from datetime import datetime, timedelta
from http import HTTPStatus
import pytest
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.querier import (
RequestType,
build_builder_query,
build_order_by,
build_raw_query,
get_all_series,
get_column_data_from_response,
make_query_request,
)
from fixtures.queriercommon import (
GOLD,
METRIC_LABEL,
METRIC_NAME,
NONE,
NUMBER_KEY,
PREFIX,
SILVER,
STRING_KEY,
)
STRING_MATRIX = [
pytest.param("{key} = 'gold'", {GOLD}, id="eq_excludes_keyless"),
pytest.param("{key} != 'gold'", {SILVER, NONE}, id="neq_includes_keyless"),
pytest.param("{key} NOT IN ['gold', 'silver']", {NONE}, id="not_in_includes_keyless"),
pytest.param("NOT {key} LIKE '%gold%'", {SILVER, NONE}, id="not_like_includes_keyless"),
pytest.param("{key} NOT CONTAINS 'gol'", {SILVER, NONE}, id="not_contains_includes_keyless"),
pytest.param("{key} EXISTS", {GOLD, SILVER}, id="exists_partitions"),
pytest.param("{key} NOT EXISTS", {NONE}, id="not_exists_partitions"),
# "Present and not X" is a composition. It is not a new operator
# semantic.
pytest.param("{key} != 'gold' AND {key} EXISTS", {SILVER}, id="neq_composed_with_exists"),
]
NUMBER_MATRIX = [
pytest.param("{key} = 0", {GOLD}, id="numeric_eq_zero_excludes_keyless"),
pytest.param("{key} != 5", {GOLD, NONE}, id="numeric_neq_includes_keyless"),
pytest.param("{key} != 0", {SILVER}, id="numeric_neq_zero_sentinel_conflation"),
]
SIGNALS = [
pytest.param("traces", "span.name", "name", id="traces"),
pytest.param("logs", "body", "body", id="logs"),
]
@pytest.mark.parametrize("expression_template,expected", STRING_MATRIX)
@pytest.mark.parametrize("context", ["resource", "attribute"])
@pytest.mark.parametrize("signal,identity_field,identity_column", SIGNALS)
def test_negative_operators_include_keyless_rows(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
keyless_rows: datetime,
signal: str,
identity_field: str,
identity_column: str,
context: str,
expression_template: str,
expected: set[str],
) -> None:
"""A negative operator is a set complement over all rows. Presence is an
explicit EXISTS opt-in. The contract holds the same way for resource and
attribute contexts, on traces and on logs."""
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
expression = expression_template.format(key=f"{context}.{STRING_KEY}")
response = make_query_request(
signoz,
token,
start_ms=int((keyless_rows - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((keyless_rows + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
signal,
limit=100,
filter_expression=expression,
order=[build_order_by("timestamp", "asc")],
select_fields=[{"name": identity_field}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
# Sets keep the assertion stable when the shared stack is reused and
# older rows with the same identities remain.
matched = {name for name in get_column_data_from_response(response.json(), identity_column) if name.startswith(PREFIX)}
assert matched == expected, expression
@pytest.mark.parametrize("expression_template,expected", NUMBER_MATRIX)
@pytest.mark.parametrize("signal,identity_field,identity_column", SIGNALS)
def test_numeric_sentinel_semantics(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
keyless_rows: datetime,
signal: str,
identity_field: str,
identity_column: str,
expression_template: str,
expected: set[str],
) -> None:
"""A missing numeric key reads as the map default 0. `!= 0` therefore
excludes rows without the key, and every other negative comparison
includes them. This is inherited sentinel behavior, pinned on purpose."""
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
expression = expression_template.format(key=f"attribute.{NUMBER_KEY}")
response = make_query_request(
signoz,
token,
start_ms=int((keyless_rows - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((keyless_rows + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
signal,
limit=100,
filter_expression=expression,
order=[build_order_by("timestamp", "asc")],
select_fields=[{"name": identity_field}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
matched = {name for name in get_column_data_from_response(response.json(), identity_column) if name.startswith(PREFIX)}
assert matched == expected, expression
@pytest.mark.parametrize("expression_template,expected", STRING_MATRIX)
def test_metrics_negative_operators_include_keyless_series(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
keyless_series: tuple[int, int],
expression_template: str,
expected: set[str],
) -> None:
"""The same contract holds for metric labels: a series without the label
matches every negative filter on it, and EXISTS opts into presence."""
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
expression = expression_template.format(key=METRIC_LABEL)
start, end = keyless_series
response = make_query_request(
signoz,
token,
start_ms=start * 1000,
end_ms=end * 1000,
queries=[
build_builder_query(
"A",
METRIC_NAME,
"avg",
"sum",
group_by=["service"],
filter_expression=expression,
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
matched = {label["value"] for series in get_all_series(response.json(), "A") for label in series["labels"] if label["key"]["name"] == "service" and label["value"].startswith(PREFIX)}
assert matched == expected, expression

View File

@@ -0,0 +1,239 @@
"""Phase 1 end-to-end checks for semantic-convention name evolution.
The fixture models a fleet split across SDK generations and deliberately includes
a dual-emitting conflict. Both request spellings must address one logical field,
with the current spelling winning when a row contains both.
"""
from collections.abc import Callable, Generator
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
from typing import Any
import pytest
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.querier import Aggregation, BuilderQuery, OrderBy, RequestType, TelemetryFieldKey, make_query_request
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
CURRENT = "deployment.environment.name"
OLD = "deployment.environment"
PREFIX = "semconv-phase1"
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 _span(timestamp: datetime, suffix: str, environment: dict[str, str]) -> Traces:
service = f"{PREFIX}-{suffix}"
return 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),
)
@pytest.fixture(name="semconv_phase1_data")
def semconv_phase1_data(
insert_traces: Callable[[list[Traces]], None],
clickhouse: types.TestContainerClickhouse,
) -> Generator[datetime]:
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=2)
insert_traces(
[
_span(now - timedelta(seconds=5), "old", {OLD: "production"}),
_span(now - timedelta(seconds=4), "current", {CURRENT: "production"}),
_span(now - timedelta(seconds=3), "both", {OLD: "production", CURRENT: "production"}),
_span(now - timedelta(seconds=2), "conflict", {OLD: "staging", CURRENT: "production"}),
_span(now - timedelta(seconds=1), "staging", {OLD: "staging"}),
_span(now, "missing", {}),
]
)
# Service-map rows are derived by the collector in production. Seed the
# derived table directly here so the backend alias allowlist is tested in
# isolation; 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
'{PREFIX}-map-{suffix}', '{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}' "
f"DELETE WHERE startsWith(src, '{PREFIX}-map-') SETTINGS mutations_sync = 1"
)
def _result(response: requests.Response) -> dict[str, Any]:
assert response.status_code == HTTPStatus.OK, response.text
results = response.json()["data"]["data"]["results"]
assert len(results) == 1
return results[0]
def _raw_names(
signoz: types.SigNoz,
token: str,
now: datetime,
expression: str,
) -> set[str]:
response = make_query_request(
signoz,
token,
start_ms=int((now - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((now + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
BuilderQuery(
signal="traces",
name="A",
limit=100,
filter_expression=expression,
select_fields=[TelemetryFieldKey("span.name")],
order=[OrderBy(TelemetryFieldKey("timestamp"), "asc")],
).to_dict()
],
)
return {row["data"]["name"] for row in (_result(response).get("rows") or [])}
def _metadata_values(signoz: types.SigNoz, token: str, name: str, context: str) -> set[str]:
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/fields/values"),
timeout=5,
headers={"authorization": f"Bearer {token}"},
params={
"signal": "traces",
"name": name,
"fieldContext": context,
"fieldDataType": "string",
},
)
assert response.status_code == HTTPStatus.OK, response.text
return set(response.json()["data"]["values"].get("stringValues") or [])
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)
# 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}"
assert _raw_names(signoz, token, now, f"{field} = 'production'") == PRODUCTION_SPANS
assert _raw_names(signoz, token, now, f"{field} = 'staging'") == STAGING_SPANS
assert _raw_names(signoz, token, now, f"{field} != 'production'") == STAGING_SPANS
assert _raw_names(signoz, token, now, f"{field} EXISTS") == PRODUCTION_SPANS | STAGING_SPANS
assert _raw_names(signoz, token, now, f"{field} NOT EXISTS") == MISSING_SPANS
grouped = make_query_request(
signoz,
token,
start_ms=int((now - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((now + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
BuilderQuery(
signal="traces",
name="A",
filter_expression=f"{field} EXISTS",
aggregations=[Aggregation("count()")],
group_by=[TelemetryFieldKey(requested, "string", context)],
order=[OrderBy(TelemetryFieldKey(requested, "string", context), "asc")],
).to_dict()
],
)
result = _result(grouped)
assert result["columns"][0]["name"] == requested, "response identity must match the request spelling"
assert result["data"] == [["production", 4], ["staging", 1]]
assert _metadata_values(signoz, token, CURRENT, context) == {"production", "staging"}
assert _metadata_values(signoz, token, OLD, context) == {"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 not 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"}

Some files were not shown because too many files have changed in this diff Show More