Compare commits

...

3 Commits

Author SHA1 Message Date
Tushar Vats
7663ab02d4 feat(logparsingpipeline): add json_body_dual_ingestion flag for dual body ingestion (#12831)
#### Description

Server counterpart to SigNoz/signoz-otel-collector#891, which makes the
collector write each log body to both the legacy `body` column and
`body_v2` while a `json_body_dual_ingestion` flag is on.

- Adds the `json_body_dual_ingestion` feature flag (experimental, off by
default).
- Normalize is placed by read mode, so user pipelines always see the
body the explorer shows. With `use_json_body` on it stays ahead of user
pipelines. With dual ingestion alone, reads are still on the legacy
`body`, so it runs after them and only feeds `body_v2`.
- Whenever dual ingestion is on, the operator carries
`json_body_dual_ingestion: true` so it stashes the original body for the
exporter to restore. Running last under dual makes that stash the
post-pipeline body, exactly what legacy ingestion stores today.
- The pipeline preview follows the same rule and normalizes only under
`use_json_body`.

Design notes: [Normalize Operator and
Pipelines](https://app.notion.com/p/signoz/Normalize-Operator-and-Pipelines-3d7fcc6bcd19802396b9e8e817e30381),
[Dual JSON body
ingestion](https://app.notion.com/p/signoz/Dual-JSON-body-ingestion-3c6fcc6bcd198049963bce6ce2648af8)

#### Additional Information

- Collectors must run a build containing
SigNoz/signoz-otel-collector#891 before this flag is turned on. Verified
locally against v0.144.9: an older collector does not reject the unknown
operator key, it silently ignores it (operator configs are decoded with
`confmap.WithIgnoreUnused()`), runs normalize, and writes the normalized
body into the legacy `body` column until it is upgraded. No collector
release includes #891 yet.
- The exporter's own `json_body_dual_ingestion` key is collector deploy
config and is flipped together with this flag; the server does not set
it.
- Toggling either flag does not bump the pipeline config version, so
connected agents need a new pipeline save to pick up the operator or its
position. Same caveat as `use_json_body` today.
- Under dual, normalize is last among the SigNoz pipelines; custom
collector processors placed after them still see the normalized map.

🤖 Generated with [Claude Code](https://claude.com/claude-code)
2026-10-02 10:56:38 +00:00
Srikanth Chekuri
1643df620b feat: resolve semconv families across logs and metrics (#12870)
Some checks failed
build-staging / prepare (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
#### Description

Phase 2 of #6143: semantic-convention families resolve on logs and
metrics, behind the `resolve_semconv_families` flag (default off), on
the storage contract of #12802.

- Registry: each member carries the scope of its own rename edges, so a
fan-out keeps one membership per target and an ambiguous name stays
literal.
- Logs and metrics families need no family code of their own. The gate
applies to every signal, and `LogicalRead` merges members through each
storage's `Read`.
- Metric-name families union the storage names in every `metric_name`
filter, and the querier reads type, temporality, and the reduced flag
across the family.
- Span-metrics labels: the metrics the processor emits, listed by name,
also read each family member with the `resource_` prefix. Requested
names are never rewritten.
- Values suggestions and related values cover every spelling of the
family.
- `deployment.environment.name` resolves on all three signals.
`db.system.name` stays off until a value-mapping reader exists.

#### Additional Information

- `pkg/semconv.Family` fields are now unexported, and `transition.go` is
removed. #12446 reads the old API and needs an update when stacked.
- A target that emits both names of a metric-name family double-counts
in `sum()` during the overlap window. Reading both names is the feature.
Pinned by a test.
- A family of metrics labels keeps the keyless contract of a single
label: no guard, no NULL group.
2026-10-01 20:28:15 +00:00
Aditya Singh
fd21f8b760 feat(bottom-strip): show the page count on the left for noz, alerts and home (#12963)
#### Description

- noz page, alert rules, home, exceptions and services now push their
count to the left of the strip. each page has its own `useXStripInfo`
hook that builds the config and publishes it.. same pattern trace
details already uses, five more times.
- the hook does the whole thing, so the page call is one line. noz and
home fetch or subscribe inside the hook so the page does not re-render
just to keep the strip current.. the rest pass values they already hold.
- alert rules and exceptions show "N of M".. first number is what is on
the page right now. both read the same two values the table hands its
own pagination, so the strip cannot disagree with the table. services,
home and noz are a plain count.
- services renders one of two tables on `use_span_metrics`, so the hook
is called from both.. the count text lives in one place either way.
- dashboards is left out for now, it already shows this on its own
strip.

#### Issues closed by this PR

Part of https://github.com/SigNoz/events-pod/issues/53
Part of https://github.com/SigNoz/events-pod/issues/55


#### Screenshots

Ai Assistant  Bottom strip

<img width="3456" height="1968" alt="image"
src="https://github.com/user-attachments/assets/0042f71b-1b56-417a-8283-af9ba9596351"
/>


Alert rules

<img width="3452" height="1992" alt="image"
src="https://github.com/user-attachments/assets/a674de13-9841-4717-89ce-a656c9df646c"
/>

Exceptions 

<img width="3456" height="1970" alt="image"
src="https://github.com/user-attachments/assets/6333ea6b-2342-46dc-985d-2855bebe4003"
/>



Services

<img width="3456" height="1996" alt="image"
src="https://github.com/user-attachments/assets/abe0c5cb-cfe3-48b6-9591-ee71e677a1a6"
/>




#### Additional Information
2026-10-01 19:30:00 +00:00
64 changed files with 3507 additions and 574 deletions

View File

@@ -67,7 +67,7 @@ jobs:
with:
go-version: "1.24"
- name: check-semconv-generated-files
run: go run ./scripts/semconv -check
run: make semconv-check
build:
if: |
github.event_name == 'merge_group' ||

View File

@@ -262,6 +262,10 @@ py-clean: ## Clear all pycache and pytest cache from tests directory recursively
semconv-generate: ## Regenerate semantic-convention families for Go and TypeScript
@go run ./scripts/semconv
.PHONY: semconv-check
semconv-check: ## Fail if the generated semantic-convention files are stale
@go run ./scripts/semconv -check
.PHONY: gen-mocks
gen-mocks:
@echo ">> Generating mocks"

View File

@@ -84,9 +84,9 @@ A storage answers four questions and nothing else:
| WhenAbsent | Absent row reads | Positive filter | Raw select | Multi-candidate column | Field keys |
|---|---|---|---|---|---|
| `AlwaysPresent` | a real value | no guard | no guard | no branch, ends the candidate list | table columns |
| `AbsentIsSentinel` | `''`, 0, false, and that is not a value | exists guard | exists guard | presence branch | map attributes, cast JSON paths, string families |
| `AbsentIsSentinel` | `''`, 0, false, and that is not a value | exists guard | exists guard | presence branch | map attributes, cast JSON paths, string families of such members |
| `AbsentIsNull` | NULL | no guard | no guard | presence branch | multi-era folds, body JSON paths, numeric families |
| `AbsentIsValue` | `''`, and that is the keyless contract | no guard | no guard | no presence branch | metrics labels, rule state history labels |
| `AbsentIsValue` | `''`, and that is the keyless contract | no guard | no guard | no presence branch | metrics labels, rule state history labels, and families of such members |
### The generic layer
@@ -109,7 +109,7 @@ The functions, from the outside in:
| `RejectsBodyFunction(traits, operator)` | Runs before resolution. A storage without body functions (`has`, `hasAny`, `hasAll`, `hasToken`, `search`) errors. The fingerprint side of a split skips the term, because the main query evaluates it. After resolution, `Condition` errors when `has`, `hasAny`, `hasAll`, or `hasToken` lands on a map-backed key (resource, attribute, scope), before the split can drop it. |
| `SharedCondition(...)` | The `Compile` of every storage without its own condition language: `LogicalRead`, the shared data-type collision cast, `OperatorCondition`, then the guard rule. |
| `OperatorCondition(...)` | The operator switch over an already cast read. A storage with its own cast policy composes with it. |
| `LogicalRead(...)` | The only place family expressions are built. A single-member field reads through its member. A family merges the member reads, current member first: `COALESCE(NULLIF(m1, ''), NULLIF(m2, ''), '')` for strings, `multiIf` with a NULL tail for numbers. It ORs the member presence tests. A row without any member reads what the tail of the merge reads. A member with a value map reads through `TransformRead`. `NOT EXISTS` is the read's `Absence`, the storage's own negated form. |
| `LogicalRead(...)` | The only place family expressions are built. A single-member field reads through its member. A family merges the member reads, current member first: `COALESCE(NULLIF(m1, ''), NULLIF(m2, ''), '')` for strings, `multiIf` with a NULL tail for numbers. It ORs the member presence tests. A row without any member reads what the tail of the merge reads. When every member reads its sentinel as a value, so does the family. A member with a value map reads through `TransformRead`. `NOT EXISTS` is the read's `Absence`, the storage's own negated form. |
### A resolved key

View File

@@ -1,32 +1,82 @@
// Code generated by scripts/semconv. DO NOT EDIT.
export type SemconvFamily = {
readonly current: string;
readonly old: readonly string[];
readonly kind: 'attribute' | 'metric';
// An empty contexts/signals/applyToMetrics array places no constraint on
// that axis.
export type SemconvMember = {
readonly name: string;
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly applyToMetrics: readonly string[];
};
export type SemconvFamily = {
readonly current: string;
readonly kind: 'attribute' | 'metric';
readonly members: readonly SemconvMember[];
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly valueMap: Readonly<Record<string, string>>;
};
export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
{
current: 'db.system.name',
old: ['db.system'],
kind: 'attribute',
current: 'container.cpu.usage',
kind: 'metric',
members: [
{
name: 'container.cpu.utilization',
contexts: [],
signals: [],
applyToMetrics: [],
},
],
contexts: [],
signals: [],
applyToMetrics: [],
valueMap: {},
},
{
current: 'deployment.environment.name',
old: ['deployment.environment'],
kind: 'attribute',
members: [
{
name: 'deployment.environment',
contexts: [],
signals: [],
applyToMetrics: [],
},
],
contexts: ['attribute', 'resource'],
signals: ['logs', 'metrics', 'traces'],
valueMap: {},
},
{
current: 'k8s.node.cpu.usage',
kind: 'metric',
members: [
{
name: 'k8s.node.cpu.utilization',
contexts: [],
signals: [],
applyToMetrics: [],
},
],
contexts: [],
signals: [],
valueMap: {},
},
{
current: 'k8s.pod.cpu.usage',
kind: 'metric',
members: [
{
name: 'k8s.pod.cpu.utilization',
contexts: [],
signals: [],
applyToMetrics: [],
},
],
contexts: [],
signals: [],
applyToMetrics: [],
valueMap: {},
},
] as const;

View File

@@ -0,0 +1,53 @@
import { renderHook } from '@testing-library/react';
import { useBottomStripStore } from 'container/BottomStrip/store/useBottomStripStore';
import { StripItemKind } from 'container/BottomStrip/types';
import { useExceptionsStripInfo } from '../useExceptionsStripInfo';
describe('useExceptionsStripInfo', () => {
beforeEach(() => {
useBottomStripStore.setState({ left: null, ownerId: null });
});
it('shows the rows on the page against the total', () => {
renderHook(() => useExceptionsStripInfo({ shownCount: 25, totalCount: 500 }));
expect(useBottomStripStore.getState().left).toMatchObject([
{ kind: StripItemKind.Text, text: '25 of 500 exceptions' },
]);
});
it('still says n of m when the whole list fits on one page', () => {
renderHook(() => useExceptionsStripInfo({ shownCount: 42, totalCount: 42 }));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '42 of 42 exceptions',
});
});
it('says exception, not exceptions, when there is one', () => {
renderHook(() => useExceptionsStripInfo({ shownCount: 1, totalCount: 1 }));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '1 of 1 exception',
});
});
it('shows zero before the counts land', () => {
renderHook(() => useExceptionsStripInfo({ shownCount: 0, totalCount: 0 }));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '0 of 0 exceptions',
});
});
it('clears the strip when the page unmounts', () => {
const { unmount } = renderHook(() =>
useExceptionsStripInfo({ shownCount: 25, totalCount: 500 }),
);
unmount();
expect(useBottomStripStore.getState().left).toBeNull();
});
});

View File

@@ -38,6 +38,7 @@ import { Exception, PayloadProps } from 'types/api/errors/getAll';
import { GlobalReducer } from 'types/reducer/globalTime';
import { FilterDropdownExtendsProps } from './types';
import { useExceptionsStripInfo } from './useExceptionsStripInfo';
import {
extractFilterValues,
getDefaultFilterValue,
@@ -160,6 +161,11 @@ function AllErrors(): JSX.Element {
},
]);
useExceptionsStripInfo({
shownCount: data?.payload?.length ?? 0,
totalCount: errorCountResponse.data?.payload ?? 0,
});
const isFetching = isErrorsFetching || errorCountResponse.isFetching;
useEffect(() => {
setIsFetching(isFetching);

View File

@@ -0,0 +1,26 @@
import { useMemo } from 'react';
import { useBottomStrip } from 'container/BottomStrip/useBottomStrip';
import { type StripItem, StripItemKind } from 'container/BottomStrip/types';
import { pluralize } from 'utils/pluralize';
interface UseExceptionsStripInfoArgs {
shownCount: number;
totalCount: number;
}
export function useExceptionsStripInfo({
shownCount,
totalCount,
}: UseExceptionsStripInfoArgs): void {
const items = useMemo<StripItem[]>(
() => [
{
kind: StripItemKind.Text,
text: `${shownCount} of ${pluralize(totalCount, 'exception')}`,
},
],
[shownCount, totalCount],
);
useBottomStrip(items);
}

View File

@@ -854,7 +854,9 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
<ChangelogModal changelog={changelog} onClose={toggleChangelogModal} />
)}
<Toaster />
<Toaster
offset={{ bottom: 'calc(var(--bottom-strip-height, 0px) + 24px)' }}
/>
</Layout>
</TooltipProvider>
);

View File

@@ -26,6 +26,7 @@ import { initialQueriesMap, PANEL_TYPES } from 'constants/queryBuilder';
import { REACT_QUERY_KEY } from 'constants/reactQueryKeys';
import ROUTES from 'constants/routes';
import { DEFAULT_TIME_RANGE } from 'container/TopNav/DateTimeSelectionV2/constants';
import { useHomeStripInfo } from 'container/Home/useHomeStripInfo';
import { useGetQueryRange } from 'hooks/queryBuilder/useGetQueryRange';
import { useIsDarkMode } from 'hooks/useDarkMode';
import { useSafeNavigate } from 'hooks/useSafeNavigate';
@@ -64,6 +65,8 @@ const homeInterval = 30 * 60 * 1000;
// eslint-disable-next-line sonarjs/cognitive-complexity
export default function Home(): JSX.Element {
useHomeStripInfo();
const { user } = useAppContext();
const { safeNavigate } = useSafeNavigate();
const isDarkMode = useIsDarkMode();

View File

@@ -0,0 +1,58 @@
import { renderHook } from '@testing-library/react';
import { useGetAlerts } from 'api/generated/services/alerts';
import { useBottomStripStore } from 'container/BottomStrip/store/useBottomStripStore';
import { StripItemKind } from 'container/BottomStrip/types';
import { useHomeStripInfo } from '../useHomeStripInfo';
jest.mock('api/generated/services/alerts', () => ({
useGetAlerts: jest.fn(),
}));
const mockUseGetAlerts = useGetAlerts as jest.Mock;
describe('useHomeStripInfo', () => {
beforeEach(() => {
useBottomStripStore.setState({ left: null, ownerId: null });
});
it('counts the firing alert instances', () => {
mockUseGetAlerts.mockReturnValue({ data: { data: [{}, {}, {}] } });
renderHook(() => useHomeStripInfo());
expect(useBottomStripStore.getState().left).toMatchObject([
{ kind: StripItemKind.Text, text: '3 alerts firing' },
]);
});
it('says alert, not alerts, when only one is firing', () => {
mockUseGetAlerts.mockReturnValue({ data: { data: [{}] } });
renderHook(() => useHomeStripInfo());
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '1 alert firing',
});
});
it('shows zero before the response lands', () => {
mockUseGetAlerts.mockReturnValue({ data: undefined });
renderHook(() => useHomeStripInfo());
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '0 alerts firing',
});
});
it('clears the strip when the page unmounts', () => {
mockUseGetAlerts.mockReturnValue({ data: { data: [{}, {}, {}] } });
const { unmount } = renderHook(() => useHomeStripInfo());
unmount();
expect(useBottomStripStore.getState().left).toBeNull();
});
});

View File

@@ -0,0 +1,21 @@
import { useMemo } from 'react';
import { useGetAlerts } from 'api/generated/services/alerts';
import { useBottomStrip } from 'container/BottomStrip/useBottomStrip';
import { type StripItem, StripItemKind } from 'container/BottomStrip/types';
import { pluralize } from 'utils/pluralize';
export function useHomeStripInfo(): void {
// Firing instances, not rules, matching the triggered alerts page.
const { data } = useGetAlerts();
const count = data?.data?.length ?? 0;
const items = useMemo<StripItem[]>(
() => [
{ kind: StripItemKind.Text, text: `${pluralize(count, 'alert')} firing` },
],
[count],
);
useBottomStrip(items);
}

View File

@@ -0,0 +1,53 @@
import { renderHook } from '@testing-library/react';
import { useBottomStripStore } from 'container/BottomStrip/store/useBottomStripStore';
import { StripItemKind } from 'container/BottomStrip/types';
import { useAlertRulesStripInfo } from '../useAlertRulesStripInfo';
describe('useAlertRulesStripInfo', () => {
beforeEach(() => {
useBottomStripStore.setState({ left: null, ownerId: null });
});
it('shows the rows on the page against the total', () => {
renderHook(() => useAlertRulesStripInfo({ shownCount: 15, totalCount: 17 }));
expect(useBottomStripStore.getState().left).toMatchObject([
{ kind: StripItemKind.Text, text: '15 of 17 rules' },
]);
});
it('still says n of m when the whole list fits on one page', () => {
renderHook(() => useAlertRulesStripInfo({ shownCount: 17, totalCount: 17 }));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '17 of 17 rules',
});
});
it('says rule, not rules, when there is only one', () => {
renderHook(() => useAlertRulesStripInfo({ shownCount: 1, totalCount: 1 }));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '1 of 1 rule',
});
});
it('shows zero when nothing matched', () => {
renderHook(() => useAlertRulesStripInfo({ shownCount: 0, totalCount: 12 }));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '0 of 12 rules',
});
});
it('clears the strip when the page unmounts', () => {
const { unmount } = renderHook(() =>
useAlertRulesStripInfo({ shownCount: 15, totalCount: 17 }),
);
unmount();
expect(useBottomStripStore.getState().left).toBeNull();
});
});

View File

@@ -20,6 +20,7 @@ import { ALERT_RULES_PARAMS, useAlertRulesFilters } from './hooks';
import styles from './ListAlertRules.module.scss';
import { getAlertRuleColumns } from './table.config';
import type { AlertRule } from './types';
import { useAlertRulesStripInfo } from './useAlertRulesStripInfo';
import { useAlertRulesData } from './useAlertRulesData';
import { useAlertRulesHandlers } from './useAlertRulesHandlers';
@@ -87,6 +88,11 @@ function ListAlertRules(): JSX.Element {
return filteredRules.slice(start, start + limit);
}, [filteredRules, page, limit]);
useAlertRulesStripInfo({
shownCount: paginatedRules.length,
totalCount: filteredRules.length,
});
const columnsWithActions = useMemo(() => {
if (!action) {
return columns;

View File

@@ -0,0 +1,27 @@
import { useMemo } from 'react';
import { useBottomStrip } from 'container/BottomStrip/useBottomStrip';
import { type StripItem, StripItemKind } from 'container/BottomStrip/types';
import { pluralize } from 'utils/pluralize';
interface UseAlertRulesStripInfoArgs {
/** Rows on the current page, matching the table's own footer. */
shownCount: number;
totalCount: number;
}
export function useAlertRulesStripInfo({
shownCount,
totalCount,
}: UseAlertRulesStripInfoArgs): void {
const items = useMemo<StripItem[]>(
() => [
{
kind: StripItemKind.Text,
text: `${shownCount} of ${pluralize(totalCount, 'rule')}`,
},
],
[shownCount, totalCount],
);
useBottomStrip(items);
}

View File

@@ -21,6 +21,7 @@ import { getTotalRPS } from 'utils/services';
import { getColumns } from '../Columns/ServiceColumn';
import { ServiceMetricsTableProps } from '../types';
import { useServicesStripInfo } from '../useServicesStripInfo';
import { getServiceListFromQuery } from '../utils';
function ServiceMetricTable({
@@ -67,6 +68,8 @@ function ServiceMetricTable({
[isLoading, queries, topLevelOperations],
);
useServicesStripInfo(services.length);
const { search } = useLocation();
const tableColumns = useMemo(() => getColumns(search, true), [search]);
const [RPS, setRPS] = useState(0);

View File

@@ -18,6 +18,7 @@ import { GlobalReducer } from 'types/reducer/globalTime';
import { Tags } from 'hooks/useResourceAttribute/types';
import SkipOnBoardingModal from '../SkipOnBoardModal';
import { useServicesStripInfo } from '../useServicesStripInfo';
import ServiceTraceTable from './ServiceTracesTable';
function ServiceTraces(): JSX.Element {
@@ -42,6 +43,8 @@ function ServiceTraces(): JSX.Element {
const services = data || [];
useServicesStripInfo(services.length);
const [skipOnboarding, setSkipOnboarding] = useState(
localStorageGet(SKIP_ONBOARDING) === 'true',
);

View File

@@ -0,0 +1,43 @@
import { renderHook } from '@testing-library/react';
import { useBottomStripStore } from 'container/BottomStrip/store/useBottomStripStore';
import { StripItemKind } from 'container/BottomStrip/types';
import { useServicesStripInfo } from '../useServicesStripInfo';
describe('useServicesStripInfo', () => {
beforeEach(() => {
useBottomStripStore.setState({ left: null, ownerId: null });
});
it('shows how many services are listed', () => {
renderHook(() => useServicesStripInfo(18));
expect(useBottomStripStore.getState().left).toMatchObject([
{ kind: StripItemKind.Text, text: '18 services' },
]);
});
it('says service, not services, when there is one', () => {
renderHook(() => useServicesStripInfo(1));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '1 service',
});
});
it('shows zero when there are none', () => {
renderHook(() => useServicesStripInfo(0));
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '0 services',
});
});
it('clears the strip when the page unmounts', () => {
const { unmount } = renderHook(() => useServicesStripInfo(18));
unmount();
expect(useBottomStripStore.getState().left).toBeNull();
});
});

View File

@@ -0,0 +1,13 @@
import { useMemo } from 'react';
import { useBottomStrip } from 'container/BottomStrip/useBottomStrip';
import { type StripItem, StripItemKind } from 'container/BottomStrip/types';
import { pluralize } from 'utils/pluralize';
export function useServicesStripInfo(count: number): void {
const items = useMemo<StripItem[]>(
() => [{ kind: StripItemKind.Text, text: pluralize(count, 'service') }],
[count],
);
useBottomStrip(items);
}

View File

@@ -11,6 +11,8 @@ import { useAIAssistantStore } from 'container/AIAssistant/store/useAIAssistantS
import { VariantContext } from 'container/AIAssistant/VariantContext';
import Noz from 'components/Noz/Noz';
import { useAIAssistantStripInfo } from './useAIAssistantStripInfo';
import styles from './AIAssistantPage.module.scss';
import ConversationsList from 'container/AIAssistant/components/ConversationsList';
@@ -41,6 +43,8 @@ export default function AIAssistantPage(): JSX.Element {
// eslint-disable-next-line react-hooks/exhaustive-deps
}, []);
useAIAssistantStripInfo();
const conversations = useAIAssistantStore((s) => s.conversations);
const activeConversationId = useAIAssistantStore(
(s) => s.activeConversationId,

View File

@@ -0,0 +1,60 @@
import { renderHook } from '@testing-library/react';
import { useAIAssistantStore } from 'container/AIAssistant/store/useAIAssistantStore';
import { useBottomStripStore } from 'container/BottomStrip/store/useBottomStripStore';
import { StripItemKind } from 'container/BottomStrip/types';
import { useAIAssistantStripInfo } from '../useAIAssistantStripInfo';
function seed(conversations: Record<string, unknown>): void {
useAIAssistantStore.setState({ conversations } as never);
}
describe('useAIAssistantStripInfo', () => {
beforeEach(() => {
useBottomStripStore.setState({ left: null, ownerId: null });
});
it('counts only the conversations that are not archived', () => {
seed({
a: { id: 'a' },
b: { id: 'b' },
c: { id: 'c', archived: true },
});
renderHook(() => useAIAssistantStripInfo());
expect(useBottomStripStore.getState().left).toMatchObject([
{ kind: StripItemKind.Text, text: '2 conversations' },
]);
});
it('says one conversation, not 1 conversations', () => {
seed({ a: { id: 'a' } });
renderHook(() => useAIAssistantStripInfo());
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '1 conversation',
});
});
it('shows zero when there are none', () => {
seed({});
renderHook(() => useAIAssistantStripInfo());
expect(useBottomStripStore.getState().left?.[0]).toMatchObject({
text: '0 conversations',
});
});
it('clears the strip when the page unmounts', () => {
seed({ a: { id: 'a' } });
const { unmount } = renderHook(() => useAIAssistantStripInfo());
unmount();
expect(useBottomStripStore.getState().left).toBeNull();
});
});

View File

@@ -243,3 +243,11 @@ export const TooltipsInApprovalDiff: Story = {
args: { tooltipsOpen: true, agent: 'awaiting-approval', contents: BRIEF },
play: openApprovalDiff,
};
/**
* The conversation count in the bottom strip, in place of the build version.
* Archived threads are left out of it.
*/
export const BottomStrip: Story = {
args: { bottomStrip: true },
};

View File

@@ -0,0 +1,20 @@
import { useMemo } from 'react';
import { useAIAssistantStore } from 'container/AIAssistant/store/useAIAssistantStore';
import { useBottomStrip } from 'container/BottomStrip/useBottomStrip';
import { type StripItem, StripItemKind } from 'container/BottomStrip/types';
import { pluralize } from 'utils/pluralize';
export function useAIAssistantStripInfo(): void {
const conversations = useAIAssistantStore((state) => state.conversations);
const count = Object.values(conversations).filter(
(conversation) => !conversation.archived,
).length;
const items = useMemo<StripItem[]>(
() => [{ kind: StripItemKind.Text, text: pluralize(count, 'conversation') }],
[count],
);
useBottomStrip(items);
}

View File

@@ -133,3 +133,19 @@ export const ColumnPicker: Story = {
export const Tooltips: Story = {
args: { tooltipsOpen: true },
};
/**
* The rule count in the bottom strip, in place of the build version: the rows on
* the page against the total, the same pair the table's own footer prints.
*/
export const BottomStrip: Story = {
args: { bottomStrip: true },
};
/** The same count on a second page, where the two numbers come apart. */
export const BottomStripPaginated: Story = {
args: { bottomStrip: true, rules: RULE_MAX },
parameters: {
signoz: { route: '/alerts?tab=AlertRules&page=2&limit=10' },
},
};

View File

@@ -125,3 +125,11 @@ export const QuickFiltersSettingsWithBanner: Story = {
args: { banner: 'trial-expiry' },
play: dirtyQuickFiltersSettings,
};
/**
* The exception count in the bottom strip, in place of the build version: the
* rows on the page against the total the count query returns.
*/
export const BottomStrip: Story = {
args: { bottomStrip: true },
};

View File

@@ -97,3 +97,11 @@ export const NavSettingsMenu: Story = {
await screen.findByRole('menu');
},
};
/**
* The firing alert count in the bottom strip, in place of the build version. It
* counts firing instances, so it does not match the alert rules widget above.
*/
export const BottomStrip: Story = {
args: { bottomStrip: true },
};

View File

@@ -61,3 +61,13 @@ export const Pagination: Story = {
export const OverTrialLimit: Story = {
args: { traffic: 'over-trial-limit', banner: 'trial-expiry' },
};
/** The service count in the bottom strip, in place of the build version. */
export const BottomStrip: Story = {
args: { bottomStrip: true },
};
/** The same count on the span metrics table, which is the page's other path. */
export const BottomStripSpanMetrics: Story = {
args: { bottomStrip: true, mode: 'span-metrics' },
};

View File

@@ -1,5 +1,6 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import ROUTES from 'constants/routes';
import { toast } from '@signozhq/ui/sonner';
import { screen, userEvent, waitFor } from 'storybook/test';
import type { GlobalMockArgs } from '../globals';
@@ -108,3 +109,24 @@ export const WithoutStripAddCardModal: Story = {
});
},
};
/** Raises a toast and waits for it, so the shot is taken with it on screen. */
const raiseToast: NonNullable<Story['play']> = async () => {
toast.success('Service account deleted');
await screen.findByText('Service account deleted', undefined, untilLoaded);
};
/**
* A toast with the strip on. Sonner pins itself to the foot of the viewport, so
* without an offset it lands on top of Ask Noz and Support.
*/
export const WithToast: Story = {
args: { noz: true, support: 'pylon' },
play: raiseToast,
};
/** The same toast with the flag off, back at sonner's own inset from the viewport. */
export const WithoutStripToast: Story = {
args: { bottomStrip: false, noz: true, support: 'pylon' },
play: raiseToast,
};

View File

@@ -9,6 +9,7 @@ var (
FeaturePutMetersInZeus = featuretypes.MustNewName("put_meters_in_zeus")
FeatureUseMeterReporter = featuretypes.MustNewName("use_meter_reporter")
FeatureUseJSONBody = featuretypes.MustNewName("use_json_body")
FeatureJSONBodyDualIngestion = featuretypes.MustNewName("json_body_dual_ingestion")
FeatureEnableMetricsReduction = featuretypes.MustNewName("enable_metrics_reduction")
FeatureResolveSemconvFamilies = featuretypes.MustNewName("resolve_semconv_families")
FeatureUseTraceAttributesJSON = featuretypes.MustNewName("use_trace_attributes_json")
@@ -64,6 +65,14 @@ func MustNewRegistry() featuretypes.Registry {
DefaultVariant: featuretypes.MustNewName("disabled"),
Variants: featuretypes.NewBooleanVariants(),
},
&featuretypes.Feature{
Name: FeatureJSONBodyDualIngestion,
Kind: featuretypes.KindBoolean,
Stage: featuretypes.StageExperimental,
Description: "Controls whether the collector's normalize operator keeps the original log body so it is ingested into both the legacy body and the JSON body columns",
DefaultVariant: featuretypes.MustNewName("disabled"),
Variants: featuretypes.NewBooleanVariants(),
},
&featuretypes.Feature{
Name: FeatureEnableMetricsReduction,
Kind: featuretypes.KindBoolean,
@@ -76,7 +85,7 @@ func MustNewRegistry() featuretypes.Registry {
Name: FeatureResolveSemconvFamilies,
Kind: featuretypes.KindBoolean,
Stage: featuretypes.StageExperimental,
Description: "Controls whether trace queries resolve a semantic-convention name to all the spellings of its family",
Description: "Controls whether trace, log, and metric queries resolve a semantic-convention name to all the spellings of its family",
DefaultVariant: featuretypes.MustNewName("disabled"),
Variants: featuretypes.NewBooleanVariants(),
},

View File

@@ -435,6 +435,7 @@ func (m *module) buildFilterClause(ctx context.Context, orgID valuer.UUID, filte
whereClauseSelectors[idx].SelectorMatchType = telemetrytypes.FieldSelectorMatchTypeExact
}
whereClauseSelectors = querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, m.fl, whereClauseSelectors)
keys, _, err := m.telemetryMetadataStore.GetKeysMulti(ctx, orgID, whereClauseSelectors)
if err != nil {
return nil, err

View File

@@ -936,6 +936,7 @@ func (m *module) buildFilterClause(ctx context.Context, orgID valuer.UUID, filte
// whereClauseSelectors[idx].Source = query.Source
}
whereClauseSelectors = querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, m.fl, whereClauseSelectors)
keys, _, err := m.telemetryMetadataStore.GetKeysMulti(ctx, orgID, whereClauseSelectors)
if err != nil {
return nil, err

View File

@@ -0,0 +1,102 @@
package querier
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/types/metrictypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// The metric metadata of a query on one name of a metric-name family comes
// from every name of the family: the temporality is Multiple when the names
// differ, and the reduced flag is set when any name has reduced data.
func TestResolveMetricMetadataReadsTheFamily(t *testing.T) {
testCases := []struct {
name string
temporalities map[string]metrictypes.Temporality
reduced map[string]bool
expectedTemporality metrictypes.Temporality
expectedReduced bool
}{
{
name: "SameTemporality_KeepsIt",
temporalities: map[string]metrictypes.Temporality{
"k8s.pod.cpu.usage": metrictypes.Cumulative,
"k8s.pod.cpu.utilization": metrictypes.Cumulative,
},
expectedTemporality: metrictypes.Cumulative,
},
{
name: "DifferentTemporalities_ReadAsMultiple",
temporalities: map[string]metrictypes.Temporality{
"k8s.pod.cpu.usage": metrictypes.Delta,
"k8s.pod.cpu.utilization": metrictypes.Cumulative,
},
expectedTemporality: metrictypes.Multiple,
},
{
name: "OnlyOldNameKnown_TakesItsTemporality",
temporalities: map[string]metrictypes.Temporality{
"k8s.pod.cpu.utilization": metrictypes.Delta,
},
expectedTemporality: metrictypes.Delta,
},
{
name: "OldNameReduced_MarksTheAggregationReduced",
temporalities: map[string]metrictypes.Temporality{
"k8s.pod.cpu.usage": metrictypes.Cumulative,
},
reduced: map[string]bool{"k8s.pod.cpu.utilization": true},
expectedTemporality: metrictypes.Cumulative,
expectedReduced: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
metadataStore := telemetrytypestest.NewMockMetadataStore()
metadataStore.TemporalityMap = testCase.temporalities
metadataStore.TypeMap = map[string]metrictypes.Type{
"k8s.pod.cpu.usage": metrictypes.GaugeType,
"k8s.pod.cpu.utilization": metrictypes.GaugeType,
}
metadataStore.ReducedMap = testCase.reduced
q := &querier{
logger: instrumentationtest.New().Logger(),
metadataStore: metadataStore,
fl: flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
}),
}
queries := []qbtypes.QueryEnvelope{{
Type: qbtypes.QueryTypeBuilder,
Spec: qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]{
Name: "A",
Signal: telemetrytypes.SignalMetrics,
Aggregations: []qbtypes.MetricAggregation{{
MetricName: "k8s.pod.cpu.usage",
TimeAggregation: metrictypes.TimeAggregationAvg,
SpaceAggregation: metrictypes.SpaceAggregationAvg,
}},
},
}}
missing, warnings, err := q.resolveMetricMetadata(context.Background(), valuer.UUID{}, queries, 0, 0, qbtypes.RequestTypeTimeSeries)
require.NoError(t, err)
assert.Empty(t, missing)
assert.Empty(t, warnings)
spec := queries[0].Spec.(qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation])
assert.Equal(t, testCase.expectedTemporality, spec.Aggregations[0].Temporality)
assert.Equal(t, testCase.expectedReduced, spec.Aggregations[0].Reduced)
})
}
}

View File

@@ -369,6 +369,8 @@ func (q *querier) populateQBEvent(event *qbtypes.QBEvent, queries []qbtypes.Quer
// resolved: never-seen metrics and dormant metrics (seen but no data in
// the query window).
// - err: Internal when a metadata fetch fails.
//
// Metric metadata resolves through every name of a metric-name family.
func (q *querier) resolveMetricMetadata(ctx context.Context, orgID valuer.UUID, queries []qbtypes.QueryEnvelope, start, end uint64, requestType qbtypes.RequestType) (missingMetricQueries []string, metricWarnings []string, err error) {
metricNames := make([]string, 0)
for idx := range queries {
@@ -381,7 +383,7 @@ func (q *querier) resolveMetricMetadata(ctx context.Context, orgID valuer.UUID,
}
for _, agg := range spec.Aggregations {
if agg.MetricName != "" {
metricNames = append(metricNames, agg.MetricName)
metricNames = append(metricNames, querybuilder.FamilyMetricNames(ctx, orgID, q.fl, agg.MetricName)...)
}
}
}
@@ -409,14 +411,16 @@ func (q *querier) resolveMetricMetadata(ctx context.Context, orgID valuer.UUID,
presentAggregations := make([]qbtypes.MetricAggregation, 0, len(spec.Aggregations))
for i := range spec.Aggregations {
familyNames := querybuilder.FamilyMetricNames(ctx, orgID, q.fl, spec.Aggregations[i].MetricName)
if spec.Aggregations[i].MetricName != "" && spec.Aggregations[i].Temporality == metrictypes.Unknown {
if temp, ok := metricTemporality[spec.Aggregations[i].MetricName]; ok && temp != metrictypes.Unknown {
spec.Aggregations[i].Temporality = temp
}
spec.Aggregations[i].Temporality = familyTemporality(metricTemporality, familyNames)
}
if spec.Aggregations[i].MetricName != "" && spec.Aggregations[i].Type == metrictypes.UnspecifiedType {
if foundMetricType, ok := metricTypes[spec.Aggregations[i].MetricName]; ok && foundMetricType != metrictypes.UnspecifiedType {
spec.Aggregations[i].Type = foundMetricType
for _, member := range familyNames {
if foundMetricType, ok := metricTypes[member]; ok && foundMetricType != metrictypes.UnspecifiedType {
spec.Aggregations[i].Type = foundMetricType
break
}
}
}
if spec.Aggregations[i].Type == metrictypes.UnspecifiedType {
@@ -434,8 +438,11 @@ func (q *querier) resolveMetricMetadata(ctx context.Context, orgID valuer.UUID,
return nil, nil, err
}
}
if reducedMetricsSet[spec.Aggregations[i].MetricName] {
spec.Aggregations[i].Reduced = true
for _, member := range familyNames {
if reducedMetricsSet[member] {
spec.Aggregations[i].Reduced = true
break
}
}
presentAggregations = append(presentAggregations, spec.Aggregations[i])
}
@@ -505,6 +512,26 @@ func (q *querier) resolveMetricMetadata(ctx context.Context, orgID valuer.UUID,
return missingMetricQueries, warnings, nil
}
// familyTemporality is the temporality the family names share, or Multiple
// when they differ.
func familyTemporality(temporalities map[string]metrictypes.Temporality, names []string) metrictypes.Temporality {
found := metrictypes.Unknown
for _, name := range names {
temporality, ok := temporalities[name]
if !ok || temporality == metrictypes.Unknown {
continue
}
if found == metrictypes.Unknown {
found = temporality
continue
}
if found != temporality {
return metrictypes.Multiple
}
}
return found
}
func (q *querier) QueryRawStream(ctx context.Context, orgID valuer.UUID, req *qbtypes.QueryRangeRequest, client *qbtypes.RawStream) {
// Coerce the window to epoch milliseconds up front (End may be 0 for the

View File

@@ -52,10 +52,10 @@ import (
"github.com/SigNoz/signoz/pkg/query-service/constants"
chErrors "github.com/SigNoz/signoz/pkg/query-service/errors"
"github.com/SigNoz/signoz/pkg/query-service/metrics"
"github.com/SigNoz/signoz/pkg/query-service/model"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/query-service/utils"
"github.com/SigNoz/signoz/pkg/semconv"
)
const (
@@ -3202,7 +3202,14 @@ 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))
current := semconv.Current(semconv.KindMetric, telemetrytypes.FieldKeySelector{
Name: req.AggregateAttribute,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
})
if current != req.AggregateAttribute {
names = append(names, current)
}
rows, err = r.db.Query(ctx, query, req.FilterAttributeKey, names, req.FilterAttributeKey, fmt.Sprintf("%%%s%%", req.SearchText), common.PastDayRoundOff())

View File

@@ -220,7 +220,26 @@ func (ic *LogParsingPipelineController) ValidatePipelines(ctx context.Context,
return err
}
func (ic *LogParsingPipelineController) getNormalizePipeline() pipelinetypes.GettablePipeline {
// withNormalizePipeline places normalize where the read path dictates. Ahead of user pipelines
// when queries run on body_v2 (use_json_body), so operators see the body the explorer shows.
// After them when dual ingestion alone writes body_v2, so operators keep seeing the raw body
// users still query. Absent when neither flag is on.
func (ic *LogParsingPipelineController) withNormalizePipeline(ctx context.Context, orgID valuer.UUID, pipelines []pipelinetypes.GettablePipeline) []pipelinetypes.GettablePipeline {
evalCtx := featuretypes.NewFlaggerEvaluationContext(orgID)
dualIngestion := ic.fl.BooleanOrEmpty(ctx, flagger.FeatureJSONBodyDualIngestion, evalCtx)
switch {
case ic.fl.BooleanOrEmpty(ctx, flagger.FeatureUseJSONBody, evalCtx):
return append([]pipelinetypes.GettablePipeline{getNormalizePipeline(dualIngestion)}, pipelines...)
case dualIngestion:
return append(slices.Clone(pipelines), getNormalizePipeline(true))
default:
return pipelines
}
}
// stashOriginalBody makes normalize carry the pre-normalization body in an internal attribute
// for the exporter to restore into the legacy body column.
func getNormalizePipeline(stashOriginalBody bool) pipelinetypes.GettablePipeline {
return pipelinetypes.GettablePipeline{
StoreablePipeline: pipelinetypes.StoreablePipeline{
Name: "Default Pipeline - PreProcessing Body",
@@ -239,10 +258,11 @@ func (ic *LogParsingPipelineController) getNormalizePipeline() pipelinetypes.Get
},
Config: []pipelinetypes.PipelineOperator{
{
ID: uuid.NewString(),
Type: "normalize",
Enabled: true,
If: "body != nil",
ID: uuid.NewString(),
Type: "normalize",
Enabled: true,
If: "body != nil",
JSONBodyDualIngestion: stashOriginalBody,
},
},
}
@@ -351,8 +371,11 @@ func (ic *LogParsingPipelineController) PreviewLogsPipelines(
}
// The collector gets the same pipeline prepended over opamp; see RecommendAgentConfig.
// Under dual ingestion alone it runs after user operators and only feeds body_v2, which
// the explorer does not show yet, so the preview leaves it out. The original-body stash
// is left off: the preview has no exporter to restore and strip it.
if ic.fl.BooleanOrEmpty(ctx, flagger.FeatureUseJSONBody, featuretypes.NewFlaggerEvaluationContext(orgID)) {
pipelines = append([]pipelinetypes.GettablePipeline{ic.getNormalizePipeline()}, pipelines...)
pipelines = append([]pipelinetypes.GettablePipeline{getNormalizePipeline(false)}, pipelines...)
}
result, collectorLogs, err := SimulatePipelinesProcessing(ctx, pipelines, request.Logs)
@@ -373,16 +396,16 @@ func (pc *LogParsingPipelineController) AgentFeatureType() agentConf.AgentFeatur
// Implements agentConf.AgentFeature interface.
// RecommendAgentConfig generates the collector config to be sent to agents.
// The normalize pipeline (when use_json_body feature flag is on) is injected here, after
// rawPipelineData is serialized. So it is only present in the config sent to
// The normalize pipeline (when use_json_body or json_body_dual_ingestion is on) is placed
// here, after rawPipelineData is serialized. So it is only present in the config sent to
// the collector and never persisted to the database as part of the user's pipeline list.
//
// NOTE: The configId sent to agents is derived from the pipeline version number
// (e.g. "LogPipelines:5"), not the YAML content. If server-side logic changes
// the generated YAML without bumping the version (e.g. toggling the use_json_body
// flag or updating operator IfExpressions), agents that already applied that version will
// not re-apply the new config. In such cases, users must save a new pipeline version
// via the API to force agents to pick up the change.
// the generated YAML without bumping the version (e.g. toggling the use_json_body or
// json_body_dual_ingestion flags or updating operator IfExpressions), agents that already
// applied that version will not re-apply the new config. In such cases, users must save a
// new pipeline version via the API to force agents to pick up the change.
func (pc *LogParsingPipelineController) RecommendAgentConfig(
orgId valuer.UUID,
currentConfYaml []byte,
@@ -408,10 +431,8 @@ func (pc *LogParsingPipelineController) RecommendAgentConfig(
return nil, "", err
}
if pc.fl.BooleanOrEmpty(ctx, flagger.FeatureUseJSONBody, featuretypes.NewFlaggerEvaluationContext(orgId)) {
// add default normalize pipeline at the beginning, only for sending to collector
enrichedPipelines = append([]pipelinetypes.GettablePipeline{pc.getNormalizePipeline()}, enrichedPipelines...)
}
// normalize is only for sending to the collector, never persisted
enrichedPipelines = pc.withNormalizePipeline(ctx, orgId, enrichedPipelines)
updatedConf, err := GenerateCollectorConfigWithPipelines(currentConfYaml, enrichedPipelines)
if err != nil {

View File

@@ -1,14 +0,0 @@
package metrics
var MetricsUnderTransition = map[string]string{
"k8s.pod.cpu.utilization": "k8s.pod.cpu.usage",
"k8s.node.cpu.utilization": "k8s.node.cpu.usage",
"container.cpu.utilization": "container.cpu.usage",
}
func GetTransitionedMetric(metric string) string {
if transitionedMetric, ok := MetricsUnderTransition[metric]; ok {
return transitionedMetric
}
return metric
}

View File

@@ -10,8 +10,9 @@ import (
"log/slog"
"github.com/SigNoz/signoz/pkg/query-service/constants"
"github.com/SigNoz/signoz/pkg/query-service/metrics"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// ValidateAndCastValue validates and casts the value of a key to the corresponding data type of the key
@@ -234,12 +235,12 @@ func ClickHouseFormattedValue(v interface{}) string {
func ClickHouseFormattedMetricNames(v interface{}) string {
if name, ok := v.(string); ok {
transitionedMetrics := metrics.GetTransitionedMetric(name)
if transitionedMetrics != name {
return ClickHouseFormattedValue([]interface{}{transitionedMetrics})
} else {
return ClickHouseFormattedValue([]interface{}{name})
}
current := semconv.Current(semconv.KindMetric, telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
})
return ClickHouseFormattedValue([]interface{}{current})
}
return ClickHouseFormattedValue(v)

View File

@@ -1,6 +1,7 @@
package utils
import (
"github.com/stretchr/testify/assert"
"reflect"
"testing"
@@ -483,3 +484,21 @@ func TestGetEpochNanoSecs(t *testing.T) {
})
}
}
// The legacy readers redirect an old metric name to its current name.
func TestClickHouseFormattedMetricNames(t *testing.T) {
testCases := []struct {
name string
metric string
expected string
}{
{name: "OldName_RedirectsToCurrent", metric: "k8s.pod.cpu.utilization", expected: "['k8s.pod.cpu.usage']"},
{name: "CurrentName_Unchanged", metric: "k8s.pod.cpu.usage", expected: "['k8s.pod.cpu.usage']"},
{name: "OutsideFamily_Unchanged", metric: "http.server.duration", expected: "['http.server.duration']"},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.expected, ClickHouseFormattedMetricNames(testCase.metric))
})
}
}

View File

@@ -4,57 +4,61 @@ import (
"context"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
// semconvFamiliesEnabled evaluates the resolve_semconv_families flag for the
// SemconvFamiliesEnabled evaluates the resolve_semconv_families flag for the
// org. A nil flagger means off, so a caller without family support stays
// literal by default.
func semconvFamiliesEnabled(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger) bool {
func SemconvFamiliesEnabled(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger) bool {
if fl == nil {
return false
}
return fl.BooleanOrEmpty(ctx, flagger.FeatureResolveSemconvFamilies, featuretypes.NewFlaggerEvaluationContext(orgID))
}
// ExpandKeySelectorsForFamilies adds selectors for the other members of each
// semantic-convention family that a selector names. The metadata fetched for
// a query then contains each spelling that MatchingLogicalFields can group.
// This function is the prefetch of the resolution layer: statement builders
// call it after they derive the selectors, and the metadata store stays
// family-blind (autocomplete responses keep the literal spelling that the
// user typed). It does nothing when the resolve_semconv_families flag is off
// for the org. Only trace selectors expand today, because that matches the
// family support. Fuzzy (search-style) selectors never expand.
// ExpandKeySelectorsForFamilies adds a selector for each other spelling of
// the family a selector names, so the fetched metadata holds every member.
// Off, or for a fuzzy selector, it returns the selectors as they are. A call
// site without the prefetch stays literal and never merges wrong.
func ExpandKeySelectorsForFamilies(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, selectors []*telemetrytypes.FieldKeySelector) []*telemetrytypes.FieldKeySelector {
if !semconvFamiliesEnabled(ctx, orgID, fl) {
if !SemconvFamiliesEnabled(ctx, orgID, fl) {
return selectors
}
// The same name under another context, signal, or data type needs its
// own siblings. The metric context is not part of the key: metric callers
// duplicate the selectors per metric name after this expansion.
type identity struct {
signal telemetrytypes.Signal
fieldContext telemetrytypes.FieldContext
fieldDataType telemetrytypes.FieldDataType
name string
}
out := selectors
seen := make(map[string]bool, len(selectors))
seen := make(map[identity]bool, len(selectors))
for _, selector := range selectors {
seen[selector.Name] = true
seen[identity{selector.Signal, selector.FieldContext, selector.FieldDataType, selector.Name}] = true
}
for _, selector := range selectors {
if selector.Signal != telemetrytypes.SignalTraces ||
selector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeFuzzy {
if selector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeFuzzy {
continue
}
members := semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: selector.Name,
Signal: selector.Signal,
FieldContext: selector.FieldContext,
members := familySpellings(telemetrytypes.FieldKeySelector{
Name: selector.Name,
Signal: selector.Signal,
FieldContext: selector.FieldContext,
MetricContext: selector.MetricContext,
})
for _, member := range members {
if seen[member] {
id := identity{selector.Signal, selector.FieldContext, selector.FieldDataType, member}
if seen[id] {
continue
}
seen[member] = true
seen[id] = true
expanded := *selector
expanded.Name = member
out = append(out, &expanded)

View File

@@ -20,8 +20,9 @@ import (
// member reads, current member first. It is present when any member is
// present, and absent when no member is present. A row without any member
// reads what the tail of the merge reads: the sentinel for a string family,
// NULL for the others. A member with a value map reads in the current
// vocabulary.
// NULL for the others. When every member reads its sentinel as a value, so
// does the family, and the keyless contract of the signal survives the
// merge. A member with a value map reads in the current vocabulary.
func LogicalRead(ctx context.Context, q qbtypes.QueryInfo, storage qbtypes.Storage, logical *telemetrytypes.LogicalField) (qbtypes.Read, error) {
if !logical.IsFamily() {
return memberRead(ctx, q, storage, logical, 0)
@@ -35,7 +36,7 @@ func LogicalRead(ctx context.Context, q qbtypes.QueryInfo, storage qbtypes.Stora
reads = append(reads, read)
}
merged := qbtypes.Read{WhenAbsent: familyAbsence(logical)}
merged := qbtypes.Read{WhenAbsent: familyAbsence(logical, reads)}
guards := make([]string, 0, len(reads))
for _, read := range reads {
guards = append(guards, read.Presence)
@@ -97,10 +98,16 @@ func clickHouseStringArray(values []string) string {
}
// familyAbsence is what the merged read yields for a row without any
// member: the sentinel tail of a string family, NULL for the others.
func familyAbsence(logical *telemetrytypes.LogicalField) qbtypes.Absent {
if logical.FieldDataType == telemetrytypes.FieldDataTypeString {
return qbtypes.AbsentIsSentinel
// member: the sentinel tail of a string family, NULL for the others. When
// every member's sentinel is a value, the tail is one too.
func familyAbsence(logical *telemetrytypes.LogicalField, reads []qbtypes.Read) qbtypes.Absent {
if logical.FieldDataType != telemetrytypes.FieldDataTypeString {
return qbtypes.AbsentIsNull
}
return qbtypes.AbsentIsNull
for _, read := range reads {
if read.WhenAbsent != qbtypes.AbsentIsValue {
return qbtypes.AbsentIsSentinel
}
}
return qbtypes.AbsentIsValue
}

View File

@@ -48,7 +48,7 @@ func TestFamiliesOffByDefault(t *testing.T) {
}},
}
fields := matchingLogicalFields(false, telemetrytypes.SignalUnspecified, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
fields := matchingLogicalFields(false, telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 1)
assert.False(t, fields[0].IsFamily())
assert.Equal(t, []string{"deployment.environment.name"}, memberNames(fields[0]))
@@ -76,7 +76,7 @@ func TestMatchingLogicalFieldsGroupsFamilyMembers(t *testing.T) {
}
for _, requested := range []string{"deployment.environment.name", "deployment.environment"} {
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, &telemetrytypes.TelemetryFieldKey{Name: requested}, fieldKeys)
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{Name: requested}, fieldKeys)
require.Len(t, fields, 1, "a family is one logical field, requested via %s", requested)
logical := fields[0]
assert.Equal(t, requested, logical.Name, "response identity is the requested spelling")
@@ -106,7 +106,7 @@ func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
}},
}
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, &telemetrytypes.TelemetryFieldKey{
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
}, fieldKeys)
@@ -115,26 +115,92 @@ func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
assert.Equal(t, []string{"resource.deployment.environment.name", "deployment.environment"}, memberNames(fields[0]))
}
// Non-trace signals have no family support: the requested spelling stays
// literal, and a family member name never pulls in its siblings.
func TestMatchingLogicalFieldsKeepsLogsLiteral(t *testing.T) {
logsKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
// Log entries group into families exactly like trace entries.
func TestMatchingLogicalFieldsGroupsLogEntries(t *testing.T) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {logsKey("deployment.environment.name")},
"deployment.environment": {logsKey("deployment.environment")},
}},
"deployment.environment": {{
Name: "deployment.environment",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
}
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
fields := matchingLogicalFields(true, telemetrytypes.SignalLogs, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 1)
assert.False(t, fields[0].IsFamily())
assert.Equal(t, []string{"deployment.environment.name"}, memberNames(fields[0]))
assert.True(t, fields[0].IsFamily())
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, memberNames(fields[0]))
}
// Metric entries of a span-metrics metric group across the plain and the
// resource_ spellings of the family, in member-major order: every spelling
// of the current name precedes the first spelling of the old one.
func TestMatchingLogicalFieldsGroupsMetricSpellings(t *testing.T) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
"resource_deployment.environment.name": {{
Name: "resource_deployment.environment.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
"deployment.environment": {{
Name: "deployment.environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
"resource_deployment.environment": {{
Name: "resource_deployment.environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
}
fields := matchingLogicalFields(true, telemetrytypes.SignalMetrics, &telemetrytypes.MetricContext{MetricName: "signoz_calls_total"}, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment"}, fieldKeys)
require.Len(t, fields, 1)
assert.True(t, fields[0].IsFamily())
assert.Equal(t, []string{
"deployment.environment.name", "resource_deployment.environment.name",
"deployment.environment", "resource_deployment.environment",
}, memberNames(fields[0]))
}
// A non-string entry never joins a family: the merged read has no common
// ClickHouse type across the storages.
func TestMatchingLogicalFieldsKeepsNumberEntriesSingle(t *testing.T) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeNumber,
}},
"deployment.environment": {{
Name: "deployment.environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeNumber,
}},
}
fields := matchingLogicalFields(true, telemetrytypes.SignalMetrics, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment"}, fieldKeys)
require.Len(t, fields, 2)
for _, logical := range fields {
assert.False(t, logical.IsFamily())
}
}
// A family and a genuine same-name collision stack cleanly: the family stays
@@ -165,7 +231,7 @@ func TestResolveLogicalFieldsKeepsFamilyThroughAmbiguity(t *testing.T) {
}
requested := &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, requested, fieldKeys)
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, nil, requested, fieldKeys)
require.Len(t, fields, 2, "resource family + attribute collision")
resolved, warning := ResolveLogicalFields(requested, fields)
@@ -229,7 +295,7 @@ func TestMatchingLogicalFieldsNeverMergesAcrossDataTypes(t *testing.T) {
}},
}
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
fields := matchingLogicalFields(true, telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 2)
for _, logical := range fields {
assert.False(t, logical.IsFamily())
@@ -254,11 +320,15 @@ func TestExpandKeySelectorsForFamilies(t *testing.T) {
"service.name",
"deployment.environment.name",
"deployment.environment",
}, names, "one sibling selector for the trace family member; logs and non-family names untouched")
"deployment.environment",
}, names, "each selector identity gets its own sibling, and a non-family name stays untouched")
sibling := expanded[len(expanded)-1]
assert.Equal(t, telemetrytypes.SignalTraces, sibling.Signal)
assert.Equal(t, telemetrytypes.FieldSelectorMatchTypeExact, sibling.SelectorMatchType)
tracesSibling := expanded[len(expanded)-2]
assert.Equal(t, telemetrytypes.SignalTraces, tracesSibling.Signal)
assert.Equal(t, telemetrytypes.FieldSelectorMatchTypeExact, tracesSibling.SelectorMatchType)
logsSibling := expanded[len(expanded)-1]
assert.Equal(t, telemetrytypes.SignalLogs, logsSibling.Signal,
"a same-named selector under another signal must not take the sibling")
}
func TestExpandKeySelectorsForFamiliesDeduplicatesAndSkipsFuzzy(t *testing.T) {

View File

@@ -0,0 +1,64 @@
package querybuilder
import (
"context"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
// The span-metrics processor in signoz-otel-collector writes resource
// attributes as labels with a resource_ prefix on these metrics only. The
// signoz_latency histogram is stored as its dotted sub-metrics.
const spanMetricsResourcePrefix = "resource_"
var spanMetrics = map[string]struct{}{
"signoz_calls_total": {},
"signoz_latency": {},
"signoz_latency.bucket": {},
"signoz_latency.sum": {},
"signoz_latency.count": {},
"signoz_latency.min": {},
"signoz_latency.max": {},
"signoz_db_latency_sum": {},
"signoz_db_latency_count": {},
"signoz_external_call_latency_sum": {},
"signoz_external_call_latency_count": {},
}
// MetricLabelSpellings returns the family members of selector.Name, current
// first, and on a span-metrics metric each member with the resource_ prefix
// too. A name outside a family, or one the selector leaves ambiguous, is
// returned as it is.
func MetricLabelSpellings(selector telemetrytypes.FieldKeySelector) []string {
members := semconv.Members(semconv.KindAttribute, selector)
if len(members) <= 1 {
return []string{selector.Name}
}
if selector.MetricContext == nil {
return members
}
if _, ok := spanMetrics[selector.MetricContext.MetricName]; !ok {
return members
}
spellings := make([]string, 0, len(members)*2)
for _, member := range members {
spellings = append(spellings, member, spanMetricsResourcePrefix+member)
}
return spellings
}
// FamilyMetricNames returns the metric-name family of metricName when the
// flag is on for the org, else the name alone.
func FamilyMetricNames(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, metricName string) []string {
if !SemconvFamiliesEnabled(ctx, orgID, fl) {
return []string{metricName}
}
return semconv.Members(semconv.KindMetric, telemetrytypes.FieldKeySelector{
Name: metricName,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
})
}

View File

@@ -0,0 +1,100 @@
package querybuilder
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
)
func TestMetricLabelSpellingsReturnsTheFamilyMembers(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "deployment.environment",
Signal: telemetrytypes.SignalMetrics,
MetricContext: &telemetrytypes.MetricContext{MetricName: "k8s.pod.cpu.usage"},
}
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, MetricLabelSpellings(selector))
}
func TestMetricLabelSpellingsAddsTheResourcePrefixForSpanMetrics(t *testing.T) {
testCases := []struct {
name string
metric string
expected []string
}{
{
name: "SpanMetric_ReadsPlainAndResourceSpellings",
metric: "signoz_calls_total",
expected: []string{
"deployment.environment.name", "resource_deployment.environment.name",
"deployment.environment", "resource_deployment.environment",
},
},
{
name: "LatencyHistogramSubMetric_ReadsPlainAndResourceSpellings",
metric: "signoz_latency.bucket",
expected: []string{
"deployment.environment.name", "resource_deployment.environment.name",
"deployment.environment", "resource_deployment.environment",
},
},
{
name: "OtherSignozMetric_ReadsPlainSpellings",
metric: "signoz_other_metric",
expected: []string{"deployment.environment.name", "deployment.environment"},
},
{
name: "UnderscoreLatencySubMetric_ReadsPlainSpellings",
metric: "signoz_latency_bucket",
expected: []string{"deployment.environment.name", "deployment.environment"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "deployment.environment",
Signal: telemetrytypes.SignalMetrics,
MetricContext: &telemetrytypes.MetricContext{MetricName: testCase.metric},
}
assert.Equal(t, testCase.expected, MetricLabelSpellings(selector))
})
}
}
// A requested name is never rewritten: a resource_ spelling that is not a
// family member stays literal, on a span metric too.
func TestMetricLabelSpellingsKeepsThePrefixedRequestLiteral(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "resource_deployment.environment",
Signal: telemetrytypes.SignalMetrics,
MetricContext: &telemetrytypes.MetricContext{MetricName: "signoz_calls_total"},
}
assert.Equal(t, []string{"resource_deployment.environment"}, MetricLabelSpellings(selector))
}
func TestMetricLabelSpellingsStaysLiteralOutsideTheVocabulary(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "http.route",
Signal: telemetrytypes.SignalMetrics,
}
assert.Equal(t, []string{"http.route"}, MetricLabelSpellings(selector))
}
func TestFamilyMetricNames(t *testing.T) {
on := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
assert.Equal(t, []string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"}, FamilyMetricNames(context.Background(), valuer.UUID{}, on, "k8s.pod.cpu.utilization"))
assert.Equal(t, []string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"}, FamilyMetricNames(context.Background(), valuer.UUID{}, on, "k8s.pod.cpu.usage"))
assert.Equal(t, []string{"http.server.duration"}, FamilyMetricNames(context.Background(), valuer.UUID{}, on, "http.server.duration"))
off := flaggertest.WithBooleanFlags(t, map[string]bool{})
assert.Equal(t, []string{"k8s.pod.cpu.utilization"}, FamilyMetricNames(context.Background(), valuer.UUID{}, off, "k8s.pod.cpu.utilization"))
}

View File

@@ -21,7 +21,7 @@ func NewQueryInfo(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, si
EndNs: endNs,
Signal: signal,
Metric: metric,
FamiliesOn: semconvFamiliesEnabled(ctx, orgID, fl),
FamiliesOn: SemconvFamiliesEnabled(ctx, orgID, fl),
}
if fl != nil {
evalCtx := featuretypes.NewFlaggerEvaluationContext(orgID)
@@ -61,7 +61,7 @@ func Resolve(
traits := storage.Traits()
lookup := key
matches := matchingLogicalFields(q.FamiliesOn, q.Signal, key, fieldKeys)
matches := matchingLogicalFields(q.FamiliesOn, q.Signal, q.Metric, key, fieldKeys)
if len(matches) == 0 && slices.Contains(traits.OwnContexts, key.FieldContext) {
// a column the storage knows under the key's own context is the key
// as written, and only a miss corrects to the bare spelling
@@ -71,7 +71,7 @@ func Resolve(
}
}
lookup = telemetrytypes.NewTelemetryFieldKey(key.Name, telemetrytypes.FieldContextUnspecified, key.FieldDataType)
matches = matchingLogicalFields(q.FamiliesOn, q.Signal, lookup, fieldKeys)
matches = matchingLogicalFields(q.FamiliesOn, q.Signal, q.Metric, lookup, fieldKeys)
}
resolved := qbtypes.Resolved{Key: key, Ambiguous: len(matches) > 1}

View File

@@ -1012,43 +1012,36 @@ func assignIfEmpty(s *string, value string) {
}
}
// familyMemberNames returns the physical spellings to look up for the
// referenced key: the semantic-convention family members (current-first) when
// families are on and the query can resolve to traces, else just the requested
// name. Only the traces storage understands families today. Logs and
// metrics keep the requested spelling until theirs land.
func familyMemberNames(familiesOn bool, signal telemetrytypes.Signal, field *telemetrytypes.TelemetryFieldKey) []string {
if !familiesOn {
return []string{field.Name}
// familySpellings returns the storage spellings for the selector. Metrics
// add the span-metrics label layout.
func familySpellings(selector telemetrytypes.FieldKeySelector) []string {
if selector.Signal == telemetrytypes.SignalMetrics {
return MetricLabelSpellings(selector)
}
if signal != telemetrytypes.SignalUnspecified && signal != telemetrytypes.SignalTraces {
return []string{field.Name}
}
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: field.FieldContext,
})
return semconv.Members(semconv.KindAttribute, selector)
}
// matchingLogicalFields resolves the referenced key against the metadata map
// into logical fields, honoring any context/data type the user specified.
//
// Physical keys that are members of one semantic-convention family (traces
// only today) group into one logical field per (signal, context, data type)
// identity, members ordered current-first. Every other matching key becomes
// its own single-member logical field. Ambiguity is the length of the
// returned slice: one family is one element and is never ambiguous with
// itself, but the slice can hold several logical fields, including several
// family fields, one per identity, when the family exists under more than
// one context or data type. Members alias the metadata map entries; nothing
// is copied or mutated.
//
// Family grouping only happens when families are on for the query. Off,
// every match stays a single-member logical field.
func matchingLogicalFields(familiesOn bool, signal telemetrytypes.Signal, field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
members := familyMemberNames(familiesOn, signal, field)
matches := collectMemberMatches(field, members, fieldKeys)
// matchingLogicalFields resolves the key against the metadata map. Members
// of one family group into one logical field per (signal, context, data
// type) identity, current first. Every other match is its own single-member
// field. The length of the result is the ambiguity: a family is never
// ambiguous with itself. Members alias the map entries. With families off,
// only the requested name is looked up. The key's own signal wins over the
// query's signal.
func matchingLogicalFields(familiesOn bool, signal telemetrytypes.Signal, metric *telemetrytypes.MetricContext, field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
members := []string{field.Name}
if familiesOn {
if field.Signal != telemetrytypes.SignalUnspecified {
signal = field.Signal
}
members = familySpellings(telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: signal,
FieldContext: field.FieldContext,
MetricContext: metric,
})
}
matches := collectMemberMatches(field, members, metric, fieldKeys)
return groupIntoLogicalFields(field.Name, len(members) > 1, matches)
}
@@ -1060,48 +1053,35 @@ type memberMatch struct {
rank int
}
// matchesRequestedIdentity reports whether the entry fits the context and data
// type that the request specified; unspecified matches any. A context-prefixed
// lookup already matched the context through the lookup key itself.
func matchesRequestedIdentity(field, item *telemetrytypes.TelemetryFieldKey, contextMatched bool) bool {
if !contextMatched && field.FieldContext != telemetrytypes.FieldContextUnspecified && field.FieldContext != item.FieldContext {
return false
}
if field.FieldDataType != telemetrytypes.FieldDataTypeUnspecified && field.FieldDataType != item.FieldDataType {
return false
}
return true
}
// inFamilyScope reports whether a match found under a sibling member name is
// legitimate: the entry must be trace metadata, and the member must be in the
// family of the requested name for the entry's context. A member lookup can
// otherwise find a same-named field in a scope where the family does not
// apply.
func inFamilyScope(field, item *telemetrytypes.TelemetryFieldKey, memberName string) bool {
if item.Signal != telemetrytypes.SignalTraces {
return false
}
return slices.Contains(semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: item.FieldContext,
}), memberName)
}
// collectMemberMatches finds the metadata entries for every member spelling:
// first under the member names, then under their context-prefixed spellings
// (a context can be a legitimate part of a stored name, e.g. `attribute.key`).
func collectMemberMatches(field *telemetrytypes.TelemetryFieldKey, members []string, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []memberMatch {
// collectMemberMatches finds the metadata entries for every member spelling,
// under the member name and under its context-prefixed spelling, because a
// context can be part of a stored name. An unspecified context or data type
// matches any.
func collectMemberMatches(field *telemetrytypes.TelemetryFieldKey, members []string, metric *telemetrytypes.MetricContext, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []memberMatch {
matches := make([]memberMatch, 0)
collect := func(lookupName string, rank int, memberName string, contextMatched bool) {
for _, item := range fieldKeys[lookupName] {
if !matchesRequestedIdentity(field, item, contextMatched) {
// A context-prefixed lookup matched the context through the key.
if !contextMatched && field.FieldContext != telemetrytypes.FieldContextUnspecified && field.FieldContext != item.FieldContext {
continue
}
if memberName != field.Name && !inFamilyScope(field, item, memberName) {
if field.FieldDataType != telemetrytypes.FieldDataTypeUnspecified && field.FieldDataType != item.FieldDataType {
continue
}
if memberName != field.Name {
// A sibling can match a same-named field where the family does
// not apply, so the member must be a spelling of the requested
// name for the entry's own signal and context.
spellings := familySpellings(telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
MetricContext: metric,
})
if !slices.Contains(spellings, memberName) {
continue
}
}
matches = append(matches, memberMatch{key: item, rank: rank})
}
}
@@ -1117,18 +1097,22 @@ func collectMemberMatches(field *telemetrytypes.TelemetryFieldKey, members []str
return matches
}
// groupIntoLogicalFields turns matches into logical fields. Trace entries in
// family mode group by their (signal, context, data type) identity; every
// other entry becomes its own single-member field. Members sort by family
// rank at the end: precedence is a property of the family, not of the order
// in which the lookups found the members.
// groupIntoLogicalFields groups a string entry of a family signal under the
// resource or attribute context by its (signal, context, data type)
// identity. Every other entry is its own single-member field. Members sort
// by family rank, not by lookup order.
func groupIntoLogicalFields(requestedName string, familyMode bool, matches []memberMatch) []*telemetrytypes.LogicalField {
fields := make([]*telemetrytypes.LogicalField, 0, len(matches))
groups := make(map[string]*telemetrytypes.LogicalField)
ranks := make(map[*telemetrytypes.TelemetryFieldKey]int)
for _, match := range matches {
if !familyMode || match.key.Signal != telemetrytypes.SignalTraces {
familySignal := match.key.Signal == telemetrytypes.SignalTraces ||
match.key.Signal == telemetrytypes.SignalLogs ||
match.key.Signal == telemetrytypes.SignalMetrics
familyContext := match.key.FieldContext == telemetrytypes.FieldContextResource ||
match.key.FieldContext == telemetrytypes.FieldContextAttribute
if !familyMode || !familySignal || !familyContext || match.key.FieldDataType != telemetrytypes.FieldDataTypeString {
fields = append(fields, telemetrytypes.SingleLogicalField(requestedName, match.key))
continue
}
@@ -1145,7 +1129,10 @@ func groupIntoLogicalFields(requestedName string, familyMode bool, matches []mem
groups[identity] = group
fields = append(fields, group)
}
if groupHasMemberNamed(group, match.key.Name) {
alreadyMember := slices.ContainsFunc(group.Members, func(member *telemetrytypes.TelemetryFieldKey) bool {
return member.Name == match.key.Name
})
if alreadyMember {
continue
}
ranks[match.key] = match.rank
@@ -1159,12 +1146,3 @@ func groupIntoLogicalFields(requestedName string, familyMode bool, matches []mem
}
return fields
}
func groupHasMemberNamed(group *telemetrytypes.LogicalField, name string) bool {
for _, member := range group.Members {
if member.Name == name {
return true
}
}
return false
}

View File

@@ -589,7 +589,7 @@ func TestVisitKey(t *testing.T) {
// and decides not-found handling. Replay that here against the generic
// builder behavior (error unless the key is ignored). The test maps carry
// no signal, so every logical field is single-member and flattens losslessly.
matching := matchingLogicalFields(false, telemetrytypes.SignalUnspecified, key, tt.fieldKeys)
matching := matchingLogicalFields(false, telemetrytypes.SignalUnspecified, nil, key, tt.fieldKeys)
resolved, warning := ResolveLogicalFields(key, matching)
keys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(resolved))
for _, logical := range resolved {

View File

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

View File

@@ -1,6 +1,7 @@
package semconv
import (
"iter"
"slices"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -14,16 +15,47 @@ type Kind struct {
valuer.String
}
// Family is one logical telemetry field. Old is ordered from the most recent
// predecessor to the oldest one and therefore also defines fallback order.
// Member is one historical spelling with the scope of its rename edges. A
// nil axis is unconstrained.
type Member struct {
name string
contexts []telemetrytypes.FieldContext
signals []telemetrytypes.Signal
applyToMetrics []string
}
func (m Member) Name() string {
return m.name
}
// Family is one logical telemetry field. Members run from the most recent
// predecessor to the oldest, which is the fallback order. The family-level
// contexts and signals come from the overlay and gate where the family
// resolves. The member scopes come from the schema edges and gate which
// members apply.
type Family struct {
Current string
Old []string
Kind Kind
Contexts []telemetrytypes.FieldContext
Signals []telemetrytypes.Signal
ApplyToMetrics []string
ValueMap map[string]string
current string
kind Kind
members []Member
contexts []telemetrytypes.FieldContext
signals []telemetrytypes.Signal
}
func (f Family) Current() string {
return f.current
}
func (f Family) Kind() Kind {
return f.kind
}
// Old returns the historical spellings in fallback order.
func (f Family) Old() []string {
names := make([]string, len(f.members))
for i, member := range f.members {
names[i] = member.name
}
return names
}
var (
@@ -38,90 +70,143 @@ func (Kind) Enum() []any {
return []any{KindAttribute, KindMetric}
}
// Lookup returns the enabled family containing selector.Name for kind. The
// returned family must not be modified.
func Lookup(kind Kind, selector telemetrytypes.FieldKeySelector) (Family, bool) {
idx, ok := lookupIndex(kind, selector)
if !ok {
return Family{}, false
}
return families[idx], true
}
// Members returns the current name first, followed by historical names in
// fallback order. A name outside an enabled family is returned unchanged. The
// returned slice must not be modified.
// Members returns the current name first, followed by the historical spellings
// admitted for the selector, in fallback order. A name outside an enabled
// family is returned unchanged. So is a name the selector leaves ambiguous.
// The returned slice must not be modified.
func Members(kind Kind, selector telemetrytypes.FieldKeySelector) []string {
idx, ok := lookupIndex(kind, selector)
if !ok {
return []string{selector.Name}
}
return familyMembers[idx]
return admittedMembers(idx, selector)
}
func admittedMembers(idx int, selector telemetrytypes.FieldKeySelector) []string {
admitted := 0
for _, member := range families[idx].members {
if memberAdmits(member, selector) {
admitted++
}
}
if admitted == len(families[idx].members) {
return familyMembers[idx]
}
names := make([]string, 0, admitted+1)
names = append(names, families[idx].current)
for _, member := range families[idx].members {
if memberAdmits(member, selector) {
names = append(names, member.name)
}
}
return names
}
// Current returns the current name for selector.Name, or the input name when
// it does not belong to an enabled family.
// it does not resolve to a family.
func Current(kind Kind, selector telemetrytypes.FieldKeySelector) string {
idx, ok := lookupIndex(kind, selector)
if !ok {
return selector.Name
}
return families[idx].Current
return families[idx].current
}
// All returns every enabled family. The returned slice and families must not be
// modified.
func All() []Family {
return families
func All() iter.Seq[Family] {
return func(yield func(Family) bool) {
for _, family := range families {
if !yield(family) {
return
}
}
}
}
func buildIndexes() (map[string][]int, [][]string) {
index := make(map[string][]int)
members := make([][]string, len(families))
add := func(name string, i int) {
if !slices.Contains(index[name], i) {
index[name] = append(index[name], i)
}
}
for i, family := range families {
members[i] = make([]string, 0, len(family.Old)+1)
members[i] = append(members[i], family.Current)
members[i] = append(members[i], family.Old...)
index[family.Current] = append(index[family.Current], i)
for _, old := range family.Old {
index[old] = append(index[old], i)
members[i] = make([]string, 0, len(family.members)+1)
members[i] = append(members[i], family.current)
add(family.current, i)
for _, member := range family.members {
members[i] = append(members[i], member.name)
add(member.name, i)
}
}
return index, members
}
// lookupIndex returns the one family that admits selector.Name for kind.
// An axis the selector leaves empty constrains nothing. When more than one
// family is admitted, the name stays literal.
func lookupIndex(kind Kind, selector telemetrytypes.FieldKeySelector) (int, bool) {
found, foundIdx := 0, 0
for _, idx := range memberToFamilies[selector.Name] {
if matchesSelector(families[idx], kind, selector) {
return idx, true
if familyAdmits(families[idx], kind, selector) {
found++
foundIdx = idx
}
}
return 0, false
if found != 1 {
return 0, false
}
return foundIdx, true
}
func matchesSelector(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if family.Kind != kind {
// familyAdmits reports whether the family gate admits the selector and the
// name is the current name or an admitted member.
func familyAdmits(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if family.kind != kind {
return false
}
if selector.Signal != telemetrytypes.SignalUnspecified && len(family.Signals) > 0 {
if !slices.Contains(family.Signals, selector.Signal) {
return false
if !axisAdmits(family.signals, selector.Signal, telemetrytypes.SignalUnspecified) {
return false
}
if !axisAdmits(family.contexts, selector.FieldContext, telemetrytypes.FieldContextUnspecified) {
return false
}
if selector.Name == family.current {
for _, member := range family.members {
if memberAdmits(member, selector) {
return true
}
}
return false
}
for _, member := range family.members {
if member.name == selector.Name && memberAdmits(member, selector) {
return true
}
}
return false
}
if selector.FieldContext != telemetrytypes.FieldContextUnspecified && len(family.Contexts) > 0 {
if !slices.Contains(family.Contexts, selector.FieldContext) {
return false
}
// memberAdmits reports whether the member applies for the selector. An empty
// axis on either side admits.
func memberAdmits(member Member, selector telemetrytypes.FieldKeySelector) bool {
if !axisAdmits(member.signals, selector.Signal, telemetrytypes.SignalUnspecified) {
return false
}
if selector.Signal == telemetrytypes.SignalMetrics && len(family.ApplyToMetrics) > 0 {
if selector.MetricContext == nil {
return false
}
return slices.Contains(family.ApplyToMetrics, selector.MetricContext.MetricName)
if !axisAdmits(member.contexts, selector.FieldContext, telemetrytypes.FieldContextUnspecified) {
return false
}
if len(member.applyToMetrics) > 0 &&
selector.MetricContext != nil && selector.MetricContext.MetricName != "" &&
!slices.Contains(member.applyToMetrics, selector.MetricContext.MetricName) {
return false
}
return true
}
func axisAdmits[T comparable](scope []T, value T, unspecified T) bool {
if len(scope) == 0 || value == unspecified {
return true
}
return slices.Contains(scope, value)
}

View File

@@ -78,3 +78,116 @@ func TestMembersReturnsInputWhenKindDoesNotMatch(t *testing.T) {
"an attribute family must not match a metric-name lookup",
)
}
func TestFamilySignalsGateResolution(t *testing.T) {
swapFamilies(t, []Family{{
current: "gated.current",
kind: KindAttribute,
members: []Member{{name: "gated.old"}},
signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalTraces},
}})
metrics := telemetrytypes.FieldKeySelector{Name: "gated.old", Signal: telemetrytypes.SignalMetrics}
logs := telemetrytypes.FieldKeySelector{Name: "gated.old", Signal: telemetrytypes.SignalLogs}
assert.Equal(t, []string{"gated.old"}, Members(KindAttribute, metrics),
"a family gated to traces and logs must stay literal for metrics")
assert.Equal(t, []string{"gated.current", "gated.old"}, Members(KindAttribute, logs),
"the gate admits the signals it lists")
}
func TestMetricNameFamilyResolves(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{Name: "k8s.pod.cpu.utilization", Signal: telemetrytypes.SignalMetrics}
assert.Equal(t, []string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"}, Members(KindMetric, selector))
assert.Equal(t, "k8s.pod.cpu.usage", Current(KindMetric, selector))
assert.Equal(t, []string{"k8s.pod.cpu.utilization"}, Members(KindAttribute, selector),
"a metric-name family must not match an attribute lookup")
}
func TestMembersReturnsSharedSliceForUnscopedFamily(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{Name: "deployment.environment", Signal: telemetrytypes.SignalTraces}
first := Members(KindAttribute, selector)
second := Members(KindAttribute, selector)
assert.Equal(t, &first[0], &second[0],
"a family whose members all admit must return the precomputed slice, not a copy")
}
func TestAllIteratesEnabledFamilies(t *testing.T) {
currents := []string{}
for family := range All() {
currents = append(currents, family.Current())
}
assert.Contains(t, currents, "deployment.environment.name")
assert.Contains(t, currents, "k8s.pod.cpu.usage")
}
// swapFamilies replaces the generated table for one test so scoped-member and
// fan-out behavior can be pinned without enabling such families for real.
func swapFamilies(t *testing.T, replacement []Family) {
t.Helper()
prevFamilies, prevIndex, prevMembers := families, memberToFamilies, familyMembers
families = replacement
memberToFamilies, familyMembers = buildIndexes()
t.Cleanup(func() {
families, memberToFamilies, familyMembers = prevFamilies, prevIndex, prevMembers
})
}
func TestFanOutResolvesOnlyWithEnoughInformation(t *testing.T) {
swapFamilies(t, []Family{
{
current: "cpu.mode",
kind: KindAttribute,
members: []Member{{name: "state", applyToMetrics: []string{"system.cpu.time"}}},
},
{
current: "db.client.connection.state",
kind: KindAttribute,
members: []Member{{name: "state", applyToMetrics: []string{"db.client.connections.usage"}}},
},
})
ambiguous := telemetrytypes.FieldKeySelector{Name: "state", Signal: telemetrytypes.SignalMetrics}
assert.Equal(t, []string{"state"}, Members(KindAttribute, ambiguous),
"without a metric name, a fanned-out member admits several families and must stay literal")
pinned := ambiguous
pinned.MetricContext = &telemetrytypes.MetricContext{MetricName: "system.cpu.time"}
assert.Equal(t, []string{"cpu.mode", "state"}, Members(KindAttribute, pinned),
"the metric name disambiguates the fan-out")
outside := ambiguous
outside.MetricContext = &telemetrytypes.MetricContext{MetricName: "http.server.duration"}
assert.Equal(t, []string{"state"}, Members(KindAttribute, outside),
"a metric outside every apply_to_metrics list resolves no family")
}
func TestMemberScopesFilterMembers(t *testing.T) {
swapFamilies(t, []Family{{
current: "user_agent.original",
kind: KindAttribute,
members: []Member{
{name: "http.user_agent", contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextAttribute}, signals: []telemetrytypes.Signal{telemetrytypes.SignalTraces}},
{name: "browser.user_agent", contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextResource}},
},
}})
resource := telemetrytypes.FieldKeySelector{
Name: "user_agent.original",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, []string{"user_agent.original", "browser.user_agent"}, Members(KindAttribute, resource),
"a strict resource lookup must not include the span-only member")
attribute := resource
attribute.FieldContext = telemetrytypes.FieldContextAttribute
assert.Equal(t, []string{"user_agent.original", "http.user_agent"}, Members(KindAttribute, attribute),
"a strict attribute lookup must not include the resource-only member")
strictResourceOldSpan := resource
strictResourceOldSpan.Name = "http.user_agent"
assert.Equal(t, []string{"http.user_agent"}, Members(KindAttribute, strictResourceOldSpan),
"an old spelling outside its own scope stays literal")
}

View File

@@ -0,0 +1,256 @@
package logsstatementbuilder
import (
"context"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/require"
)
// A filter on either spelling of an enabled family compiles to one merged
// condition over the log resource maps. The flag default keeps it literal.
func TestStatementBuilderResolvesLogFamilies(t *testing.T) {
releaseTime := time.Date(2024, 1, 15, 10, 0, 0, 0, time.UTC)
releaseTimeNano := uint64(releaseTime.UnixNano())
testCases := []struct {
name string
familyOn bool
expected string
}{
{
name: "FamiliesOn",
familyOn: true,
expected: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC",
},
{
name: "FamiliesOff",
familyOn: false,
expected: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'deployment.environment') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): testCase.familyOn,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys := logstelemetryschema.BuildCompleteFieldKeyMap(releaseTime)
for _, name := range []string{"deployment.environment.name", "deployment.environment"} {
keys[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
mockMetadataStore.KeysMap = keys
storage := logstelemetryschema.NewStorage()
aggExprRewriter := querybuilder.NewAggExprRewriter(instrumentationtest.New().ToProviderSettings(), nil, storage, fl, telemetrytypes.SignalLogs)
statementBuilder := NewLogQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore, storage, aggExprRewriter,
logstelemetryschema.DefaultFullTextColumn, fl, nil,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
query := qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{
Signal: telemetrytypes.SignalLogs,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.LogAggregation{{Expression: "count()"}},
Filter: &qbtypes.Filter{Expression: "resource.deployment.environment = 'production'"},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{},
releaseTimeNano+uint64(24*time.Hour.Nanoseconds()),
releaseTimeNano+uint64(48*time.Hour.Nanoseconds()),
qbtypes.RequestTypeScalar, query, nil)
require.NoError(t, err)
require.Equal(t, testCase.expected, q.Query)
})
}
}
// The predicate of a filtered aggregation resolves the family exactly like
// the main WHERE clause.
func TestStatementBuilderResolvesLogFamilyFilteredAggregation(t *testing.T) {
releaseTime := time.Date(2024, 1, 15, 10, 0, 0, 0, time.UTC)
releaseTimeNano := uint64(releaseTime.UnixNano())
testCases := []struct {
name string
familyOn bool
expected string
}{
{name: "FamiliesOn", familyOn: true, expected: "SELECT countIf((COALESCE(NULLIF(multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL), ''), NULLIF(multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, mapContains(resources_string, 'deployment.environment'), resources_string['deployment.environment'], NULL), ''), '') = ? AND (multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL) IS NOT NULL OR multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, mapContains(resources_string, 'deployment.environment'), resources_string['deployment.environment'], NULL) IS NOT NULL))) AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC"},
{name: "FamiliesOff", familyOn: false, expected: "SELECT countIf(multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL) = ?) AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC"},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): testCase.familyOn,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys := logstelemetryschema.BuildCompleteFieldKeyMap(releaseTime)
for _, name := range []string{"deployment.environment.name", "deployment.environment"} {
keys[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
mockMetadataStore.KeysMap = keys
storage := logstelemetryschema.NewStorage()
aggExprRewriter := querybuilder.NewAggExprRewriter(instrumentationtest.New().ToProviderSettings(), nil, storage, fl, telemetrytypes.SignalLogs)
statementBuilder := NewLogQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore, storage, aggExprRewriter,
logstelemetryschema.DefaultFullTextColumn, fl, nil,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
query := qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{
Signal: telemetrytypes.SignalLogs,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.LogAggregation{{Expression: "countIf(deployment.environment.name = 'production')"}},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{},
releaseTimeNano+uint64(24*time.Hour.Nanoseconds()),
releaseTimeNano+uint64(48*time.Hour.Nanoseconds()),
qbtypes.RequestTypeScalar, query, nil)
require.NoError(t, err)
require.Equal(t, testCase.expected, q.Query)
})
}
}
// The mid-migration state: metadata holds one spelling of the family, and
// the query names the other. The filter and the group by both read the one
// stored spelling.
func TestStatementBuilderResolvesSingleSpellingAcrossNames(t *testing.T) {
releaseTime := time.Date(2024, 1, 15, 10, 0, 0, 0, time.UTC)
releaseTimeNano := uint64(releaseTime.UnixNano())
testCases := []struct {
name string
stored string
queried string
expected string
}{
{name: "OldDataQueriedByCurrentName", stored: "deployment.environment", queried: "deployment.environment.name", expected: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'deployment.environment') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT toString(multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, mapContains(resources_string, 'deployment.environment'), resources_string['deployment.environment'], NULL)) AS `__GROUP_BY_KEY_0_deployment.environment.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY __result_0 DESC"},
{name: "CurrentDataQueriedByOldName", stored: "deployment.environment.name", queried: "deployment.environment", expected: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'deployment.environment.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT toString(multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL)) AS `__GROUP_BY_KEY_0_deployment.environment`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_deployment.environment` ORDER BY __result_0 DESC"},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys := logstelemetryschema.BuildCompleteFieldKeyMap(releaseTime)
keys[testCase.stored] = []*telemetrytypes.TelemetryFieldKey{{
Name: testCase.stored,
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
mockMetadataStore.KeysMap = keys
storage := logstelemetryschema.NewStorage()
aggExprRewriter := querybuilder.NewAggExprRewriter(instrumentationtest.New().ToProviderSettings(), nil, storage, fl, telemetrytypes.SignalLogs)
statementBuilder := NewLogQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore, storage, aggExprRewriter,
logstelemetryschema.DefaultFullTextColumn, fl, nil,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
query := qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{
Signal: telemetrytypes.SignalLogs,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.LogAggregation{{Expression: "count()"}},
Filter: &qbtypes.Filter{Expression: testCase.queried + " = 'production'"},
GroupBy: []qbtypes.GroupByKey{
{TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: testCase.queried}},
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{},
releaseTimeNano+uint64(24*time.Hour.Nanoseconds()),
releaseTimeNano+uint64(48*time.Hour.Nanoseconds()),
qbtypes.RequestTypeScalar, query, nil)
require.NoError(t, err)
require.Equal(t, testCase.expected, q.Query)
})
}
}
// Group by resolves the family exactly like the filter. The merged column
// reads the spellings current-first with empty falling through, and a row
// with no member keeps the NULL group of a single key.
func TestStatementBuilderResolvesLogFamilyGroupBy(t *testing.T) {
releaseTime := time.Date(2024, 1, 15, 10, 0, 0, 0, time.UTC)
releaseTimeNano := uint64(releaseTime.UnixNano())
testCases := []struct {
name string
familyOn bool
expected string
}{
{name: "FamiliesOn", familyOn: true, expected: "SELECT toString(multiIf((multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL) IS NOT NULL OR multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, mapContains(resources_string, 'deployment.environment'), resources_string['deployment.environment'], NULL) IS NOT NULL), COALESCE(NULLIF(multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL), ''), NULLIF(multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, mapContains(resources_string, 'deployment.environment'), resources_string['deployment.environment'], NULL), ''), ''), NULL)) AS `__GROUP_BY_KEY_0_deployment.environment.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY __result_0 DESC"},
{name: "FamiliesOff", familyOn: false, expected: "SELECT toString(multiIf(resource.`deployment.environment.name` IS NOT NULL, resource.`deployment.environment.name`::String, mapContains(resources_string, 'deployment.environment.name'), resources_string['deployment.environment.name'], NULL)) AS `__GROUP_BY_KEY_0_deployment.environment.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY __result_0 DESC"},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): testCase.familyOn,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys := logstelemetryschema.BuildCompleteFieldKeyMap(releaseTime)
for _, name := range []string{"deployment.environment.name", "deployment.environment"} {
keys[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
mockMetadataStore.KeysMap = keys
storage := logstelemetryschema.NewStorage()
aggExprRewriter := querybuilder.NewAggExprRewriter(instrumentationtest.New().ToProviderSettings(), nil, storage, fl, telemetrytypes.SignalLogs)
statementBuilder := NewLogQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore, storage, aggExprRewriter,
logstelemetryschema.DefaultFullTextColumn, fl, nil,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
query := qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{
Signal: telemetrytypes.SignalLogs,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.LogAggregation{{Expression: "count()"}},
GroupBy: []qbtypes.GroupByKey{
{TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}},
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{},
releaseTimeNano+uint64(24*time.Hour.Nanoseconds()),
releaseTimeNano+uint64(48*time.Hour.Nanoseconds()),
qbtypes.RequestTypeScalar, query, nil)
require.NoError(t, err)
require.Equal(t, testCase.expected, q.Query)
})
}
}

View File

@@ -125,6 +125,7 @@ func (b *logQueryStatementBuilder) Build(
bodyJSONEnabled := b.fl.BooleanOrEmpty(ctx, flagger.FeatureUseJSONBody, featuretypes.NewFlaggerEvaluationContext(orgID))
keySelectors, warnings := getKeySelectors(query, bodyJSONEnabled)
keySelectors = querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, b.fl, keySelectors)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, err

View File

@@ -0,0 +1,175 @@
package metricsstatementbuilder
import (
"context"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
"github.com/SigNoz/signoz/pkg/types/metrictypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/require"
)
// The flag merges the label spellings of the family and unions the storage
// names of the metric-name family, in the filter, the group by column, and
// every metric_name filter.
func TestStatementBuilderResolvesFamilies(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
mockMetadataStore.KeysMap = map[string][]*telemetrytypes.TelemetryFieldKey{}
for _, name := range []string{"deployment.environment.name", "deployment.environment"} {
mockMetadataStore.KeysMap[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
statementBuilder := NewMetricQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore,
metricstelemetryschema.NewStorage(),
fl,
)
query := qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]{
Signal: telemetrytypes.SignalMetrics,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.MetricAggregation{
{
MetricName: "k8s.pod.cpu.utilization",
Type: metrictypes.GaugeType,
Temporality: metrictypes.Unspecified,
TimeAggregation: metrictypes.TimeAggregationAvg,
SpaceAggregation: metrictypes.SpaceAggregationAvg,
},
},
Filter: &qbtypes.Filter{
Expression: "deployment.environment = 'production'",
},
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
},
},
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, qbtypes.RequestTypeTimeSeries, query, nil)
require.NoError(t, err)
require.Equal(t, "WITH __temporal_aggregation_cte AS (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, COALESCE(NULLIF(JSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(JSONExtractString(labels, 'deployment.environment'), ''), '') AS `__GROUP_BY_KEY_0_deployment.environment.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?, ?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND COALESCE(NULLIF(JSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(JSONExtractString(labels, 'deployment.environment'), ''), '') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_deployment.environment.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?, ?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY fingerprint, ts), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_deployment.environment.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_deployment.environment.name`, ts", q.Query)
require.Equal(t, []any{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization", uint64(1747936800000), uint64(1747983420000), "unspecified", "production", "k8s.pod.cpu.usage", "k8s.pod.cpu.utilization", uint64(1747947390000), uint64(1747983420000), 0}, q.Args)
}
// The mid-migration state: metadata holds only the old label spelling, and
// the query names the current one. The filter and the group by both read
// the one stored label.
func TestStatementBuilderResolvesSingleSpellingAcrossNames(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
mockMetadataStore.KeysMap = map[string][]*telemetrytypes.TelemetryFieldKey{}
for _, name := range []string{"deployment.environment"} {
mockMetadataStore.KeysMap[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
statementBuilder := NewMetricQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore,
metricstelemetryschema.NewStorage(),
fl,
)
query := qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]{
Signal: telemetrytypes.SignalMetrics,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.MetricAggregation{
{
MetricName: "k8s.pod.cpu.utilization",
Type: metrictypes.GaugeType,
Temporality: metrictypes.Unspecified,
TimeAggregation: metrictypes.TimeAggregationAvg,
SpaceAggregation: metrictypes.SpaceAggregationAvg,
},
},
Filter: &qbtypes.Filter{
Expression: "deployment.environment.name = 'production'",
},
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
},
},
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, qbtypes.RequestTypeTimeSeries, query, nil)
require.NoError(t, err)
require.Equal(t, "WITH __temporal_aggregation_cte AS (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, JSONExtractString(labels, 'deployment.environment') AS `__GROUP_BY_KEY_0_deployment.environment.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?, ?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND JSONExtractString(labels, 'deployment.environment') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_deployment.environment.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?, ?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY fingerprint, ts), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_deployment.environment.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_deployment.environment.name`, ts", q.Query)
}
// With the flag at its default, both the labels and the metric name stay
// literal.
func TestStatementBuilderKeepsFamiliesLiteralByDefault(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): false,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
mockMetadataStore.KeysMap = map[string][]*telemetrytypes.TelemetryFieldKey{}
for _, name := range []string{"deployment.environment.name", "deployment.environment"} {
mockMetadataStore.KeysMap[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
statementBuilder := NewMetricQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore,
metricstelemetryschema.NewStorage(),
fl,
)
query := qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]{
Signal: telemetrytypes.SignalMetrics,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.MetricAggregation{
{
MetricName: "k8s.pod.cpu.utilization",
Type: metrictypes.GaugeType,
Temporality: metrictypes.Unspecified,
TimeAggregation: metrictypes.TimeAggregationAvg,
SpaceAggregation: metrictypes.SpaceAggregationAvg,
},
},
Filter: &qbtypes.Filter{
Expression: "deployment.environment = 'production'",
},
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
},
},
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, qbtypes.RequestTypeTimeSeries, query, nil)
require.NoError(t, err)
require.Equal(t, "WITH __temporal_aggregation_cte AS (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, JSONExtractString(labels, 'deployment.environment.name') AS `__GROUP_BY_KEY_0_deployment.environment.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND JSONExtractString(labels, 'deployment.environment') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_deployment.environment.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY fingerprint, ts), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_deployment.environment.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_deployment.environment.name`, ts", q.Query)
require.Equal(t, []any{"k8s.pod.cpu.utilization", uint64(1747936800000), uint64(1747983420000), "unspecified", "production", "k8s.pod.cpu.utilization", uint64(1747947390000), uint64(1747983420000), 0}, q.Args)
}

View File

@@ -117,7 +117,9 @@ func (b *StatementBuilder) Build(
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
keySelectors := GetKeySelectors(query)
keySelectors := querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, b.flagger, GetKeySelectors(query))
metricNames := querybuilder.FamilyMetricNames(ctx, orgID, b.flagger, query.Aggregations[0].MetricName)
keySelectors = expandSelectorsForMetricNames(keySelectors, metricNames)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, err
@@ -125,7 +127,30 @@ func (b *StatementBuilder) Build(
start, end = querybuilder.AdjustedMetricTimeRange(start, end, uint64(query.StepInterval.Seconds()), query)
return b.buildPipelineStatement(ctx, orgID, start, end, requestType, query, keys, variables)
return b.buildPipelineStatement(ctx, orgID, start, end, requestType, query, keys, metricNames, variables)
}
// expandSelectorsForMetricNames duplicates the selectors per family metric
// name. Label-key metadata is filtered by the exact metric_name.
func expandSelectorsForMetricNames(selectors []*telemetrytypes.FieldKeySelector, metricNames []string) []*telemetrytypes.FieldKeySelector {
if len(metricNames) <= 1 {
return selectors
}
out := selectors
for _, selector := range selectors {
if selector.MetricContext == nil {
continue
}
for _, metricName := range metricNames {
if metricName == selector.MetricContext.MetricName {
continue
}
expanded := *selector
expanded.MetricContext = &telemetrytypes.MetricContext{MetricName: metricName}
out = append(out, &expanded)
}
}
return out
}
func (b *StatementBuilder) buildPipelineStatement(
@@ -135,6 +160,7 @@ func (b *StatementBuilder) buildPipelineStatement(
requestType qbtypes.RequestType,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
keys map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []string,
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
var (
@@ -166,13 +192,13 @@ func (b *StatementBuilder) buildPipelineStatement(
var filterWarnings []string
var err error
if timeSeriesCTE, timeSeriesCTEArgs, filterWarnings, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, variables, tsTable); err != nil {
if timeSeriesCTE, timeSeriesCTEArgs, filterWarnings, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, metricNames, variables, tsTable); err != nil {
return nil, err
}
if qbtypes.CanShortCircuitDelta(agg) {
// spatial_aggregation_cte directly for certain delta queries
if frag, args, err := b.buildTemporalAggDeltaFastPath(start, end, cteQuery, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
if frag, args, err := b.buildTemporalAggDeltaFastPath(start, end, cteQuery, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
return nil, err
} else if frag != "" {
cteFragments = append(cteFragments, frag)
@@ -180,7 +206,7 @@ func (b *StatementBuilder) buildPipelineStatement(
}
} else {
// temporal_aggregation_cte
if frag, args, err := b.buildTemporalAggregationCTE(ctx, start, end, cteQuery, keys, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
if frag, args, err := b.buildTemporalAggregationCTE(ctx, start, end, cteQuery, keys, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
return nil, err
} else if frag != "" {
cteFragments = append(cteFragments, frag)
@@ -201,16 +227,16 @@ func (b *StatementBuilder) buildPipelineStatement(
var tsArgs []any
// time series rows are written on hour boundaries
tsStart := start - (start % metricstelemetryschema.OneHourInMilliseconds)
if tsCTE, tsArgs, err = b.buildReducedTimeSeriesCTE(ctx, orgID, tsStart, end, cteQuery, keys, variables); err != nil {
if tsCTE, tsArgs, err = b.buildReducedTimeSeriesCTE(ctx, orgID, tsStart, end, cteQuery, keys, metricNames, variables); err != nil {
return nil, err
}
if qbtypes.CanShortCircuitReduced(agg) {
// spatial_aggregation_cte directly, no per-series level
if spatialFrag, spatialArgs, ok := b.buildReducedSpatialAggFastPath(start, end, cteQuery, tsCTE, tsArgs); ok {
if spatialFrag, spatialArgs, ok := b.buildReducedSpatialAggFastPath(start, end, cteQuery, metricNames, tsCTE, tsArgs); ok {
reducedFragments = []string{spatialFrag}
reducedArgs = [][]any{spatialArgs}
}
} else if temporalFrag, temporalArgs, ok := b.buildReducedTemporalAggregationCTE(start, end, cteQuery, tsCTE, tsArgs); ok {
} else if temporalFrag, temporalArgs, ok := b.buildReducedTemporalAggregationCTE(start, end, cteQuery, metricNames, tsCTE, tsArgs); ok {
spatialFrag, spatialArgs := b.buildReducedSpatialAggregationCTE(cteQuery)
reducedFragments = []string{temporalFrag, spatialFrag}
reducedArgs = [][]any{temporalArgs, spatialArgs}
@@ -268,6 +294,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
keys map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []string,
variables map[string]qbtypes.VariableItem,
) (string, []any, error) {
sb := sqlbuilder.NewSelectBuilder()
@@ -300,7 +327,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
sb.SelectMore(sqlbuilder.Escape(fmt.Sprintf("%s AS %s", col, GroupByColumnAlias(i, g.Name))))
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -325,6 +352,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
func (b *StatementBuilder) buildReducedSpatialAggFastPath(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
) (string, []any, bool) {
@@ -345,7 +373,7 @@ func (b *StatementBuilder) buildReducedSpatialAggFastPath(
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", agg.MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -359,6 +387,7 @@ func (b *StatementBuilder) buildReducedSpatialAggFastPath(
func (b *StatementBuilder) buildReducedTemporalAggregationCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
) (string, []any, bool) {
@@ -387,7 +416,7 @@ func (b *StatementBuilder) buildReducedTemporalAggregationCTE(
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", agg.MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -428,6 +457,7 @@ func (b *StatementBuilder) buildReducedSpatialAggregationCTE(
func (b *StatementBuilder) buildTemporalAggDeltaFastPath(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []string,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -469,7 +499,7 @@ func (b *StatementBuilder) buildTemporalAggDeltaFastPath(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -486,6 +516,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
keys map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []string,
variables map[string]qbtypes.VariableItem,
tsTable string,
) (string, []any, []string, error) {
@@ -522,7 +553,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -554,22 +585,24 @@ func (b *StatementBuilder) buildTemporalAggregationCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
_ map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []string,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
) (string, []any, error) {
if query.Aggregations[0].Temporality == metrictypes.Delta {
return b.buildTemporalAggDelta(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
return b.buildTemporalAggDelta(ctx, start, end, query, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
} else if query.Aggregations[0].Temporality != metrictypes.Multiple {
return b.buildTemporalAggCumulativeOrUnspecified(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
return b.buildTemporalAggCumulativeOrUnspecified(ctx, start, end, query, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
}
return b.buildTemporalAggForMultipleTemporalities(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
return b.buildTemporalAggForMultipleTemporalities(ctx, start, end, query, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
}
func (b *StatementBuilder) buildTemporalAggDelta(
_ context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []string,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -601,7 +634,7 @@ func (b *StatementBuilder) buildTemporalAggDelta(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -617,6 +650,7 @@ func (b *StatementBuilder) buildTemporalAggCumulativeOrUnspecified(
_ context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []string,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -642,7 +676,7 @@ func (b *StatementBuilder) buildTemporalAggCumulativeOrUnspecified(
baseSb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
baseSb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
baseSb.Where(
baseSb.In("metric_name", query.Aggregations[0].MetricName),
baseSb.In("metric_name", sqlbuilder.List(metricNames)),
baseSb.GTE("unix_milli", start),
baseSb.LT("unix_milli", end),
)
@@ -683,6 +717,7 @@ func (b *StatementBuilder) buildTemporalAggForMultipleTemporalities(
_ context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []string,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -733,7 +768,7 @@ func (b *StatementBuilder) buildTemporalAggForMultipleTemporalities(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", sqlbuilder.List(metricNames)),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)

View File

@@ -2062,3 +2062,51 @@ func TestStatementBuilderSemconvFamilies(t *testing.T) {
})
}
}
// The mid-migration state: metadata holds only the old spelling, and the
// query names the current one. The resource-filter condition and the group
// by column both read the stored spelling, so the filter and the groups
// agree.
func TestStatementBuilderSemconvSingleSpelling(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
storage := tracestelemetryschema.NewStorage()
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
mockMetadataStore.KeysMap = map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment": {{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
}
aggExprRewriter := querybuilder.NewAggExprRewriter(instrumentationtest.New().ToProviderSettings(), nil, storage, fl, telemetrytypes.SignalTraces)
statementBuilder := NewTraceQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore,
storage,
aggExprRewriter,
nil,
fl,
false,
100000,
)
query := qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.TraceAggregation{{Expression: "count()"}},
Filter: &qbtypes.Filter{
Expression: "deployment.environment.name = 'production'",
},
GroupBy: []qbtypes.GroupByKey{
{TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}},
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, qbtypes.RequestTypeScalar, query, nil)
require.NoError(t, err)
require.Equal(t, "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'deployment.environment') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT toString(multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, mapContains(resources_string, 'deployment.environment'), resources_string['deployment.environment'], NULL)) AS `__GROUP_BY_KEY_0_deployment.environment.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY __result_0 DESC", q.Query)
}

View File

@@ -0,0 +1,150 @@
package telemetrymetadata
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/querybuilder"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// The flagger provider registration is process-global and keyed by provider
// name, so each flagger must be used before the next one is created.
func TestFamilyValueNames(t *testing.T) {
selector := &telemetrytypes.FieldValueSelector{
FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "deployment.environment"},
}
off := &telemetryMetaStore{fl: flaggertest.WithBooleanFlags(t, map[string]bool{})}
assert.Equal(t,
[]string{"deployment.environment"},
off.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalLogs, selector),
"the flag default keeps values literal")
on := &telemetryMetaStore{fl: flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})}
assert.Equal(t,
[]string{"deployment.environment.name", "deployment.environment"},
on.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalLogs, selector),
"values for one spelling must cover the whole family")
assert.Equal(t,
[]string{"deployment.environment.name", "deployment.environment"},
on.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalMetrics, selector),
"metric values cover the family members")
spanMetric := &telemetrytypes.FieldValueSelector{
FieldKeySelector: &telemetrytypes.FieldKeySelector{
Name: "deployment.environment",
MetricContext: &telemetrytypes.MetricContext{MetricName: "signoz_calls_total"},
},
}
assert.Equal(t,
[]string{
"deployment.environment.name", "resource_deployment.environment.name",
"deployment.environment", "resource_deployment.environment",
},
on.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalMetrics, spanMetric),
"span-metrics values cover the resource_ spellings too")
}
// A family condition on the related values table follows the shared guard
// rule: the operator applies to the current-first merge, a positive operator
// takes the presence guard, and a negative operator keeps the keyless rows.
func TestConditionForFamilyMergedSemantics(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
storage := NewStorage()
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
"deployment.environment": {{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
}
testCases := []struct {
name string
operator qbtypes.FilterOperator
expected string
}{
{
name: "Equal_TakesPresenceGuard",
operator: qbtypes.FilterOperatorEqual,
expected: "SELECT 1 WHERE (COALESCE(NULLIF(resource_attributes['deployment.environment.name'], ''), NULLIF(resource_attributes['deployment.environment'], ''), '') = ? AND (mapContains(resource_attributes, 'deployment.environment.name') OR mapContains(resource_attributes, 'deployment.environment')))",
},
{
name: "NotEqual_KeepsKeylessRows",
operator: qbtypes.FilterOperatorNotEqual,
expected: "SELECT 1 WHERE COALESCE(NULLIF(resource_attributes['deployment.environment.name'], ''), NULLIF(resource_attributes['deployment.environment'], ''), '') <> ?",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
sb := sqlbuilder.NewSelectBuilder()
q := querybuilder.NewQueryInfo(context.Background(), valuer.UUID{}, fl, telemetrytypes.SignalTraces, nil, 0, 0)
conds, _, err := querybuilder.Conditions(context.Background(), q, storage,
&telemetrytypes.TelemetryFieldKey{Name: "deployment.environment"}, testCase.operator, "production", fieldKeys, false, sb)
require.NoError(t, err)
sb.Select("1").Where(conds...)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Equal(t, testCase.expected, sql)
})
}
}
// The related values search applies to the current-first merge that the
// suggestion shows, also when the request carries no data type.
func TestContainsConditionsSearchesFamilyAsOneField(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
q := querybuilder.NewQueryInfo(context.Background(), valuer.UUID{}, fl, telemetrytypes.SignalTraces, nil, 0, 0)
store := &telemetryMetaStore{storage: NewStorage()}
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
}
testCases := []struct {
name string
names []string
expected string
}{
{
name: "FamilySpellings_MergeIntoOneCondition",
names: []string{"deployment.environment.name", "deployment.environment"},
expected: "SELECT 1 WHERE (LOWER(COALESCE(NULLIF(resource_attributes['deployment.environment.name'], ''), NULLIF(resource_attributes['deployment.environment'], ''), '')) LIKE LOWER(?) AND (mapContains(resource_attributes, 'deployment.environment.name') OR mapContains(resource_attributes, 'deployment.environment')))",
},
{
name: "SingleSpelling_StaysLiteral",
names: []string{"deployment.environment"},
expected: "SELECT 1 WHERE if(mapContains(resource_attributes, ?), LOWER(resource_attributes['deployment.environment']) LIKE LOWER(?), false)",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
sb := sqlbuilder.NewSelectBuilder()
conds, err := store.containsConditions(context.Background(), q, key, testCase.names, "prod", sb)
require.NoError(t, err)
sb.Select("1").Where(conds...)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Equal(t, testCase.expected, sql)
})
}
}

View File

@@ -15,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"
@@ -1292,25 +1293,41 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
FieldDataType: fieldValueSelector.FieldDataType,
}
q := querybuilder.NewQueryInfo(ctx, orgID, nil, fieldValueSelector.Signal, nil, 0, 0)
selectRead, err := t.storage.Read(ctx, q, key)
selectColumn := selectRead.SQL
if err != nil {
// we don't have a explicit column to select from the related metadata table
// so we will select either from resource_attributes or attributes table
// in that order
resourceRead, _ := t.storage.Read(ctx, q, &telemetrytypes.TelemetryFieldKey{
Name: key.Name,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
attributeRead, _ := t.storage.Read(ctx, q, &telemetrytypes.TelemetryFieldKey{
Name: key.Name,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
selectColumn = fmt.Sprintf("if(notEmpty(%s), %s, %s)", resourceRead.SQL, resourceRead.SQL, attributeRead.SQL)
q := querybuilder.NewQueryInfo(ctx, orgID, t.fl, fieldValueSelector.Signal, nil, 0, 0)
// One column per family spelling, merged current-first, so the
// suggestions cover rows that carry only an old spelling.
names := t.familyValueNames(ctx, orgID, fieldValueSelector.Signal, fieldValueSelector)
memberColumns := make([]string, 0, len(names))
for _, name := range names {
memberKey := &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: fieldValueSelector.Signal,
FieldContext: fieldValueSelector.FieldContext,
FieldDataType: fieldValueSelector.FieldDataType,
}
memberRead, err := t.storage.Read(ctx, q, memberKey)
memberColumn := memberRead.SQL
if err != nil {
// we don't have a explicit column to select from the related metadata table
// so we will select either from resource_attributes or attributes table
// in that order
resourceRead, _ := t.storage.Read(ctx, q, &telemetrytypes.TelemetryFieldKey{
Name: name,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
attributeRead, _ := t.storage.Read(ctx, q, &telemetrytypes.TelemetryFieldKey{
Name: name,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
memberColumn = fmt.Sprintf("if(notEmpty(%s), %s, %s)", resourceRead.SQL, resourceRead.SQL, attributeRead.SQL)
}
memberColumns = append(memberColumns, memberColumn)
}
selectColumn := memberColumns[len(memberColumns)-1]
for i := len(memberColumns) - 2; i >= 0; i-- {
selectColumn = fmt.Sprintf("if(notEmpty(%s), %s, %s)", memberColumns[i], memberColumns[i], selectColumn)
}
sb := sqlbuilder.Select("DISTINCT " + selectColumn).From(t.relatedMetadataDBName + "." + t.relatedMetadataTblName)
@@ -1320,6 +1337,7 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
for _, keySelector := range keySelectors {
keySelector.Signal = fieldValueSelector.Signal
}
keySelectors = querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, t.fl, keySelectors)
keys, _, err := t.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, false, err
@@ -1361,20 +1379,20 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
// search on attributes
key.FieldContext = telemetrytypes.FieldContextAttribute
attrConds, err := t.containsConditions(ctx, q, key, fieldValueSelector.Value, sb)
attrConds, err := t.containsConditions(ctx, q, key, names, fieldValueSelector.Value, sb)
if err == nil {
conds = append(conds, attrConds...)
}
// search on resource
key.FieldContext = telemetrytypes.FieldContextResource
resourceConds, err := t.containsConditions(ctx, q, key, fieldValueSelector.Value, sb)
resourceConds, err := t.containsConditions(ctx, q, key, names, fieldValueSelector.Value, sb)
if err == nil {
conds = append(conds, resourceConds...)
}
key.FieldContext = origContext
} else {
keyConds, err := t.containsConditions(ctx, q, key, fieldValueSelector.Value, sb)
keyConds, err := t.containsConditions(ctx, q, key, names, fieldValueSelector.Value, sb)
if err == nil {
conds = append(conds, keyConds...)
}
@@ -1432,7 +1450,7 @@ func (t *telemetryMetaStore) GetRelatedValues(ctx context.Context, orgID valuer.
return t.getRelatedValues(ctx, orgID, fieldValueSelector)
}
func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalTraces.StringValue(),
instrumentationtypes.CodeNamespace: "metadata",
@@ -1443,11 +1461,12 @@ func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueS
return values, true, nil
}
knownBool := isKnownBoolField(fieldValueSelector, tracestelemetryschema.IntrinsicFields, tracestelemetryschema.CalculatedFields)
names := t.familyValueNames(ctx, orgID, telemetrytypes.SignalTraces, fieldValueSelector)
// unix_milli is the hour of the span start
return t.getTagTableValues(ctx, t.tracesDBName+"."+t.tracesFieldsTblName, fieldValueSelector, knownBool)
return t.getTagTableValues(ctx, t.tracesDBName+"."+t.tracesFieldsTblName, fieldValueSelector, names, knownBool)
}
func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalLogs.StringValue(),
instrumentationtypes.CodeNamespace: "metadata",
@@ -1455,8 +1474,18 @@ func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, fieldValueSe
})
knownBool := isKnownBoolField(fieldValueSelector, logstelemetryschema.IntrinsicFields)
names := t.familyValueNames(ctx, orgID, telemetrytypes.SignalLogs, fieldValueSelector)
// unix_milli is the hour the log was ingested, not the log's own timestamp
return t.getTagTableValues(ctx, t.logsDBName+"."+t.logsFieldsTblName, fieldValueSelector, knownBool)
return t.getTagTableValues(ctx, t.logsDBName+"."+t.logsFieldsTblName, fieldValueSelector, names, knownBool)
}
// tagKeyCondition matches the requested key, or every spelling of its
// family when there is more than one.
func tagKeyCondition(sb *sqlbuilder.SelectBuilder, name string, names []string) string {
if len(names) > 1 {
return sb.In("tag_key", sqlbuilder.List(names))
}
return sb.E("tag_key", name)
}
// tagTableSinceDay restricts rows to the tag table's day partitions from the
@@ -1472,9 +1501,9 @@ func tagTableSinceDay(sb *sqlbuilder.SelectBuilder, startUnixMilli int64) {
// tagTableHasBoolRows reports whether the tag table holds a bool row for the
// key. Bool rows carry no value, so one row is enough to know the key takes
// the values true and false.
func (t *telemetryMetaStore) tagTableHasBoolRows(ctx context.Context, table string, selector *telemetrytypes.FieldValueSelector) (bool, error) {
func (t *telemetryMetaStore) tagTableHasBoolRows(ctx context.Context, table string, selector *telemetrytypes.FieldValueSelector, names []string) (bool, error) {
sb := sqlbuilder.Select("1").From(table)
sb.Where(sb.E("tag_key", selector.Name))
sb.Where(tagKeyCondition(sb, selector.Name, names))
sb.Where(sb.E("tag_data_type", telemetrytypes.FieldDataTypeBool.TagDataType()))
if selector.FieldContext != telemetrytypes.FieldContextUnspecified {
sb.Where(sb.E("tag_type", selector.FieldContext.TagType()))
@@ -1494,7 +1523,7 @@ func (t *telemetryMetaStore) tagTableHasBoolRows(ctx context.Context, table stri
// getTagTableValues returns the string and number values of the key from a
// tag table, and true and false when the key is a known bool field or the
// table holds bool rows for it. Bool rows do not count towards the limit.
func (t *telemetryMetaStore) getTagTableValues(ctx context.Context, table string, fieldValueSelector *telemetrytypes.FieldValueSelector, knownBool bool) (*telemetrytypes.TelemetryFieldValues, bool, error) {
func (t *telemetryMetaStore) getTagTableValues(ctx context.Context, table string, fieldValueSelector *telemetrytypes.FieldValueSelector, names []string, knownBool bool) (*telemetrytypes.TelemetryFieldValues, bool, error) {
limit := fieldValueSelector.Limit
if limit == 0 {
limit = 50
@@ -1507,7 +1536,7 @@ func (t *telemetryMetaStore) getTagTableValues(ctx context.Context, table string
return values, true, nil
}
} else if fieldValueSelector.FieldDataType == telemetrytypes.FieldDataTypeUnspecified {
hasBoolRows, err := t.tagTableHasBoolRows(ctx, table, fieldValueSelector)
hasBoolRows, err := t.tagTableHasBoolRows(ctx, table, fieldValueSelector, names)
if err != nil {
return nil, false, err
}
@@ -1519,7 +1548,7 @@ func (t *telemetryMetaStore) getTagTableValues(ctx context.Context, table string
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(table)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
sb.Where(tagKeyCondition(sb, fieldValueSelector.Name, names))
}
sb.Where(sb.NE("tag_data_type", telemetrytypes.FieldDataTypeBool.TagDataType()))
@@ -1719,7 +1748,11 @@ func (t *telemetryMetaStore) getMetricFieldValues(ctx context.Context, orgID val
From(t.metricsDBName + "." + t.metricsFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("attr_name", fieldValueSelector.Name))
if names := t.familyValueNames(ctx, orgID, telemetrytypes.SignalMetrics, fieldValueSelector); len(names) > 1 {
sb.Where(sb.In("attr_name", sqlbuilder.List(names)))
} else {
sb.Where(sb.E("attr_name", fieldValueSelector.Name))
}
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1731,7 +1764,11 @@ func (t *telemetryMetaStore) getMetricFieldValues(ctx context.Context, orgID val
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
if metricNames := querybuilder.FamilyMetricNames(ctx, orgID, t.fl, fieldValueSelector.MetricContext.MetricName); len(metricNames) > 1 {
sb.Where(sb.In("metric_name", sqlbuilder.List(metricNames)))
} else {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
}
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", clickhousesql.LikePattern(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2053,12 +2090,12 @@ func (t *telemetryMetaStore) GetAllValues(ctx context.Context, orgID valuer.UUID
switch fieldValueSelector.Signal {
case telemetrytypes.SignalTraces:
values, complete, err = t.getSpanFieldValues(ctx, fieldValueSelector)
values, complete, err = t.getSpanFieldValues(ctx, orgID, fieldValueSelector)
case telemetrytypes.SignalLogs:
if fieldValueSelector.Source == telemetrytypes.SourceAudit {
values, complete, err = t.getAuditFieldValues(ctx, fieldValueSelector)
} else {
values, complete, err = t.getLogFieldValues(ctx, fieldValueSelector)
values, complete, err = t.getLogFieldValues(ctx, orgID, fieldValueSelector)
}
case telemetrytypes.SignalMetrics:
if fieldValueSelector.Source == telemetrytypes.SourceMeter {
@@ -2071,13 +2108,13 @@ func (t *telemetryMetaStore) GetAllValues(ctx context.Context, orgID valuer.UUID
mapOfRelatedValues := make(map[any]bool)
allUnspecifiedValues := &telemetrytypes.TelemetryFieldValues{}
tracesValues, tracesComplete, err := t.getSpanFieldValues(ctx, fieldValueSelector)
tracesValues, tracesComplete, err := t.getSpanFieldValues(ctx, orgID, fieldValueSelector)
if err == nil {
populateComplete := populateAllUnspecifiedValues(allUnspecifiedValues, mapOfValues, mapOfRelatedValues, tracesValues, limit)
complete = complete && tracesComplete && populateComplete
}
logsValues, logsComplete, err := t.getLogFieldValues(ctx, fieldValueSelector)
logsValues, logsComplete, err := t.getLogFieldValues(ctx, orgID, fieldValueSelector)
if err == nil {
populateComplete := populateAllUnspecifiedValues(allUnspecifiedValues, mapOfValues, mapOfRelatedValues, logsValues, limit)
complete = complete && logsComplete && populateComplete
@@ -2589,10 +2626,42 @@ func (t *telemetryMetaStore) fetchLastSeenInfoForTable(ctx context.Context, tabl
return lastSeenInfo, nil
}
// containsConditions compiles a contains search on one key of the related
// values table. The key is its own metadata.
func (t *telemetryMetaStore) containsConditions(ctx context.Context, q qbtypes.QueryInfo, key *telemetrytypes.TelemetryFieldKey, value string, sb *sqlbuilder.SelectBuilder) ([]string, error) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}}
// containsConditions compiles a contains search over the key and its family
// spellings, which stand as their own metadata.
func (t *telemetryMetaStore) containsConditions(ctx context.Context, q qbtypes.QueryInfo, key *telemetrytypes.TelemetryFieldKey, names []string, value string, sb *sqlbuilder.SelectBuilder) ([]string, error) {
// A family forms only over string keys with a signal. The table stores
// only string maps, so an unspecified data type is a string here.
dataType := key.FieldDataType
if dataType == telemetrytypes.FieldDataTypeUnspecified {
dataType = telemetrytypes.FieldDataTypeString
}
fieldKeys := make(map[string][]*telemetrytypes.TelemetryFieldKey, len(names))
for _, name := range names {
fieldKeys[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: key.Signal,
FieldContext: key.FieldContext,
FieldDataType: dataType,
}}
}
conds, _, err := querybuilder.Conditions(ctx, q, t.storage, key, qbtypes.FilterOperatorContains, value, fieldKeys, false, sb)
return conds, err
}
// familyValueNames returns the spellings whose values merge into the
// suggestions. With the flag off, the requested name alone.
func (t *telemetryMetaStore) familyValueNames(ctx context.Context, orgID valuer.UUID, signal telemetrytypes.Signal, fieldValueSelector *telemetrytypes.FieldValueSelector) []string {
if !querybuilder.SemconvFamiliesEnabled(ctx, orgID, t.fl) {
return []string{fieldValueSelector.Name}
}
selector := telemetrytypes.FieldKeySelector{
Name: fieldValueSelector.Name,
Signal: signal,
FieldContext: fieldValueSelector.FieldContext,
MetricContext: fieldValueSelector.MetricContext,
}
if signal == telemetrytypes.SignalMetrics {
return querybuilder.MetricLabelSpellings(selector)
}
return semconv.Members(semconv.KindAttribute, selector)
}

View File

@@ -122,6 +122,9 @@ type PipelineOperator struct {
EnablePaths bool `json:"enable_paths,omitempty" yaml:"enable_paths,omitempty"`
PathPrefix string `json:"path_prefix,omitempty" yaml:"path_prefix,omitempty"`
// normalize fields, set by the server only
JSONBodyDualIngestion bool `json:"-" yaml:"json_body_dual_ingestion,omitempty"`
// Used in Severity Parsing and JSON Flattening mapping
Mapping map[string][]string `json:"mapping,omitempty" yaml:"mapping,omitempty"`
// severity parser fields

View File

@@ -8,6 +8,7 @@ import (
"go/format"
"os"
"path/filepath"
"slices"
"sort"
"strconv"
"strings"
@@ -64,6 +65,8 @@ type overlayFile struct {
Families map[string]overlayFamily `yaml:"families"`
}
// Contexts and Signals set the family-level gate. The Add* fields widen the
// member scopes the schema derived.
type overlayFamily struct {
Enabled *bool `yaml:"enabled"`
Kind string `yaml:"kind"`
@@ -79,27 +82,36 @@ type overlayFamily struct {
ValueMap map[string]string `yaml:"value_map"`
}
type edge struct {
old string
current string
kind string
// scope is the constraint a rename edge carries. A nil axis is unconstrained.
type scope struct {
contexts []string
signals []string
allContexts bool
allSignals bool
applyToMetrics []string
}
type edge struct {
old string
current string
kind string
scope scope
}
type graphKey struct{ kind, name string }
type generatedFamily struct {
Current string
Old []string
Kind string
type generatedMember struct {
Name string
Contexts []string
Signals []string
ApplyToMetrics []string
ValueMap map[string]string
}
type generatedFamily struct {
Current string
Kind string
Members []generatedMember
Contexts []string
Signals []string
ValueMap map[string]string
}
func main() {
@@ -232,25 +244,31 @@ func collectEdges(schemas []schemaFile) ([]edge, error) {
{name: "metrics", section: version.Metrics},
}
for _, scoped := range sections {
contexts, signals, allContexts, allSignals, err := scopeForSection(scoped.name)
sectionScope, err := scopeForSection(scoped.name)
if err != nil {
return nil, err
}
for _, change := range scoped.section.Changes {
if change.RenameAttributes != nil {
if change.RenameAttributes.ApplyToMetrics != nil && len(change.RenameAttributes.ApplyToMetrics) == 0 {
return nil, fmt.Errorf(
"schema version %q has an explicitly empty apply_to_metrics; an empty list would be emitted as unconstrained",
versionName,
)
}
edgeScope := sectionScope
edgeScope.applyToMetrics = change.RenameAttributes.ApplyToMetrics
for _, old := range sortedMapKeys(change.RenameAttributes.AttributeMap) {
versionEdges = append(versionEdges, edge{
old: old, current: change.RenameAttributes.AttributeMap[old], kind: kindAttribute,
contexts: contexts, signals: signals,
allContexts: allContexts, allSignals: allSignals,
applyToMetrics: change.RenameAttributes.ApplyToMetrics,
scope: edgeScope,
})
}
}
for _, old := range sortedMapKeys(change.RenameMetrics) {
versionEdges = append(versionEdges, edge{
old: old, current: change.RenameMetrics[old], kind: kindMetric,
contexts: []string{"metric"}, signals: []string{"metrics"},
scope: scope{contexts: []string{"metric"}, signals: []string{"metrics"}},
})
}
}
@@ -311,99 +329,177 @@ func compareVersionParts(left, right [3]int) int {
return 0
}
func scopeForSection(section string) (contexts, signals []string, allContexts, allSignals bool, err error) {
func scopeForSection(section string) (scope, error) {
switch section {
case "all":
return nil, nil, true, true, nil
return scope{}, nil
case "resources":
return []string{"resource"}, nil, false, true, nil
return scope{contexts: []string{"resource"}}, nil
case "spans":
return []string{"attribute"}, []string{"traces"}, false, false, nil
return scope{contexts: []string{"attribute"}, signals: []string{"traces"}}, nil
case "logs":
return []string{"attribute"}, []string{"logs"}, false, false, nil
return scope{contexts: []string{"attribute"}, signals: []string{"logs"}}, nil
case "metrics":
return []string{"attribute"}, []string{"metrics"}, false, false, nil
return scope{contexts: []string{"attribute"}, signals: []string{"metrics"}}, nil
default:
return nil, nil, false, false, fmt.Errorf("unsupported schema section %q", section)
return scope{}, fmt.Errorf("unsupported schema section %q", section)
}
}
// unionScope widens per axis. A nil side widens to nil. The metric-name axis
// takes only edges that can apply to metrics, so a non-metrics rename does
// not erase a scoped list.
func unionScope(left, right scope) scope {
return scope{
contexts: unionAxis(left.contexts, right.contexts),
signals: unionAxis(left.signals, right.signals),
applyToMetrics: applyToMetricsUnion(left, right),
}
}
func applyToMetricsUnion(left, right scope) []string {
// A nil signal axis admits every signal, metrics included.
leftApplies := left.signals == nil || slices.Contains(left.signals, "metrics")
rightApplies := right.signals == nil || slices.Contains(right.signals, "metrics")
switch {
case leftApplies && rightApplies:
return unionAxis(left.applyToMetrics, right.applyToMetrics)
case leftApplies:
return sortedCopy(left.applyToMetrics)
case rightApplies:
return sortedCopy(right.applyToMetrics)
}
return nil
}
func unionAxis(left, right []string) []string {
if left == nil || right == nil {
return nil
}
merged := appendUnique(append([]string(nil), left...), right...)
sort.Strings(merged)
return merged
}
// pathResult is one path from a name to a family root. Hops carry only
// reachability. A member's scope comes from its own rename edges, because a
// later rename can be filed under another section.
type pathResult struct {
root string
distance int
}
func rootsFor(next map[graphKey][]edge, kind, name string, distance int, seen map[string]bool) ([]pathResult, error) {
if seen[name] {
return nil, fmt.Errorf("rename cycle for %s %q", kind, name)
}
outgoing := next[graphKey{kind: kind, name: name}]
if len(outgoing) == 0 {
return []pathResult{{root: name, distance: distance}}, nil
}
seen[name] = true
defer delete(seen, name)
var results []pathResult
for _, hop := range outgoing {
hopResults, err := rootsFor(next, kind, hop.current, distance+1, seen)
if err != nil {
return nil, err
}
results = append(results, hopResults...)
}
return results, nil
}
func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily, error) {
edges, err := collectEdges(schemas)
if err != nil {
return nil, err
}
next := make(map[graphKey]string)
// An old name can fan out into several families, so the graph keeps every
// successor.
next := make(map[graphKey][]edge)
for _, item := range edges {
key := graphKey{kind: item.kind, name: item.old}
if existing, ok := next[key]; ok && existing == item.current {
// Repeated entries are common in chained schema histories. Treat an
// identical edge as a no-op so it cannot sever a later edge in the
// same chain (A -> B, B -> C, then a repeated A -> B).
merged := false
for i, existing := range next[key] {
if existing.current == item.current {
// A repeated edge merges its scope, so it cannot sever a later edge
// in the chain.
next[key][i].scope = unionScope(existing.scope, item.scope)
merged = true
break
}
}
if merged {
continue
}
// Schema history occasionally repeats an old name with a newer direct
// destination or rolls a rename back. Edges are collected
// oldest-to-newest, so the latest published current name must be a root.
// Edges run oldest to newest, so the latest published current name
// must be a root. A family only reachable through it is orphaned,
// which is correct for the rollbacks the vendored history holds.
delete(next, graphKey{kind: item.kind, name: item.current})
next[key] = item.current
next[key] = append(next[key], item)
}
type memberState struct {
sc scope
distance int
}
type familyState struct {
family generatedFamily
distance map[string]int
allContexts bool
allSignals bool
family generatedFamily
members map[string]*memberState
}
states := map[graphKey]*familyState{}
for _, item := range edges {
root, distance, err := rootFor(next, item.kind, item.old)
if err != nil {
return nil, err
}
key := graphKey{kind: item.kind, name: root}
state := states[key]
if state == nil {
state = &familyState{
family: generatedFamily{Current: root, Kind: item.kind},
distance: map[string]int{},
for key, outgoing := range next {
for _, item := range outgoing {
results, err := rootsFor(next, key.kind, item.current, 1, map[string]bool{key.name: true})
if err != nil {
return nil, err
}
for _, result := range results {
rootKey := graphKey{kind: key.kind, name: result.root}
state := states[rootKey]
if state == nil {
state = &familyState{
family: generatedFamily{Current: result.root, Kind: key.kind},
members: map[string]*memberState{},
}
states[rootKey] = state
}
member := state.members[key.name]
if member == nil {
state.members[key.name] = &memberState{sc: item.scope, distance: result.distance}
continue
}
member.sc = unionScope(member.sc, item.scope)
if result.distance < member.distance {
member.distance = result.distance
}
}
states[key] = state
}
if prior, ok := state.distance[item.old]; !ok || distance < prior {
state.distance[item.old] = distance
}
state.allContexts = state.allContexts || item.allContexts
state.allSignals = state.allSignals || item.allSignals
state.family.Contexts = appendUnique(state.family.Contexts, item.contexts...)
state.family.Signals = appendUnique(state.family.Signals, item.signals...)
state.family.ApplyToMetrics = appendUnique(state.family.ApplyToMetrics, item.applyToMetrics...)
}
for _, state := range states {
for old := range state.distance {
if old != state.family.Current {
state.family.Old = append(state.family.Old, old)
}
names := make([]string, 0, len(state.members))
for name := range state.members {
names = append(names, name)
}
sort.Slice(state.family.Old, func(i, j int) bool {
left, right := state.family.Old[i], state.family.Old[j]
if state.distance[left] != state.distance[right] {
return state.distance[left] < state.distance[right]
sort.Slice(names, func(i, j int) bool {
left, right := state.members[names[i]], state.members[names[j]]
if left.distance != right.distance {
return left.distance < right.distance
}
return left < right
return names[i] < names[j]
})
if state.allContexts {
state.family.Contexts = nil
} else {
sort.Strings(state.family.Contexts)
for _, name := range names {
member := state.members[name]
state.family.Members = append(state.family.Members, generatedMember{
Name: name,
Contexts: sortedCopy(member.sc.contexts),
Signals: sortedCopy(member.sc.signals),
ApplyToMetrics: sortedCopy(member.sc.applyToMetrics),
})
}
if state.allSignals {
state.family.Signals = nil
} else {
sort.Strings(state.family.Signals)
}
sort.Strings(state.family.ApplyToMetrics)
}
for _, current := range sortedMapKeys(overlay.Families) {
@@ -414,6 +510,12 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
}
policy.Kind = kind
overlay.Families[current] = policy
if policy.ApplyToMetrics != nil && len(policy.ApplyToMetrics) == 0 {
return nil, fmt.Errorf(
"overlay family %q has an explicitly empty apply_to_metrics; an empty list would be emitted as unconstrained",
current,
)
}
key := graphKey{kind: kind, name: current}
state := states[key]
if state == nil {
@@ -424,10 +526,7 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
kind,
)
}
state = &familyState{
family: generatedFamily{Current: current, Kind: kind, Old: append([]string(nil), policy.Old...)},
distance: map[string]int{},
}
state = &familyState{family: generatedFamily{Current: current, Kind: kind}}
states[key] = state
}
applyOverlay(&state.family, policy)
@@ -446,7 +545,7 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
if !enabled {
continue
}
if len(state.family.Old) == 0 {
if len(state.family.Members) == 0 {
return nil, fmt.Errorf(
"enabled family %q with kind %q has no old members",
state.family.Current,
@@ -455,7 +554,6 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
}
sort.Strings(state.family.Contexts)
sort.Strings(state.family.Signals)
sort.Strings(state.family.ApplyToMetrics)
result = append(result, state.family)
}
@@ -468,21 +566,13 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
return result, nil
}
func rootFor(next map[graphKey]string, kind, name string) (string, int, error) {
seen := map[string]bool{}
distance := 0
for {
if seen[name] {
return "", 0, fmt.Errorf("rename cycle for %s %q", kind, name)
}
seen[name] = true
current, ok := next[graphKey{kind: kind, name: name}]
if !ok {
return name, distance, nil
}
name = current
distance++
func sortedCopy(values []string) []string {
if values == nil {
return nil
}
out := append([]string(nil), values...)
sort.Strings(out)
return out
}
func normalizedOverlayKind(current string, policy overlayFamily) (string, error) {
@@ -501,15 +591,29 @@ func applyOverlay(family *generatedFamily, policy overlayFamily) {
family.Kind = policy.Kind
}
if policy.Old != nil {
family.Old = append([]string(nil), policy.Old...)
family.Members = nil
for _, old := range policy.Old {
family.Members = append(family.Members, generatedMember{Name: old})
}
}
for _, old := range policy.AddOld {
if slices.ContainsFunc(family.Members, func(member generatedMember) bool { return member.Name == old }) {
continue
}
family.Members = append(family.Members, generatedMember{Name: old})
}
family.Old = appendUnique(family.Old, policy.AddOld...)
if len(policy.ExcludeOld) > 0 {
excluded := make(map[string]bool, len(policy.ExcludeOld))
for _, old := range policy.ExcludeOld {
excluded[old] = true
}
family.Old = deleteMatching(family.Old, excluded)
kept := family.Members[:0]
for _, member := range family.Members {
if !excluded[member.Name] {
kept = append(kept, member)
}
}
family.Members = kept
}
if policy.Contexts != nil {
family.Contexts = append([]string(nil), policy.Contexts...)
@@ -517,12 +621,27 @@ func applyOverlay(family *generatedFamily, policy overlayFamily) {
if policy.Signals != nil {
family.Signals = append([]string(nil), policy.Signals...)
}
family.Contexts = appendUnique(family.Contexts, policy.AddContexts...)
family.Signals = appendUnique(family.Signals, policy.AddSignals...)
if policy.ApplyToMetrics != nil {
family.ApplyToMetrics = append([]string(nil), policy.ApplyToMetrics...)
// Add* fields widen a gate the overlay set. A nil gate already admits all.
if family.Contexts != nil {
family.Contexts = appendUnique(family.Contexts, policy.AddContexts...)
}
if family.Signals != nil {
family.Signals = appendUnique(family.Signals, policy.AddSignals...)
}
for i := range family.Members {
if len(policy.AddContexts) > 0 {
family.Members[i].Contexts = unionAxis(family.Members[i].Contexts, sortedCopy(policy.AddContexts))
}
if len(policy.AddSignals) > 0 {
family.Members[i].Signals = unionAxis(family.Members[i].Signals, sortedCopy(policy.AddSignals))
}
if policy.ApplyToMetrics != nil {
family.Members[i].ApplyToMetrics = sortedCopy(policy.ApplyToMetrics)
}
if len(policy.AddApplyToMetrics) > 0 && family.Members[i].ApplyToMetrics != nil {
family.Members[i].ApplyToMetrics = unionAxis(family.Members[i].ApplyToMetrics, sortedCopy(policy.AddApplyToMetrics))
}
}
family.ApplyToMetrics = appendUnique(family.ApplyToMetrics, policy.AddApplyToMetrics...)
if policy.ValueMap != nil {
family.ValueMap = make(map[string]string, len(policy.ValueMap))
for old, current := range policy.ValueMap {
@@ -546,16 +665,6 @@ func appendUnique(values []string, additions ...string) []string {
return values
}
func deleteMatching(values []string, excluded map[string]bool) []string {
result := values[:0]
for _, value := range values {
if !excluded[value] {
result = append(result, value)
}
}
return result
}
func renderGo(families []generatedFamily) ([]byte, error) {
var out bytes.Buffer
out.WriteString("// Code generated by scripts/semconv. DO NOT EDIT.\n\n")
@@ -564,7 +673,11 @@ func renderGo(families []generatedFamily) ([]byte, error) {
for _, family := range families {
if len(family.Contexts) > 0 || len(family.Signals) > 0 {
needsTelemetryTypes = true
break
}
for _, member := range family.Members {
if len(member.Contexts) > 0 || len(member.Signals) > 0 {
needsTelemetryTypes = true
}
}
}
if needsTelemetryTypes {
@@ -572,32 +685,36 @@ func renderGo(families []generatedFamily) ([]byte, error) {
}
out.WriteString("var families = []Family{\n")
for _, family := range families {
contexts, err := goFieldContextSlice(family.Contexts)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
signals, err := goSignalSlice(family.Signals)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
out.WriteString("\t{\n")
fmt.Fprintf(&out, "\t\tCurrent: %s,\n", strconv.Quote(family.Current))
fmt.Fprintf(&out, "\t\tOld: %s,\n", goStringSlice(family.Old))
fmt.Fprintf(&out, "\t\tcurrent: %s,\n", strconv.Quote(family.Current))
if family.Kind == kindMetric {
out.WriteString("\t\tKind: KindMetric,\n")
out.WriteString("\t\tkind: KindMetric,\n")
} else {
out.WriteString("\t\tKind: KindAttribute,\n")
out.WriteString("\t\tkind: KindAttribute,\n")
}
fmt.Fprintf(&out, "\t\tContexts: %s,\n", contexts)
fmt.Fprintf(&out, "\t\tSignals: %s,\n", signals)
fmt.Fprintf(&out, "\t\tApplyToMetrics: %s,\n", goStringSlice(family.ApplyToMetrics))
if len(family.ValueMap) > 0 {
out.WriteString("\t\tValueMap: map[string]string{\n")
keys := sortedMapKeys(family.ValueMap)
for _, key := range keys {
fmt.Fprintf(&out, "\t\t\t%s: %s,\n", strconv.Quote(key), strconv.Quote(family.ValueMap[key]))
out.WriteString("\t\tmembers: []Member{\n")
for _, member := range family.Members {
if err := writeGoMember(&out, family.Current, member); err != nil {
return nil, err
}
out.WriteString("\t\t},\n")
}
out.WriteString("\t\t},\n")
if len(family.Contexts) > 0 {
contexts, err := goFieldContextSlice(family.Contexts)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
fmt.Fprintf(&out, "\t\tcontexts: %s,\n", contexts)
}
if len(family.Signals) > 0 {
signals, err := goSignalSlice(family.Signals)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
fmt.Fprintf(&out, "\t\tsignals: %s,\n", signals)
}
if len(family.ValueMap) > 0 {
return nil, fmt.Errorf("family %q carries a value map, and the Go registry has no value-map reader yet", family.Current)
}
out.WriteString("\t},\n")
}
@@ -605,6 +722,29 @@ func renderGo(families []generatedFamily) ([]byte, error) {
return format.Source(out.Bytes())
}
func writeGoMember(out *bytes.Buffer, current string, member generatedMember) error {
parts := []string{fmt.Sprintf("name: %s", strconv.Quote(member.Name))}
if len(member.Contexts) > 0 {
contexts, err := goFieldContextSlice(member.Contexts)
if err != nil {
return fmt.Errorf("render family %q member %q: %w", current, member.Name, err)
}
parts = append(parts, "contexts: "+contexts)
}
if len(member.Signals) > 0 {
signals, err := goSignalSlice(member.Signals)
if err != nil {
return fmt.Errorf("render family %q member %q: %w", current, member.Name, err)
}
parts = append(parts, "signals: "+signals)
}
if len(member.ApplyToMetrics) > 0 {
parts = append(parts, "applyToMetrics: "+goStringSlice(member.ApplyToMetrics))
}
fmt.Fprintf(out, "\t\t\t{%s},\n", strings.Join(parts, ", "))
return nil
}
func goStringSlice(values []string) string {
if len(values) == 0 {
return "nil"
@@ -659,21 +799,36 @@ func goSignalSlice(values []string) (string, error) {
func renderTypeScript(families []generatedFamily) []byte {
var out bytes.Buffer
out.WriteString("// Code generated by scripts/semconv. DO NOT EDIT.\n\n")
out.WriteString("// An empty contexts/signals/applyToMetrics array places no constraint on\n")
out.WriteString("// that axis.\n")
out.WriteString("export type SemconvMember = {\n")
out.WriteString("\treadonly name: string;\n")
out.WriteString("\treadonly contexts: readonly string[];\n")
out.WriteString("\treadonly signals: readonly string[];\n")
out.WriteString("\treadonly applyToMetrics: readonly string[];\n};\n\n")
out.WriteString("export type SemconvFamily = {\n")
out.WriteString("\treadonly current: string;\n\treadonly old: readonly string[];\n")
out.WriteString("\treadonly current: string;\n")
out.WriteString("\treadonly kind: 'attribute' | 'metric';\n")
out.WriteString("\treadonly members: readonly SemconvMember[];\n")
out.WriteString("\treadonly contexts: readonly string[];\n\treadonly signals: readonly string[];\n")
out.WriteString("\treadonly applyToMetrics: readonly string[];\n")
out.WriteString("\treadonly valueMap: Readonly<Record<string, string>>;\n};\n\n")
out.WriteString("export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [\n")
for _, family := range families {
out.WriteString("\t{\n")
fmt.Fprintf(&out, "\t\tcurrent: %s,\n", tsString(family.Current))
fmt.Fprintf(&out, "\t\told: %s,\n", tsStringSlice(family.Old))
fmt.Fprintf(&out, "\t\tkind: %s,\n", tsString(family.Kind))
out.WriteString("\t\tmembers: [\n")
for _, member := range family.Members {
out.WriteString("\t\t\t{\n")
fmt.Fprintf(&out, "\t\t\t\tname: %s,\n", tsString(member.Name))
fmt.Fprintf(&out, "\t\t\t\tcontexts: %s,\n", tsStringSlice(member.Contexts))
fmt.Fprintf(&out, "\t\t\t\tsignals: %s,\n", tsStringSlice(member.Signals))
fmt.Fprintf(&out, "\t\t\t\tapplyToMetrics: %s,\n", tsStringSlice(member.ApplyToMetrics))
out.WriteString("\t\t\t},\n")
}
out.WriteString("\t\t],\n")
fmt.Fprintf(&out, "\t\tcontexts: %s,\n", tsStringSlice(family.Contexts))
fmt.Fprintf(&out, "\t\tsignals: %s,\n", tsStringSlice(family.Signals))
fmt.Fprintf(&out, "\t\tapplyToMetrics: %s,\n", tsStringSlice(family.ApplyToMetrics))
out.WriteString("\t\tvalueMap: {")
keys := sortedMapKeys(family.ValueMap)
for i, key := range keys {

View File

@@ -37,6 +37,11 @@ versions:
assert.ErrorContains(t, err, `schema version "latest"`, "malformed versions must not be silently reordered")
}
func TestParseSchemaVersionRejectsNonNumericComponent(t *testing.T) {
_, err := parseSchemaVersion("1.2.x")
assert.ErrorContains(t, err, `invalid numeric component "x"`, "non-numeric version components must fail generation")
}
func TestBuildFamiliesResolvesRenameChain(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
@@ -71,11 +76,13 @@ versions:
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "c",
Old: []string{"b", "x", "a"},
Kind: kindAttribute,
Contexts: []string{"attribute"},
Signals: []string{"traces"},
Current: "c",
Kind: kindAttribute,
Members: []generatedMember{
{Name: "b", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
{Name: "x", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
{Name: "a", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
},
}}, families, "predecessors should be ordered by distance and then name")
}
@@ -120,27 +127,239 @@ versions:
require.NoError(t, err)
assert.Equal(t, []generatedFamily{
{
Current: "all.current", Old: []string{"all.old"}, Kind: kindAttribute,
Contexts: nil, Signals: nil,
Current: "all.current", Kind: kindAttribute,
Members: []generatedMember{{Name: "all.old"}},
},
{
Current: "cpu.mode", Old: []string{"state"}, Kind: kindAttribute,
Contexts: []string{"attribute"}, Signals: []string{"metrics"},
Current: "cpu.mode", Kind: kindAttribute,
Members: []generatedMember{{
Name: "state", Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"system.cpu.time"},
}},
},
{
Current: "current.metric", Kind: kindMetric,
Members: []generatedMember{{Name: "old.metric", Contexts: []string{"metric"}, Signals: []string{"metrics"}}},
},
{
Current: "log.current", Kind: kindAttribute,
Members: []generatedMember{{Name: "log.old", Contexts: []string{"attribute"}, Signals: []string{"logs"}}},
},
{
Current: "resource.current", Kind: kindAttribute,
Members: []generatedMember{{Name: "resource.old", Contexts: []string{"resource"}}},
},
}, families, "schema sections should produce their documented per-member signal and context scopes")
}
func TestBuildFamiliesKeepsApplyToMetricsThroughNonMetricsRepeat(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
logs:
changes:
- rename_attributes:
attribute_map:
old: current
1.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
old: current
apply_to_metrics: [system.cpu.time]
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{
Name: "old", Contexts: []string{"attribute"}, Signals: []string{"logs", "metrics"},
ApplyToMetrics: []string{"system.cpu.time"},
}},
}}, families, "a repeat under a non-metrics section widens the signals and must not erase the metric scope")
}
func TestBuildFamiliesRejectsExplicitlyEmptyOverlayApplyToMetrics(t *testing.T) {
enabled := true
_, err := buildFamilies(nil, overlayFile{Families: map[string]overlayFamily{
"current": {Enabled: &enabled, Old: []string{"old"}, ApplyToMetrics: []string{}},
}})
assert.ErrorContains(t, err, "explicitly empty apply_to_metrics", "the overlay must not silently widen an empty list to every metric")
}
func TestBuildFamiliesRejectsExplicitlyEmptyApplyToMetrics(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
old: current
apply_to_metrics: []
`), &schema), "test schema must decode")
_, err := buildFamilies([]schemaFile{schema}, overlayFile{})
assert.ErrorContains(t, err, "explicitly empty apply_to_metrics", "an empty list must not silently widen to every metric")
}
func TestBuildFamiliesKeepsFanOutSeparate(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
state: cpu.mode
apply_to_metrics: [system.cpu.time]
1.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
state: db.client.connection.state
apply_to_metrics: [db.client.connections.usage]
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{
{
Current: "cpu.mode", Kind: kindAttribute,
Members: []generatedMember{{
Name: "state", Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"system.cpu.time"},
}},
},
{
Current: "current.metric", Old: []string{"old.metric"}, Kind: kindMetric,
Contexts: []string{"metric"}, Signals: []string{"metrics"},
Current: "db.client.connection.state", Kind: kindAttribute,
Members: []generatedMember{{
Name: "state", Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"db.client.connections.usage"},
}},
},
{
Current: "log.current", Old: []string{"log.old"}, Kind: kindAttribute,
Contexts: []string{"attribute"}, Signals: []string{"logs"},
}, families, "an old name with differently scoped rename targets must keep one membership per target")
}
func TestBuildFamiliesKeepsUnscopedRenameUnscoped(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
direction: network.io.direction
apply_to_metrics: [system.disk.io, system.disk.merged]
1.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
system.network.io.direction: network.io.direction
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "network.io.direction", Kind: kindAttribute,
Members: []generatedMember{
{Name: "direction", Contexts: []string{"attribute"}, Signals: []string{"metrics"}, ApplyToMetrics: []string{"system.disk.io", "system.disk.merged"}},
{Name: "system.network.io.direction", Contexts: []string{"attribute"}, Signals: []string{"metrics"}},
},
{
Current: "resource.current", Old: []string{"resource.old"}, Kind: kindAttribute,
Contexts: []string{"resource"},
}}, families, "a rename without apply_to_metrics stays unbounded; a scoped sibling must not bound it")
}
func TestBuildFamiliesKeepsMemberScopesThroughChains(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
messaging.client_id: messaging.client.id
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
messaging.kafka.client_id: messaging.client_id
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "messaging.client.id", Kind: kindAttribute,
Members: []generatedMember{
{Name: "messaging.client_id", Contexts: []string{"attribute"}, Signals: []string{"metrics"}},
{Name: "messaging.kafka.client_id", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
},
}, families, "schema sections should produce their documented signal and context scopes")
}}, families, "a member keeps the scope of its own rename edge; later hops only carry it to the root")
}
func TestBuildFamiliesRecordsCrossContextMembersSeparately(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
http.user_agent: user_agent.original
resources:
changes:
- rename_attributes:
attribute_map:
browser.user_agent: user_agent.original
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "user_agent.original", Kind: kindAttribute,
Members: []generatedMember{
{Name: "browser.user_agent", Contexts: []string{"resource"}},
{Name: "http.user_agent", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
},
}}, families, "members renamed from different contexts must keep their own context scopes")
}
func TestBuildFamiliesMergesRepeatedEdgeScopes(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
logs:
changes:
- rename_attributes:
attribute_map:
old: current
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
old: current
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{
{Name: "old", Contexts: []string{"attribute"}, Signals: []string{"logs", "traces"}},
},
}}, families, "the same rename filed under several sections widens the member scope")
}
func TestOverlayAddsFamilyWithoutSchemaHistory(t *testing.T) {
@@ -157,8 +376,8 @@ func TestOverlayAddsFamilyWithoutSchemaHistory(t *testing.T) {
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "added.current",
Old: []string{"added.old"},
Kind: kindAttribute,
Members: []generatedMember{{Name: "added.old"}},
Contexts: []string{"resource"},
Signals: []string{"traces"},
}}, families, "an explicit overlay family should not require schema history")
@@ -179,26 +398,76 @@ versions:
enabled := true
families, err := buildFamilies([]schemaFile{schema}, overlayFile{Families: map[string]overlayFamily{
"current": {
Enabled: &enabled,
AddOld: []string{"older"},
ExcludeOld: []string{"old"},
AddContexts: []string{"resource"},
AddSignals: []string{"logs"},
ValueMap: map[string]string{"legacy": "current"},
Enabled: &enabled,
AddOld: []string{"older"},
ExcludeOld: []string{"old"},
AddSignals: []string{"logs"},
ValueMap: map[string]string{"legacy": "current"},
},
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current",
Old: []string{"older"},
Kind: kindAttribute,
Contexts: []string{"attribute", "resource"},
Signals: []string{"logs", "traces"},
Members: []generatedMember{{Name: "older"}},
ValueMap: map[string]string{"legacy": "current"},
}}, families, "overlay additions and exclusions should be applied to the generated family")
}
func TestOverlayAddSignalsWidensMemberScopes(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
old: current
`), &schema), "test schema must decode")
enabled := true
families, err := buildFamilies([]schemaFile{schema}, overlayFile{Families: map[string]overlayFamily{
"current": {Enabled: &enabled, AddSignals: []string{"logs"}},
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current",
Kind: kindAttribute,
Members: []generatedMember{
{Name: "old", Contexts: []string{"attribute"}, Signals: []string{"logs", "traces"}},
},
}}, families, "add_signals widens the schema-derived member scopes and never narrows the family gate")
}
func TestOverlaySignalsSetTheFamilyGate(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
all:
changes:
- rename_attributes:
attribute_map:
old: current
`), &schema), "test schema must decode")
enabled := true
families, err := buildFamilies([]schemaFile{schema}, overlayFile{Families: map[string]overlayFamily{
"current": {Enabled: &enabled, Signals: []string{"logs", "traces"}},
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current",
Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
Signals: []string{"logs", "traces"},
}}, families, "the overlay signals list gates the family without touching member scopes")
}
func TestOverlayDisablesFamilyWhenDefaultIsEnabled(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
@@ -223,23 +492,21 @@ versions:
assert.Empty(t, families, "an explicitly disabled family must override default_enabled")
}
func TestRenderGoIsDeterministic(t *testing.T) {
func TestRenderGoRejectsValueMap(t *testing.T) {
families := []generatedFamily{{
Current: "current", Old: []string{"old"}, Kind: kindAttribute,
ValueMap: map[string]string{"b": "2", "a": "1"},
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
ValueMap: map[string]string{"a": "1"},
}}
first, err := renderGo(families)
require.NoError(t, err)
second, err := renderGo(families)
require.NoError(t, err)
assert.Equal(t, first, second, "Go generation must not depend on map iteration order")
_, err := renderGo(families)
assert.ErrorContains(t, err, "no value-map reader yet", "a value map must fail Go generation until a reader exists")
}
func TestRenderGoUsesCanonicalTelemetryTypes(t *testing.T) {
families := []generatedFamily{{
Current: "current", Old: []string{"old"}, Kind: kindAttribute,
Contexts: []string{"resource"}, Signals: []string{"traces"},
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old", Contexts: []string{"resource"}, Signals: []string{"traces"}}},
}}
output, err := renderGo(families)
@@ -248,15 +515,112 @@ func TestRenderGoUsesCanonicalTelemetryTypes(t *testing.T) {
assert.Contains(t, string(output), "telemetrytypes.SignalTraces", "generated signals should use telemetrytypes")
}
func TestRenderGoRejectsUnknownContext(t *testing.T) {
families := []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old", Contexts: []string{"bogus"}}},
}}
_, err := renderGo(families)
assert.ErrorContains(t, err, `unsupported field context "bogus"`, "a bad overlay context must fail generation, not compile")
}
func TestRenderGoRejectsUnknownSignal(t *testing.T) {
families := []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
Signals: []string{"bogus"},
}}
_, err := renderGo(families)
assert.ErrorContains(t, err, `unsupported signal "bogus"`, "a bad overlay signal must fail generation, not compile")
}
func TestRenderGoPinsOutput(t *testing.T) {
families := []generatedFamily{{
Current: "deployment.environment.name", Kind: kindAttribute,
Members: []generatedMember{{Name: "deployment.environment"}},
Signals: []string{"logs", "traces"},
}}
output, err := renderGo(families)
require.NoError(t, err)
assert.Equal(t, `// Code generated by scripts/semconv. DO NOT EDIT.
package semconv
import "github.com/SigNoz/signoz/pkg/types/telemetrytypes"
var families = []Family{
{
current: "deployment.environment.name",
kind: KindAttribute,
members: []Member{
{name: "deployment.environment"},
},
signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalTraces},
},
}
`, string(output), "the emitted Go text is a contract; regeneration must be reviewable")
}
func TestRenderTypeScriptIsDeterministic(t *testing.T) {
families := []generatedFamily{{
Current: "current", Old: []string{"old"}, Kind: kindAttribute,
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
ValueMap: map[string]string{"b": "2", "a": "1"},
}}
assert.Equal(t, renderTypeScript(families), renderTypeScript(families), "TypeScript generation must not depend on map iteration order")
}
func TestRenderTypeScriptPinsOutput(t *testing.T) {
families := []generatedFamily{{
Current: "deployment.environment.name", Kind: kindAttribute,
Members: []generatedMember{{Name: "deployment.environment"}},
Signals: []string{"logs", "traces"},
}}
assert.Equal(t, `// Code generated by scripts/semconv. DO NOT EDIT.
// An empty contexts/signals/applyToMetrics array places no constraint on
// that axis.
export type SemconvMember = {
readonly name: string;
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly applyToMetrics: readonly string[];
};
export type SemconvFamily = {
readonly current: string;
readonly kind: 'attribute' | 'metric';
readonly members: readonly SemconvMember[];
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly valueMap: Readonly<Record<string, string>>;
};
export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
{
current: 'deployment.environment.name',
kind: 'attribute',
members: [
{
name: 'deployment.environment',
contexts: [],
signals: [],
applyToMetrics: [],
},
],
contexts: [],
signals: ['logs', 'traces'],
valueMap: {},
},
] as const;
`, string(renderTypeScript(families)), "the emitted TypeScript text is a contract; regeneration must be reviewable")
}
func TestBuildFamiliesHandlesRenameRollback(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
@@ -280,11 +644,9 @@ versions:
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "original",
Old: []string{"temporary"},
Kind: kindMetric,
Contexts: []string{"metric"},
Signals: []string{"metrics"},
Current: "original",
Kind: kindMetric,
Members: []generatedMember{{Name: "temporary", Contexts: []string{"metric"}, Signals: []string{"metrics"}}},
}}, families, "the latest rollback destination should remain the family root")
}
@@ -356,11 +718,9 @@ versions:
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "shared.current",
Old: []string{"attribute.old"},
Kind: kindAttribute,
Contexts: []string{"attribute"},
Signals: []string{"traces"},
Current: "shared.current",
Kind: kindAttribute,
Members: []generatedMember{{Name: "attribute.old", Contexts: []string{"attribute"}, Signals: []string{"traces"}}},
}}, families, "a kind-less overlay policy should affect only the attribute family")
}

View File

@@ -2,10 +2,35 @@
#
# Families are keyed by their current OpenTelemetry name. Schema-derived
# families are disabled by default so rollout remains explicit and reversible.
# The signals list is the per-family rollout gate.
default_enabled: false
families:
deployment.environment.name:
enabled: true
signals: [traces, logs, metrics]
# The family names an attribute of the resource or the span/log, never a
# field inside a log body: without this gate a body-context key with the
# same path joins the family and skips the body-JSON machinery.
contexts: [resource, attribute]
db.system.name:
# The db.system value domain also renamed (mssql -> microsoft.sql_server
# and others), and no value mapping is read yet. A name-only merge matches
# half the history on every signal, so the family stays off until a
# value-mapping reader exists.
enabled: false
# These metric renames predate the vendored schema history, so the overlay
# declares the old names itself.
k8s.pod.cpu.usage:
enabled: true
kind: metric
old: [k8s.pod.cpu.utilization]
k8s.node.cpu.usage:
enabled: true
kind: metric
old: [k8s.node.cpu.utilization]
container.cpu.usage:
enabled: true
kind: metric
old: [container.cpu.utilization]

View File

@@ -4,19 +4,25 @@ from datetime import UTC, datetime, timedelta
import pytest
from fixtures.logs import Logs
from fixtures.metrics import Metrics
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
PREFIX = "semconv-fam"
CURRENT_KEY = "deployment.environment.name"
OLD_KEY = "deployment.environment"
# Row identities. The span name, the log body, and service.name are the identity.
# Tests compare identity sets filtered by PREFIX, so reruns on a reused stack
# with leftover rows stay stable.
OLD = f"{PREFIX}-old" # only the old spelling, value "production"
NEW = f"{PREFIX}-new" # only the current spelling, value "production"
# Tests compare identity sets filtered by PREFIX.
OLD = f"{PREFIX}-old"
NEW = f"{PREFIX}-new"
BOTH = f"{PREFIX}-both" # current "staging" and old "production" - the conflict row
NEITHER = f"{PREFIX}-neither" # no member at all
NEITHER = f"{PREFIX}-neither"
LABEL_METRIC = "semconv.fam.label.metric"
# A span-metrics metric, so the resource_ label layout of the span-metrics
# processor applies.
SPAN_METRIC = "signoz_calls_total"
OLD_NAME_METRIC = "k8s.pod.cpu.utilization"
CURRENT_NAME_METRIC = "k8s.pod.cpu.usage"
_ROWS = [
(OLD, {OLD_KEY: "production"}, timedelta(seconds=4)),
@@ -31,9 +37,10 @@ def family_fleet(
insert_logs: Callable[[list[Logs]], None],
insert_traces: Callable[[list[Traces]], None],
) -> Generator[datetime]:
"""Yields the base timestamp of the inserted rows."""
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=1)
"""Inserts one span and one log per identity and yields the base
timestamp. The base aligns to the minute, so no row offset crosses a 60s
time-series bucket boundary."""
now = datetime.now(tz=UTC).replace(second=0, microsecond=0) - timedelta(minutes=1)
insert_traces(
[
Traces(
@@ -62,3 +69,30 @@ def family_fleet(
]
)
yield now
@pytest.fixture(name="metric_family_fleet", scope="function")
def metric_family_fleet(insert_metrics: Callable[[list[Metrics]], None]) -> Generator[datetime]:
"""Inserts the label and metric-name family series and yields the query
end timestamp. Every series has a power-of-two value, so a missed member
is a unique wrong sum."""
now = datetime.now(tz=UTC)
# The querier clamps very recent metric samples (flux interval), so the
# fleet sits in the past.
seeded = now - timedelta(minutes=10)
gauge = {"temporality": "Unspecified", "type_": "Gauge", "is_monotonic": False}
insert_metrics(
[
Metrics(metric_name=LABEL_METRIC, labels={"deployment.environment.name": "staging"}, timestamp=seeded, value=1.0, **gauge),
Metrics(metric_name=LABEL_METRIC, labels={"deployment.environment": "production"}, timestamp=seeded, value=2.0, **gauge),
Metrics(metric_name=SPAN_METRIC, labels={"resource_deployment.environment.name": "staging", "pod": "span-current"}, timestamp=seeded, value=4.0, **gauge),
Metrics(metric_name=LABEL_METRIC, labels={"region": "keyless"}, timestamp=seeded, value=8.0, **gauge),
Metrics(metric_name=OLD_NAME_METRIC, labels={"pod": "a"}, timestamp=seeded, value=16.0, **gauge),
Metrics(metric_name=CURRENT_NAME_METRIC, labels={"pod": "b"}, timestamp=seeded, value=32.0, **gauge),
Metrics(metric_name=LABEL_METRIC, labels={"deployment.environment.name": "staging", "deployment.environment": "production"}, timestamp=seeded, value=64.0, **gauge),
Metrics(metric_name=SPAN_METRIC, labels={"resource_deployment.environment": "production", "pod": "span-old"}, timestamp=seeded, value=128.0, **gauge),
Metrics(metric_name=SPAN_METRIC, labels={"deployment.environment": "production", "pod": "span-plain"}, timestamp=seeded, value=256.0, **gauge),
Metrics(metric_name=LABEL_METRIC, labels={"deployment.environment.name": "production"}, timestamp=seeded, value=512.0, **gauge),
]
)
yield now

View File

@@ -17,7 +17,7 @@ pkg/types/querybuildertypes/querybuildertypesv5/builder_elements.go):
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.
semconvfamilies/01_family_matrix.py.
Seed data lives in fixtures/queriercommon.py: GOLD and SILVER carry the
keys, NONE carries none. Every case asserts which identities a filter

View File

@@ -6,14 +6,19 @@ import pytest
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.metadata import AttributesMetadata, get_field_keys, get_field_values
from fixtures.querier import (
RequestType,
build_aggregation,
build_group_by_field,
build_order_by,
build_raw_query,
build_scalar_query,
build_traces_scalar_query,
get_all_warnings,
get_column_data_from_response,
get_scalar_columns,
get_scalar_table_data,
make_query_request,
)
from fixtures.semconvfamilies import (
@@ -33,6 +38,11 @@ FILTER_MATRIX = [
pytest.param("{key} IN ['production', 'staging']", {OLD, NEW, BOTH}, id="in_matches_merged_value"),
pytest.param("{key} NOT IN ['production']", {BOTH, NEITHER}, id="not_in_keeps_keyless"),
pytest.param("{key} LIKE '%prod%'", {OLD, NEW}, id="like_matches_merged_value"),
pytest.param("{key} NOT LIKE '%prod%'", {BOTH, NEITHER}, id="not_like_keeps_keyless"),
pytest.param("{key} ILIKE 'PROD%'", {OLD, NEW}, id="ilike_matches_merged_value"),
pytest.param("{key} CONTAINS 'oduct'", {OLD, NEW}, id="contains_matches_merged_value"),
pytest.param("{key} REGEXP '^prod.*'", {OLD, NEW}, id="regexp_matches_merged_value"),
pytest.param("{key} NOT CONTAINS 'prod'", {BOTH, NEITHER}, id="not_contains_keeps_keyless"),
pytest.param("{key} EXISTS", {OLD, NEW, BOTH}, id="exists_is_any_member"),
pytest.param("{key} NOT EXISTS", {NEITHER}, id="not_exists_is_no_member"),
pytest.param("{key} != 'production' AND {key} EXISTS", {BOTH}, id="neq_composed_with_exists"),
@@ -82,6 +92,43 @@ def test_family_filters(
assert matched == expected, expression
@pytest.mark.parametrize("expression_template,expected", FILTER_MATRIX)
@pytest.mark.parametrize("requested_key", [CURRENT_KEY, OLD_KEY], ids=["current", "old"])
@pytest.mark.parametrize("context", ["resource", "attribute"])
def test_log_family_filters(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
context: str,
requested_key: str,
expression_template: str,
expected: set[str],
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
expression = expression_template.format(key=f"{context}.{requested_key}")
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
"logs",
limit=100,
filter_expression=expression,
order=[build_order_by("timestamp", "asc")],
select_fields=[{"name": "body"}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
matched = {name for name in get_column_data_from_response(response.json(), "body") if name.startswith(PREFIX)}
assert matched == expected, expression
@pytest.mark.parametrize("expression_template,expected", LITERAL_MATRIX)
def test_flag_off_stays_literal(
signoz_families_off: types.SigNoz,
@@ -115,40 +162,6 @@ def test_flag_off_stays_literal(
assert matched == expected, expression
@pytest.mark.parametrize("expression_template,expected", LITERAL_MATRIX)
def test_logs_stay_literal_with_flag_on(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
expression_template: str,
expected: set[str],
) -> None:
"""Only traces have family support today."""
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
expression = expression_template.format(key=f"resource.{CURRENT_KEY}")
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
"logs",
limit=100,
filter_expression=expression,
order=[build_order_by("timestamp", "asc")],
select_fields=[{"name": "body"}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
matched = {body for body in get_column_data_from_response(response.json(), "body") if body.startswith(PREFIX)}
assert matched == expected, expression
@pytest.mark.parametrize("requested_key", [CURRENT_KEY, OLD_KEY], ids=["current", "old"])
def test_group_by_merges_and_echoes_requested_spelling(
signoz: types.SigNoz,
@@ -166,7 +179,7 @@ def test_group_by_merges_and_echoes_requested_spelling(
request_type=RequestType.SCALAR,
queries=[
build_traces_scalar_query(
[build_aggregation("count()")],
[build_aggregation("count_distinct(name)")],
filter_expression=f"service.name LIKE '{PREFIX}%'",
group_by=[build_group_by_field(requested_key, "string", "resource")],
)
@@ -174,14 +187,179 @@ def test_group_by_merges_and_echoes_requested_spelling(
)
assert response.status_code == HTTPStatus.OK, response.text
result = response.json()["data"]["data"]["results"][0]
group_column = result["columns"][0]
group_column = get_scalar_columns(response.json())[0]
assert group_column["name"] == requested_key, group_column
assert group_column["columnType"] == "group", group_column
groups = {row[0] for row in result["data"]}
assert {"production", "staging"}.issubset(groups), groups
assert None in groups, groups
# Distinct identities per group make the counts rerun-safe on a reused
# stack: OLD and NEW merge into production, BOTH is staging, NEITHER has
# no spelling at all.
groups = {tuple(row) for row in get_scalar_table_data(response.json())}
assert groups == {("production", 2), ("staging", 1), (None, 1)}, groups
@pytest.mark.parametrize("requested_key", [CURRENT_KEY, OLD_KEY], ids=["current", "old"])
def test_log_group_by_merges_and_echoes_requested_spelling(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
requested_key: str,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="logs",
aggregations=[build_aggregation("count_distinct(body)")],
filter_expression=f"service.name LIKE '{PREFIX}%'",
group_by=[build_group_by_field(requested_key, "string", "resource")],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
group_column = get_scalar_columns(response.json())[0]
assert group_column["name"] == requested_key, group_column
groups = {tuple(row) for row in get_scalar_table_data(response.json())}
assert groups == {("production", 2), ("staging", 1), (None, 1)}, groups
def test_time_series_group_by_merges_family(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.TIME_SERIES,
queries=[
build_traces_scalar_query(
[build_aggregation("count_distinct(name)")],
filter_expression=f"service.name LIKE '{PREFIX}%'",
group_by=[build_group_by_field(CURRENT_KEY, "string", "resource")],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
series = response.json()["data"]["data"]["results"][0]["aggregations"][0]["series"]
values_by_group = {}
for entry in series:
for label in entry.get("labels", []):
if label["key"]["name"] == CURRENT_KEY:
values_by_group[label["value"]] = max((point["value"] for point in entry["values"]), default=None)
assert values_by_group.get("production") == 2, values_by_group
assert values_by_group.get("staging") == 1, values_by_group
def test_raw_select_reads_merged_value(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
"traces",
limit=100,
filter_expression=f"service.name LIKE '{PREFIX}%'",
order=[build_order_by("timestamp", "asc")],
select_fields=[{"name": "span.name"}, {"name": f"resource.{CURRENT_KEY}"}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = response.json()["data"]["data"]["results"][0]["rows"]
value_by_identity = {}
for row in rows:
data = row["data"]
if data.get("name", "").startswith(PREFIX):
value_by_identity[data["name"]] = data.get(CURRENT_KEY)
assert value_by_identity.get(OLD) == "production", value_by_identity
assert value_by_identity.get(NEW) == "production", value_by_identity
assert value_by_identity.get(BOTH) == "staging", value_by_identity
assert value_by_identity.get(NEITHER) in ("", None), value_by_identity
def test_order_by_family_key_sorts_merged_values(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
"traces",
limit=100,
filter_expression=f"{CURRENT_KEY} EXISTS AND service.name LIKE '{PREFIX}%'",
order=[build_order_by(f"resource.{CURRENT_KEY}", "asc")],
select_fields=[{"name": "span.name"}, {"name": f"resource.{CURRENT_KEY}"}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = response.json()["data"]["data"]["results"][0]["rows"]
merged_values = [row["data"].get(CURRENT_KEY) for row in rows if row["data"].get("name", "").startswith(PREFIX)]
assert merged_values == sorted(merged_values), merged_values
assert set(merged_values) == {"production", "staging"}, merged_values
def test_qualified_family_key_emits_no_ambiguity_warning(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((family_fleet - timedelta(minutes=2)).timestamp() * 1000),
end_ms=int((family_fleet + timedelta(minutes=1)).timestamp() * 1000),
request_type=RequestType.RAW,
queries=[
build_raw_query(
"A",
"traces",
limit=10,
filter_expression=f"resource.{CURRENT_KEY} = 'production'",
order=[build_order_by("timestamp", "asc")],
select_fields=[{"name": "span.name"}],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
assert get_all_warnings(response.json()) == [], response.json()["data"].get("warning")
def test_bare_name_prefers_resource_and_warns(
@@ -215,6 +393,75 @@ def test_bare_name_prefers_resource_and_warns(
matched = {name for name in get_column_data_from_response(response.json(), "name") if name.startswith(PREFIX)}
assert matched == {OLD, NEW}
warning = response.json()["data"].get("warning") or {}
messages = " ".join(entry.get("message", "") for entry in warning.get("warnings", []))
messages = " ".join(entry.get("message", "") for entry in get_all_warnings(response.json()))
assert "ambiguous" in messages.lower(), messages
def test_field_keys_stay_literal(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = get_field_keys(signoz, token, {"signal": "traces", "searchText": "deployment.environment"})
assert response.status_code == HTTPStatus.OK, response.text
names = set(response.json()["data"]["keys"].keys())
assert CURRENT_KEY in names, names
assert OLD_KEY in names, names
def test_field_values_union_the_family(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = get_field_values(
signoz,
token,
{"signal": "traces", "name": OLD_KEY, "fieldContext": "resource"},
)
assert response.status_code == HTTPStatus.OK, response.text
values = set(response.json()["data"]["values"].get("stringValues", []))
# staging exists only under the current spelling on the BOTH row, so only
# the family union makes it reachable from a query on the old spelling.
assert {"production", "staging"}.issubset(values), values
@pytest.mark.parametrize("requested_key", [CURRENT_KEY, OLD_KEY], ids=["current", "old"])
@pytest.mark.parametrize("context", ["resource", None], ids=["resource", "unspecified"])
def test_related_values_search_matches_merged_value(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
family_fleet: datetime,
insert_attributes_metadata: Callable[[list[AttributesMetadata]], None],
requested_key: str,
context: str | None,
) -> None:
insert_attributes_metadata(
[
AttributesMetadata(data_source="traces", resource_attributes={"service.name": OLD, OLD_KEY: "production"}, timestamp=family_fleet),
AttributesMetadata(data_source="traces", resource_attributes={"service.name": NEW, CURRENT_KEY: "production"}, timestamp=family_fleet),
AttributesMetadata(data_source="traces", resource_attributes={"service.name": BOTH, CURRENT_KEY: "staging", OLD_KEY: "production"}, timestamp=family_fleet),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
params = {
"signal": "traces",
"name": requested_key,
"existingQuery": f"service.name IN ['{OLD}', '{NEW}', '{BOTH}']",
"searchText": "prod",
"startUnixMilli": int((family_fleet - timedelta(hours=1)).timestamp() * 1000),
"endUnixMilli": int((family_fleet + timedelta(hours=1)).timestamp() * 1000),
}
if context is not None:
params["fieldContext"] = context
response = get_field_values(signoz, token, params)
assert response.status_code == HTTPStatus.OK, response.text
# The BOTH row reads as staging, so its old production value must not
# satisfy the search.
assert set(response.json()["data"]["values"].get("relatedValues") or []) == {"production"}

View File

@@ -0,0 +1,235 @@
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
import pytest
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.metrics import Metrics
from fixtures.querier import (
RequestType,
build_group_by_field,
build_metrics_aggregation,
build_scalar_query,
get_scalar_columns,
get_scalar_table_data,
make_query_request,
)
from fixtures.semconvfamilies import (
CURRENT_NAME_METRIC,
LABEL_METRIC,
OLD_NAME_METRIC,
SPAN_METRIC,
)
@pytest.mark.parametrize("requested_key", ["deployment.environment.name", "deployment.environment"], ids=["current", "old"])
def test_metric_label_filter_merges_stored_spellings(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
metric_family_fleet: datetime,
requested_key: str,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((metric_family_fleet - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(metric_family_fleet.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(LABEL_METRIC, "latest", "sum", "unspecified", reduce_to="last")],
filter_expression=f"{requested_key} = 'production'",
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = get_scalar_table_data(response.json())
# The old spelling series (2) and the current spelling series (512). The
# conflict series (64) merges current-first to staging and stays out.
assert rows and rows[0][-1] == 514.0, rows
def test_metric_label_group_by_merges_stored_spellings(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
metric_family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((metric_family_fleet - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(metric_family_fleet.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(LABEL_METRIC, "latest", "sum", "unspecified", reduce_to="last")],
group_by=[build_group_by_field("deployment.environment.name", "string", "attribute")],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
group_column = get_scalar_columns(response.json())[0]
assert group_column["name"] == "deployment.environment.name", group_column
groups = {row[0]: row[-1] for row in get_scalar_table_data(response.json())}
assert groups.get("production") == 514.0, groups
# The conflict series lands in the staging group: the current spelling
# wins the merge.
assert groups.get("staging") == 65.0, groups
assert groups.get("") == 8.0, groups
@pytest.mark.parametrize(
"expression,expected",
[
pytest.param("deployment.environment = 'production'", 384.0, id="old_spelling_reads_plain_and_resource_labels"),
pytest.param("deployment.environment.name = 'production'", 384.0, id="current_spelling_reads_the_same_series"),
pytest.param("deployment.environment.name = 'staging'", 4.0, id="resource_current_spelling"),
],
)
def test_span_metric_label_filter_reads_resource_layout(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
metric_family_fleet: datetime,
expression: str,
expected: float,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((metric_family_fleet - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(metric_family_fleet.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(SPAN_METRIC, "latest", "sum", "unspecified", reduce_to="last")],
filter_expression=expression,
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = get_scalar_table_data(response.json())
assert rows and rows[0][-1] == expected, rows
@pytest.mark.parametrize("requested", [OLD_NAME_METRIC, CURRENT_NAME_METRIC], ids=["old", "current"])
def test_metric_name_family_unions_storage_names(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
metric_family_fleet: datetime,
requested: str,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((metric_family_fleet - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(metric_family_fleet.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(requested, "latest", "sum", "unspecified", reduce_to="last")],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = get_scalar_table_data(response.json())
assert rows and rows[0][-1] == 48.0, rows
def test_metric_name_union_double_counts_dual_emission(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_metrics: Callable[[list[Metrics]], None],
) -> None:
now = datetime.now(tz=UTC)
seeded = now - timedelta(minutes=10)
gauge = {"temporality": "Unspecified", "type_": "Gauge", "is_monotonic": False}
insert_metrics(
[
Metrics(metric_name=OLD_NAME_METRIC, labels={"pod": "overlap"}, timestamp=seeded, value=16.0, **gauge),
Metrics(metric_name=CURRENT_NAME_METRIC, labels={"pod": "overlap"}, timestamp=seeded, value=16.0, **gauge),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz,
token,
start_ms=int((now - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(now.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(CURRENT_NAME_METRIC, "latest", "sum", "unspecified", reduce_to="last")],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = get_scalar_table_data(response.json())
assert rows and rows[0][-1] == 32.0, rows
def test_metric_family_stays_literal_with_flag_off(
signoz_families_off: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
metric_family_fleet: datetime,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = make_query_request(
signoz_families_off,
token,
start_ms=int((metric_family_fleet - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(metric_family_fleet.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(LABEL_METRIC, "latest", "sum", "unspecified", reduce_to="last")],
filter_expression="deployment.environment.name = 'production'",
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = get_scalar_table_data(response.json())
# Only the series stored under the requested spelling.
assert rows and rows[0][-1] == 512.0, rows
response = make_query_request(
signoz_families_off,
token,
start_ms=int((metric_family_fleet - timedelta(minutes=30)).timestamp() * 1000),
end_ms=int(metric_family_fleet.timestamp() * 1000),
request_type=RequestType.SCALAR,
queries=[
build_scalar_query(
name="A",
signal="metrics",
aggregations=[build_metrics_aggregation(OLD_NAME_METRIC, "latest", "sum", "unspecified", reduce_to="last")],
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
rows = get_scalar_table_data(response.json())
assert rows and rows[0][-1] == 16.0, rows

View File

@@ -51,5 +51,7 @@ def signoz_families_off(
request=request,
pytestconfig=pytestconfig,
cache_key="signoz-semconv-families-off",
env_overrides={},
env_overrides={
"SIGNOZ_FLAGGER_CONFIG_BOOLEAN_RESOLVE__SEMCONV__FAMILIES": False,
},
)