Compare commits

..

1 Commits

Author SHA1 Message Date
srikanthccv
3e30ef6a71 proto: semconv families on main via LogicalField resolution output
Resolution returns []*LogicalField: the slice expresses ambiguity (union
across, one condition each), the group expresses a semantic-convention
family (merge within, members current-first). Members alias metadata map
entries; nothing is copied, mutated, or annotated.

- telemetrytypes.LogicalField: requested spelling + shared identity
  triple + current-first members
- qbtypes.FieldMapper gains one method: ExistsFor, the per-key existence
  primitive every mapper already had internally
- querybuilder.LogicalValueExpr/LogicalExistsExpr: family composition
  written once, built only from FieldFor/ExistsFor — no signal
  implements family logic
- MatchingLogicalFields/ResolveLogicalFields: grouping by identity,
  rank-sorted members (precedence structural, not arrival order),
  family+collision stacking with resource preference
- ExpandKeySelectorsForFamilies: sibling prefetch at the statement
  builder layer; the metadata store stays family-blind and autocomplete
  stays literal
- traces + resourcefilter condition builders compile per logical field
  (COALESCE with '' tail preserving keyless-row semantics; presence =
  any member; fingerprint index hints widened to any member)
- traces ColumnExpressionFor upgrades legacy candidates to families
  post-hoc; candidate order and non-family behavior unchanged
- logs/metrics/audit/metadata/rulestatehistory flatten via SingleKeys —
  single-member by construction, SQL unchanged

Family scope: traces only, deployment.environment(.name) and
db.system.name enabled. All existing tests pass unchanged except three
test doubles that replayed the old resolution helpers.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-10 13:02:03 +05:30
63 changed files with 1428 additions and 1964 deletions

View File

@@ -12198,9 +12198,9 @@ paths:
- dashboard
/api/v1/resetPassword:
post:
deprecated: true
deprecated: false
description: This endpoint resets the password by token
operationId: ResetPasswordDeprecated
operationId: ResetPassword
requestBody:
content:
application/json:
@@ -15567,41 +15567,6 @@ paths:
summary: Forgot password
tags:
- users
/api/v2/factor_password/reset:
post:
deprecated: false
description: This endpoint resets the password using a single use reset password
token
operationId: ResetPassword
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/TypesPostableResetPassword'
responses:
"204":
description: No Content
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
summary: Reset password
tags:
- users
/api/v2/features:
get:
deprecated: false

View File

@@ -255,10 +255,9 @@ export const useCreateInvite = <
};
/**
* This endpoint resets the password by token
* @deprecated
* @summary Reset password
*/
export const resetPasswordDeprecated = (
export const resetPassword = (
typesPostableResetPasswordDTO?: BodyType<TypesPostableResetPasswordDTO>,
signal?: AbortSignal,
) => {
@@ -271,23 +270,23 @@ export const resetPasswordDeprecated = (
});
};
export const getResetPasswordDeprecatedMutationOptions = <
export const getResetPasswordMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
> => {
const mutationKey = ['resetPasswordDeprecated'];
const mutationKey = ['resetPassword'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
@@ -297,47 +296,45 @@ export const getResetPasswordDeprecatedMutationOptions = <
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
Awaited<ReturnType<typeof resetPassword>>,
{ data?: BodyType<TypesPostableResetPasswordDTO> }
> = (props) => {
const { data } = props ?? {};
return resetPasswordDeprecated(data);
return resetPassword(data);
};
return { mutationFn, ...mutationOptions };
};
export type ResetPasswordDeprecatedMutationResult = NonNullable<
Awaited<ReturnType<typeof resetPasswordDeprecated>>
export type ResetPasswordMutationResult = NonNullable<
Awaited<ReturnType<typeof resetPassword>>
>;
export type ResetPasswordDeprecatedMutationBody =
export type ResetPasswordMutationBody =
| BodyType<TypesPostableResetPasswordDTO>
| undefined;
export type ResetPasswordDeprecatedMutationError =
ErrorType<RenderErrorResponseDTO>;
export type ResetPasswordMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @deprecated
* @summary Reset password
*/
export const useResetPasswordDeprecated = <
export const useResetPassword = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
> => {
return useMutation(getResetPasswordDeprecatedMutationOptions(options));
return useMutation(getResetPasswordMutationOptions(options));
};
/**
* This endpoint lists all users
@@ -596,89 +593,6 @@ export const useForgotPassword = <
> => {
return useMutation(getForgotPasswordMutationOptions(options));
};
/**
* This endpoint resets the password using a single use reset password token
* @summary Reset password
*/
export const resetPassword = (
typesPostableResetPasswordDTO?: BodyType<TypesPostableResetPasswordDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v2/factor_password/reset`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: typesPostableResetPasswordDTO,
signal,
});
};
export const getResetPasswordMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
> => {
const mutationKey = ['resetPassword'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
options.mutation.mutationKey
? options
: { ...options, mutation: { ...options.mutation, mutationKey } }
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof resetPassword>>,
{ data?: BodyType<TypesPostableResetPasswordDTO> }
> = (props) => {
const { data } = props ?? {};
return resetPassword(data);
};
return { mutationFn, ...mutationOptions };
};
export type ResetPasswordMutationResult = NonNullable<
Awaited<ReturnType<typeof resetPassword>>
>;
export type ResetPasswordMutationBody =
| BodyType<TypesPostableResetPasswordDTO>
| undefined;
export type ResetPasswordMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Reset password
*/
export const useResetPassword = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof resetPassword>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
> => {
return useMutation(getResetPasswordMutationOptions(options));
};
/**
* This endpoint verifies whether a reset password token exists and is not expired
* @summary Verify a reset password token

View File

@@ -1,49 +0,0 @@
.header {
display: flex;
align-items: center;
justify-content: space-between;
width: 100%;
gap: 8px;
}
.tooltipContent {
--tooltip-z-index: 2100;
}
.dropdownContent {
--dropdown-menu-content-z-index: 2100;
}
.leftSection {
display: flex;
align-items: center;
gap: 8px;
}
.divider {
height: 16px;
margin: 0;
}
.timestamp {
font-family: 'Geist Mono', monospace;
font-size: var(--font-size-sm);
font-weight: var(--font-weight-normal);
color: var(--l1-foreground);
letter-spacing: -0.07px;
}
.actions {
display: flex;
align-items: center;
gap: 8px;
}
.arrows {
display: flex;
align-items: center;
gap: 2px;
padding: 2px 6px;
border-radius: 6px;
box-shadow: 0 1px 4px 0 rgba(0, 0, 0, 0.1);
}

View File

@@ -1,153 +0,0 @@
import { Button } from '@signozhq/ui/button';
import { Divider } from '@signozhq/ui/divider';
import { DropdownMenuSimple as Dropdown } from '@signozhq/ui/dropdown-menu';
import { Typography } from '@signozhq/ui/typography';
import { TooltipSimple } from '@signozhq/ui/tooltip';
import { DATE_TIME_FORMATS } from 'constants/dateTimeFormats';
import { aggregateAttributesResourcesToString } from 'container/LogDetailedView/utils';
import { toast } from '@signozhq/ui/sonner';
import { useCopyLogLink } from 'hooks/logs/useCopyLogLink';
import {
ChevronDown,
ChevronUp,
Compass,
Copy,
Ellipsis,
Link,
} from '@signozhq/icons';
import { useTimezone } from 'providers/Timezone';
import { ILog } from 'types/api/logs/log';
import { MouseEvent, MouseEventHandler } from 'react';
import { useCopyToClipboard } from 'react-use';
import styles from './LogDetailsHeader.module.scss';
const TOOLTIP_CONTENT_PROPS = { className: styles.tooltipContent };
interface LogDetailsHeaderProps {
log: ILog;
onNavigatePrev: () => void;
onNavigateNext: () => void;
isPrevDisabled: boolean;
isNextDisabled: boolean;
showOpenInExplorer?: boolean;
onOpenInExplorer?: MouseEventHandler;
}
function LogDetailsHeader({
log,
onNavigatePrev,
onNavigateNext,
isPrevDisabled,
isNextDisabled,
showOpenInExplorer = false,
onOpenInExplorer,
}: LogDetailsHeaderProps): JSX.Element {
const [, copyToClipboard] = useCopyToClipboard();
const { onLogCopy } = useCopyLogLink(log?.id);
const { formatTimezoneAdjustedTimestamp } = useTimezone();
const handleCopyLog = (): void => {
copyToClipboard(aggregateAttributesResourcesToString(log));
toast.success('Copied to clipboard', { position: 'bottom-right' });
};
const menuItems = [
{
key: 'copy-log',
label: 'Copy log',
icon: <Copy size={14} />,
onClick: handleCopyLog,
},
{
key: 'copy-link',
label: 'Copy link to log',
icon: <Link size={14} />,
onClick: (): void => onLogCopy(),
},
];
return (
<div className={styles.header} data-log-detail-ignore="true">
<div className={styles.leftSection}>
<Divider type="vertical" className={styles.divider} />
<Typography.Text
className={styles.timestamp}
data-testid="log-details-header-timestamp"
>
{formatTimezoneAdjustedTimestamp(
log.date ?? log.timestamp,
DATE_TIME_FORMATS.DASH_DATETIME,
)}
</Typography.Text>
</div>
<div className={styles.actions}>
{showOpenInExplorer && (
<Button
variant="outlined"
color="secondary"
prefix={<Compass size={16} />}
onClick={onOpenInExplorer}
>
Open in Explorer
</Button>
)}
<Dropdown
menu={{ items: menuItems }}
align="end"
className={styles.dropdownContent}
onClick={(e: MouseEvent): void => e.stopPropagation()}
>
<Button
variant="link"
color="secondary"
prefix={<Ellipsis size={16} />}
data-testid="log-details-header-menu"
/>
</Dropdown>
<div className={styles.arrows}>
<TooltipSimple
title="Move to previous log"
side="top"
open={isPrevDisabled ? false : undefined}
tooltipContentProps={TOOLTIP_CONTENT_PROPS}
>
<Button
variant="outlined"
color="secondary"
prefix={<ChevronUp size={14} />}
disabled={isPrevDisabled}
onClick={onNavigatePrev}
data-testid="log-details-header-prev"
/>
</TooltipSimple>
<TooltipSimple
title="Move to next log"
side="top"
open={isNextDisabled ? false : undefined}
tooltipContentProps={TOOLTIP_CONTENT_PROPS}
>
<Button
variant="outlined"
color="secondary"
prefix={<ChevronDown size={14} />}
disabled={isNextDisabled}
onClick={onNavigateNext}
data-testid="log-details-header-next"
/>
</TooltipSimple>
</div>
</div>
</div>
);
}
LogDetailsHeader.defaultProps = {
showOpenInExplorer: false,
onOpenInExplorer: undefined,
};
export default LogDetailsHeader;

View File

@@ -1,52 +0,0 @@
import { useCallback, useMemo } from 'react';
import { ILog } from 'types/api/logs/log';
interface UseLogNavigationParams {
logs?: ILog[];
activeLogId: string;
onNavigateLog?: (log: ILog) => void;
onScrollToLog?: (id: string) => void;
}
interface UseLogNavigationReturn {
goToPrev: () => void;
goToNext: () => void;
isPrevDisabled: boolean;
isNextDisabled: boolean;
}
export function useLogNavigation({
logs,
activeLogId,
onNavigateLog,
onScrollToLog,
}: UseLogNavigationParams): UseLogNavigationReturn {
const currentIndex = useMemo(
() => logs?.findIndex((l) => l.id === activeLogId) ?? -1,
[logs, activeLogId],
);
const canNavigate = !!logs?.length && !!onNavigateLog && currentIndex !== -1;
const isPrevDisabled = !canNavigate || currentIndex <= 0;
const isNextDisabled = !canNavigate || currentIndex >= (logs?.length ?? 0) - 1;
const goToPrev = useCallback((): void => {
if (isPrevDisabled || !logs) {
return;
}
const prev = logs[currentIndex - 1];
onNavigateLog?.(prev);
onScrollToLog?.(prev.id);
}, [isPrevDisabled, logs, currentIndex, onNavigateLog, onScrollToLog]);
const goToNext = useCallback((): void => {
if (isNextDisabled || !logs) {
return;
}
const next = logs[currentIndex + 1];
onNavigateLog?.(next);
onScrollToLog?.(next.id);
}, [isNextDisabled, logs, currentIndex, onNavigateLog, onScrollToLog]);
return { goToPrev, goToNext, isPrevDisabled, isNextDisabled };
}

View File

@@ -1,161 +0,0 @@
import { toast } from '@signozhq/ui/sonner';
import { LOCALSTORAGE } from 'constants/localStorage';
import { render, screen, userEvent } from 'tests/test-utils';
import { ILog } from 'types/api/logs/log';
import LogDetail from '..';
import { VIEW_TYPES } from '../constants';
import { LogDetailProps } from '../LogDetail.interfaces';
jest.mock('@signozhq/ui/sonner', () => ({
toast: { success: jest.fn(), error: jest.fn() },
}));
// The flag to be removed later
jest.mock('../constants', () => ({
...jest.requireActual('../constants'),
isLogDetailsV2: true,
}));
const mockLog: ILog = {
id: 'log-1',
timestamp: '2024-01-15T09:45:30Z',
date: '2024-01-15T09:45:30Z',
body: 'test log body',
severityText: 'INFO',
severityNumber: 9,
traceFlags: 0,
traceId: '',
spanID: '',
attributesString: {},
attributesInt: {},
attributesFloat: {},
resources_string: {},
scope_string: {},
attributes_string: {},
severity_text: 'INFO',
severity_number: 9,
};
const makeLog = (id: string): ILog => ({ ...mockLog, id });
function renderDrawer(props: Partial<LogDetailProps> = {}): void {
render(
<LogDetail
log={mockLog}
selectedTab={VIEW_TYPES.OVERVIEW}
onAddToQuery={jest.fn()}
onClickActionItem={jest.fn()}
onClose={jest.fn()}
{...props}
/>,
);
}
describe('LogDetail drawer — header (isLogDetailsV2)', () => {
afterEach(() => {
jest.clearAllMocks();
localStorage.clear();
});
it('renders the revamped header when a log is provided', () => {
renderDrawer();
expect(screen.getByTestId('log-details-header-menu')).toBeInTheDocument();
expect(screen.getByTestId('log-details-header-prev')).toBeInTheDocument();
expect(screen.getByTestId('log-details-header-next')).toBeInTheDocument();
});
it('shows the log timestamp formatted (DASH_DATETIME) in the header', () => {
// Pin the timezone to UTC so the formatted output is deterministic across
// machines/CI (Jest doesn't fix a TZ).
localStorage.setItem(LOCALSTORAGE.PREFERRED_TIMEZONE, 'UTC');
renderDrawer();
// mockLog date is 2024-01-15T09:45:30Z → DASH_DATETIME in UTC.
expect(screen.getByTestId('log-details-header-timestamp')).toHaveTextContent(
'Jan 15, 2024 ⎯ 09:45:30',
);
});
it('copies the log link from the ⋯ menu', async () => {
const user = userEvent.setup({ pointerEventsCheck: 0 });
renderDrawer();
await user.click(screen.getByTestId('log-details-header-menu'));
await user.click(await screen.findByText('Copy link to log'));
expect(toast.success).toHaveBeenCalled();
});
it('copies the log from the ⋯ menu', async () => {
const user = userEvent.setup({ pointerEventsCheck: 0 });
renderDrawer();
await user.click(screen.getByTestId('log-details-header-menu'));
await user.click(await screen.findByText('Copy log'));
expect(toast.success).toHaveBeenCalledWith('Copied to clipboard', {
position: 'bottom-right',
});
});
it('shows "Open in Explorer" when a handleOpenInExplorer handler is provided', () => {
renderDrawer({ handleOpenInExplorer: jest.fn() });
expect(screen.getByText('Open in Explorer')).toBeInTheDocument();
});
it('hides "Open in Explorer" when no handleOpenInExplorer handler is provided', () => {
renderDrawer();
expect(screen.queryByText('Open in Explorer')).not.toBeInTheDocument();
});
it('navigates to the next / previous log with the Down / Up arrow keys', async () => {
const user = userEvent.setup({ pointerEventsCheck: 0 });
const logs = [makeLog('log-0'), makeLog('log-1'), makeLog('log-2')];
const onNavigateLog = jest.fn();
const onScrollToLog = jest.fn();
// Active log is the middle one so both directions are available.
renderDrawer({ log: logs[1], logs, onNavigateLog, onScrollToLog });
await user.keyboard('{ArrowDown}');
expect(onNavigateLog).toHaveBeenLastCalledWith(logs[2]);
expect(onScrollToLog).toHaveBeenLastCalledWith('log-2');
await user.keyboard('{ArrowUp}');
expect(onNavigateLog).toHaveBeenLastCalledWith(logs[0]);
expect(onScrollToLog).toHaveBeenLastCalledWith('log-0');
});
it('does not navigate past the first log on ArrowUp', async () => {
const user = userEvent.setup({ pointerEventsCheck: 0 });
const logs = [makeLog('log-0'), makeLog('log-1')];
const onNavigateLog = jest.fn();
renderDrawer({ log: logs[0], logs, onNavigateLog });
await user.keyboard('{ArrowUp}');
expect(onNavigateLog).not.toHaveBeenCalled();
});
it('navigates via the header up / down buttons and disables them at boundaries', async () => {
const user = userEvent.setup({ pointerEventsCheck: 0 });
const logs = [makeLog('log-0'), makeLog('log-1')];
const onNavigateLog = jest.fn();
// Active log is the first one.
renderDrawer({ log: logs[0], logs, onNavigateLog });
expect(screen.getByTestId('log-details-header-prev')).toBeDisabled();
expect(screen.getByTestId('log-details-header-next')).toBeEnabled();
await user.click(screen.getByTestId('log-details-header-next'));
expect(onNavigateLog).toHaveBeenLastCalledWith(logs[1]);
});
});

View File

@@ -1,10 +1,3 @@
import getLocalStorage from 'api/browser/localstorage/get';
import { LOCALSTORAGE } from 'constants/localStorage';
// Temp feature flag before actual roll-out
export const isLogDetailsV2 =
getLocalStorage(LOCALSTORAGE.LOG_DETAILS_V2) === 'true';
export const VIEW_TYPES = {
OVERVIEW: 'OVERVIEW',
JSON: 'JSON',

View File

@@ -8,9 +8,7 @@ import { ToggleGroupSimple } from '@signozhq/ui/toggle-group';
import { Divider } from '@signozhq/ui/divider';
import { Typography } from '@signozhq/ui/typography';
import cx from 'classnames';
import LogStateIndicator, {
LogType,
} from 'components/Logs/LogStateIndicator/LogStateIndicator';
import { LogType } from 'components/Logs/LogStateIndicator/LogStateIndicator';
import QuerySearch from 'components/QueryBuilderV2/QueryV2/QuerySearch/QuerySearch';
import { convertExpressionToFilters } from 'components/QueryBuilderV2/utils';
import { FeatureKeys } from 'constants/features';
@@ -25,7 +23,6 @@ import {
} from 'container/LogDetailedView/utils';
import useInitialQuery from 'container/LogsExplorerContext/useInitialQuery';
import { useOptionsMenu } from 'container/OptionsMenu';
import { FontSize } from 'container/OptionsMenu/types';
import { useCopyLogLink } from 'hooks/logs/useCopyLogLink';
import { useQueryBuilder } from 'hooks/queryBuilder/useQueryBuilder';
import { useIsDarkMode } from 'hooks/useDarkMode';
@@ -51,10 +48,8 @@ import { ILogBody } from 'types/api/logs/log';
import { Query, TagFilter } from 'types/api/queryBuilder/queryBuilderData';
import { DataSource, StringOperators } from 'types/common/queryBuilder';
import { isLogDetailsV2, RESOURCE_KEYS, VIEW_TYPES, VIEWS } from './constants';
import { RESOURCE_KEYS, VIEW_TYPES, VIEWS } from './constants';
import { LogDetailInnerProps, LogDetailProps } from './LogDetail.interfaces';
import LogDetailsHeader from './LogDetailsHeader/LogDetailsHeader';
import { useLogNavigation } from './LogDetailsHeader/useLogNavigation';
import './LogDetails.styles.scss';
@@ -101,8 +96,7 @@ function LogDetailInner({
target.closest('[data-log-detail-ignore="true"]') ||
target.closest('.cm-tooltip-autocomplete') ||
target.closest('.drawer-popover') ||
target.closest('.query-status-popover') ||
target.closest('[data-radix-popper-content-wrapper]')
target.closest('.query-status-popover')
) {
return;
}
@@ -118,30 +112,49 @@ function LogDetailInner({
};
}, [onClose]);
const { goToPrev, goToNext, isPrevDisabled, isNextDisabled } =
useLogNavigation({
logs,
activeLogId: log.id,
onNavigateLog,
onScrollToLog,
});
// Keyboard navigation - handle up/down arrow keys. Only listen in the OVERVIEW
// tab so we don't hijack arrow keys from the JSON editor / context view.
// Keyboard navigation - handle up/down arrow keys
// Only listen when in OVERVIEW tab
// eslint-disable-next-line sonarjs/cognitive-complexity
useEffect(() => {
if (selectedView !== VIEW_TYPES.OVERVIEW) {
return undefined;
if (
!logs ||
!onNavigateLog ||
logs.length === 0 ||
selectedView !== VIEW_TYPES.OVERVIEW
) {
return;
}
const handleKeyDown = (e: KeyboardEvent): void => {
const currentIndex = logs.findIndex((l) => l.id === log.id);
if (currentIndex === -1) {
return;
}
if (e.key === 'ArrowUp') {
e.preventDefault();
e.stopPropagation();
goToPrev();
// Navigate to previous log
if (currentIndex > 0) {
const prevLog = logs[currentIndex - 1];
onNavigateLog(prevLog);
// Trigger scroll to the log element
if (onScrollToLog) {
onScrollToLog(prevLog.id);
}
}
} else if (e.key === 'ArrowDown') {
e.preventDefault();
e.stopPropagation();
goToNext();
// Navigate to next log
if (currentIndex < logs.length - 1) {
const nextLog = logs[currentIndex + 1];
onNavigateLog(nextLog);
// Trigger scroll to the log element
if (onScrollToLog) {
onScrollToLog(nextLog.id);
}
}
}
};
@@ -149,7 +162,7 @@ function LogDetailInner({
return (): void => {
document.removeEventListener('keydown', handleKeyDown);
};
}, [selectedView, goToPrev, goToNext]);
}, [log.id, logs, onNavigateLog, onScrollToLog, selectedView]);
const listQuery = useMemo(() => {
if (!stagedQuery || stagedQuery.builder.queryData.length < 1) {
@@ -290,6 +303,33 @@ function LogDetailInner({
};
const logType = log?.attributes_string?.log_level || LogType.INFO;
const currentLogIndex = logs ? logs.findIndex((l) => l.id === log.id) : -1;
const isPrevDisabled =
!logs || !onNavigateLog || logs.length === 0 || currentLogIndex <= 0;
const isNextDisabled =
!logs ||
!onNavigateLog ||
logs.length === 0 ||
currentLogIndex === logs.length - 1;
type HandleNavigateLogParams = {
direction: 'next' | 'previous';
};
const handleNavigateLog = ({ direction }: HandleNavigateLogParams): void => {
if (!logs || !onNavigateLog || currentLogIndex === -1) {
return;
}
if (direction === 'previous' && !isPrevDisabled) {
const prevLog = logs[currentLogIndex - 1];
onNavigateLog(prevLog);
onScrollToLog?.(prevLog.id);
} else if (direction === 'next' && !isNextDisabled) {
const nextLog = logs[currentLogIndex + 1];
onNavigateLog(nextLog);
onScrollToLog?.(nextLog.id);
}
};
return (
<Drawer
@@ -298,69 +338,57 @@ function LogDetailInner({
maskClosable={false}
getContainer={getContainer}
title={
isLogDetailsV2 ? (
<LogDetailsHeader
log={log}
onNavigatePrev={goToPrev}
onNavigateNext={goToNext}
isPrevDisabled={isPrevDisabled}
isNextDisabled={isNextDisabled}
showOpenInExplorer={!!handleOpenInExplorer}
onOpenInExplorer={handleOpenInExplorer}
/>
) : (
<div className="log-detail-drawer__title" data-log-detail-ignore="true">
<div className="log-detail-drawer__title-left">
<Divider type="vertical" className={cx('log-type-indicator', LogType)} />
<Typography.Text className="title">Log details</Typography.Text>
</div>
<div className="log-detail-drawer__title-right">
<div className="log-arrows">
<Tooltip
title={isPrevDisabled ? '' : 'Move to previous log'}
placement="top"
mouseLeaveDelay={0}
>
<Button
variant="outlined"
color="secondary"
prefix={<ChevronUp size={14} />}
className="log-arrow-btn log-arrow-btn-up"
disabled={isPrevDisabled}
onClick={goToPrev}
/>
</Tooltip>
<Tooltip
title={isNextDisabled ? '' : 'Move to next log'}
placement="top"
mouseLeaveDelay={0}
>
<Button
variant="outlined"
color="secondary"
prefix={<ChevronDown size={14} />}
className="log-arrow-btn log-arrow-btn-down"
disabled={isNextDisabled}
onClick={goToNext}
/>
</Tooltip>
</div>
{handleOpenInExplorer && (
<div>
<Button
variant="outlined"
color="secondary"
prefix={<Compass size={16} />}
className="open-in-explorer-btn"
onClick={handleOpenInExplorer}
>
Open in Explorer
</Button>
</div>
)}
</div>
<div className="log-detail-drawer__title" data-log-detail-ignore="true">
<div className="log-detail-drawer__title-left">
<Divider type="vertical" className={cx('log-type-indicator', LogType)} />
<Typography.Text className="title">Log details</Typography.Text>
</div>
)
<div className="log-detail-drawer__title-right">
<div className="log-arrows">
<Tooltip
title={isPrevDisabled ? '' : 'Move to previous log'}
placement="top"
mouseLeaveDelay={0}
>
<Button
variant="outlined"
color="secondary"
prefix={<ChevronUp size={14} />}
className="log-arrow-btn log-arrow-btn-up"
disabled={isPrevDisabled}
onClick={(): void => handleNavigateLog({ direction: 'previous' })}
/>
</Tooltip>
<Tooltip
title={isNextDisabled ? '' : 'Move to next log'}
placement="top"
mouseLeaveDelay={0}
>
<Button
variant="outlined"
color="secondary"
prefix={<ChevronDown size={14} />}
className="log-arrow-btn log-arrow-btn-down"
disabled={isNextDisabled}
onClick={(): void => handleNavigateLog({ direction: 'next' })}
/>
</Tooltip>
</div>
{handleOpenInExplorer && (
<div>
<Button
variant="outlined"
color="secondary"
prefix={<Compass size={16} />}
className="open-in-explorer-btn"
onClick={handleOpenInExplorer}
>
Open in Explorer
</Button>
</div>
)}
</div>
</div>
}
placement="right"
onClose={drawerCloseHandler}
@@ -379,15 +407,7 @@ function LogDetailInner({
data-testid="log-detail-drawer"
>
<div className="log-detail-drawer__log">
{isLogDetailsV2 ? (
<LogStateIndicator
severityText={log.severity_text}
severityNumber={log.severity_number}
fontSize={options?.fontSize ?? FontSize.MEDIUM}
/>
) : (
<Divider type="vertical" className={cx('log-type-indicator', logType)} />
)}
<Divider type="vertical" className={cx('log-type-indicator', logType)} />
<Tooltip
title={removeEscapeCharacters(logBody)}
placement="left"
@@ -463,25 +483,22 @@ function LogDetailInner({
</Tooltip>
)}
{/* V2 moves copy actions into the header ⋯ menu */}
{!isLogDetailsV2 && (
<Tooltip
title={selectedView === VIEW_TYPES.JSON ? 'Copy JSON' : 'Copy Log Link'}
placement="topLeft"
aria-label={
selectedView === VIEW_TYPES.JSON ? 'Copy JSON' : 'Copy Log Link'
}
mouseLeaveDelay={0}
>
<Button
variant="link"
color="secondary"
size="sm"
prefix={<Copy size={12} />}
onClick={selectedView === VIEW_TYPES.JSON ? handleJSONCopy : onLogCopy}
/>
</Tooltip>
)}
<Tooltip
title={selectedView === VIEW_TYPES.JSON ? 'Copy JSON' : 'Copy Log Link'}
placement="topLeft"
aria-label={
selectedView === VIEW_TYPES.JSON ? 'Copy JSON' : 'Copy Log Link'
}
mouseLeaveDelay={0}
>
<Button
variant="link"
color="secondary"
size="sm"
prefix={<Copy size={12} />}
onClick={selectedView === VIEW_TYPES.JSON ? handleJSONCopy : onLogCopy}
/>
</Tooltip>
</div>
</div>
{isFilterVisible && contextQuery?.builder.queryData[0] && (

View File

@@ -13,7 +13,6 @@ export enum LOCALSTORAGE {
TRACES_LIST_COLUMNS = 'TRACES_LIST_COLUMNS',
LOGS_LIST_COLUMNS = 'LOGS_LIST_COLUMNS',
LOGS_LIST_COLUMN_SIZING = 'LOGS_LIST_COLUMN_SIZING',
LOG_DETAILS_V2 = 'LOG_DETAILS_V2',
LOGGED_IN_USER_NAME = 'LOGGED_IN_USER_NAME',
LOGGED_IN_USER_EMAIL = 'LOGGED_IN_USER_EMAIL',
CHAT_SUPPORT = 'CHAT_SUPPORT',

View File

@@ -1,4 +1,4 @@
import { MouseEvent } from 'react';
import { MouseEventHandler } from 'react';
import { ILog } from 'types/api/logs/log';
import { DataTypes } from 'types/api/queryBuilder/queryAutocompleteResponse';
@@ -11,7 +11,7 @@ export type UseCopyLogLink = {
isHighlighted: boolean;
isLogsExplorerPage: boolean;
activeLogId: string | null;
onLogCopy: (event?: MouseEvent<HTMLElement>) => void;
onLogCopy: MouseEventHandler<HTMLElement>;
onClearActiveLog: () => void;
};

View File

@@ -1,4 +1,10 @@
import { MouseEvent, useCallback, useEffect, useMemo, useState } from 'react';
import {
MouseEventHandler,
useCallback,
useEffect,
useMemo,
useState,
} from 'react';
// eslint-disable-next-line no-restricted-imports
import { useSelector } from 'react-redux';
import { useLocation } from 'react-router-dom';
@@ -40,14 +46,14 @@ export const useCopyLogLink = (logId?: string): UseCopyLogLink => {
[pathname],
);
const onLogCopy = useCallback(
(event?: MouseEvent<HTMLElement>): void => {
const onLogCopy: MouseEventHandler<HTMLElement> = useCallback(
(event) => {
if (!logId) {
return;
}
event?.preventDefault();
event?.stopPropagation();
event.preventDefault();
event.stopPropagation();
urlQuery.delete(QueryParams.activeLogId);
urlQuery.delete(QueryParams.relativeTime);
@@ -60,7 +66,7 @@ export const useCopyLogLink = (logId?: string): UseCopyLogLink => {
setCopy(link);
toast.success('Copied to clipboard', { position: 'bottom-right' });
toast.success('Copied to clipboard', { position: 'top-right' });
},
[logId, urlQuery, minTime, maxTime, pathname, setCopy],
);

View File

@@ -249,7 +249,7 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
}
if err := router.Handle("/api/v1/resetPassword", handler.New(provider.authzMiddleware.OpenAccess(provider.userHandler.ResetPassword), handler.OpenAPIDef{
ID: "ResetPasswordDeprecated",
ID: "ResetPassword",
Tags: []string{"users"},
Summary: "Reset password",
Description: "This endpoint resets the password by token",
@@ -259,7 +259,7 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusConflict},
Deprecated: true,
Deprecated: false,
SecuritySchemes: []handler.OpenAPISecurityScheme{},
})).Methods(http.MethodPost).GetError(); err != nil {
return err
@@ -299,23 +299,6 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/factor_password/reset", handler.New(provider.authzMiddleware.OpenAccess(provider.userHandler.ResetPassword), handler.OpenAPIDef{
ID: "ResetPassword",
Tags: []string{"users"},
Summary: "Reset password",
Description: "This endpoint resets the password using a single use reset password token",
Request: new(types.PostableResetPassword),
RequestContentType: "application/json",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: []handler.OpenAPISecurityScheme{},
})).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/users/{id}/roles", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.GetRolesByUserID), handler.OpenAPIDef{
ID: "GetRolesByUserID",
Tags: []string{"users"},

View File

@@ -40,7 +40,10 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
// Rule state history fields have no family support, so every logical field
// is single-member and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -64,6 +64,23 @@ func (m *fieldMapper) ColumnFor(ctx context.Context, _ valuer.UUID, _, _ uint64,
return []*schema.Column{col}, nil
}
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
// Real columns always exist; labels are checked for key membership.
func (m *fieldMapper) ExistsFor(ctx context.Context, _ valuer.UUID, _, _ uint64, key *telemetrytypes.TelemetryFieldKey, exists bool) (string, error) {
col, err := m.getColumn(ctx, key)
if err != nil {
return "", err
}
if col.Name != "labels" || key.Name == "labels" {
return "true", nil
}
pred := fmt.Sprintf("has(JSONExtractKeys(labels), '%s')", strings.ReplaceAll(key.Name, "'", "\\'"))
if exists {
return pred, nil
}
return "not " + pred, nil
}
func (m *fieldMapper) ColumnExpressionFor(ctx context.Context, orgID valuer.UUID, tsStart, tsEnd uint64, field *telemetrytypes.TelemetryFieldKey, _ telemetrytypes.FieldDataType, _ map[string][]*telemetrytypes.TelemetryFieldKey) (string, error) {
colName, err := m.FieldFor(ctx, orgID, tsStart, tsEnd, field)
if err != nil {

View File

@@ -391,7 +391,7 @@ func (handler *handler) ResetPassword(w http.ResponseWriter, r *http.Request) {
defer cancel()
req := new(types.PostableResetPassword)
if err := binding.JSON.BindBody(r.Body, req); err != nil {
if err := json.NewDecoder(r.Body).Decode(req); err != nil {
render.Error(w, err)
return
}

View File

@@ -0,0 +1,43 @@
package querybuilder
import (
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// ExpandKeySelectorsForFamilies appends selectors for the other members of any
// semantic-convention family a selector names, so the metadata fetched for a
// query contains every spelling MatchingLogicalFields may group. It is the
// resolution layer's prefetch: statement builders call it after deriving
// selectors, and the metadata store stays family-blind (autocomplete responses
// keep the literal spelling the user typed). Only trace selectors expand today,
// matching family support; fuzzy (search-style) selectors never do.
func ExpandKeySelectorsForFamilies(selectors []*telemetrytypes.FieldKeySelector) []*telemetrytypes.FieldKeySelector {
out := selectors
seen := make(map[string]bool, len(selectors))
for _, selector := range selectors {
seen[selector.Name] = true
}
for _, selector := range selectors {
if selector.Signal != telemetrytypes.SignalTraces ||
selector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeFuzzy {
continue
}
members := semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: selector.Name,
Signal: selector.Signal,
FieldContext: selector.FieldContext,
})
for _, member := range members {
if seen[member] {
continue
}
seen[member] = true
expanded := *selector
expanded.Name = member
out = append(out, &expanded)
}
}
return out
}

View File

@@ -21,24 +21,25 @@ const (
hasTokenFunctionDocURL = "https://signoz.io/docs/userguide/functions-reference/#hastoken-function"
)
// ResolveKeys picks which matching field keys a filter term builds conditions for.
// With 0 or 1 match it returns the input unchanged and no warning. When a name is
// ambiguous it returns a warning; a resource+attribute mix defaults to the resource
// keys (the common intent), noted in the warning.
func ResolveKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeysForName []*telemetrytypes.TelemetryFieldKey) ([]*telemetrytypes.TelemetryFieldKey, string) {
if len(fieldKeysForName) <= 1 {
return fieldKeysForName, ""
// ResolveLogicalFields picks which logical fields a filter term builds conditions
// for. With 0 or 1 field it returns the input unchanged and no warning. When a
// name is ambiguous (several logical fields — a family is one field and never
// ambiguous with itself) it returns a warning; a resource+attribute mix defaults
// to the resource fields (the common intent), noted in the warning.
func ResolveLogicalFields(field *telemetrytypes.TelemetryFieldKey, logicalFields []*telemetrytypes.LogicalField) ([]*telemetrytypes.LogicalField, string) {
if len(logicalFields) <= 1 {
return logicalFields, ""
}
warning := fmt.Sprintf(
"Key `%s` is ambiguous, found %d different combinations of field context / data type: %v.",
field.Name,
len(fieldKeysForName),
fieldKeysForName,
len(logicalFields),
logicalFields,
)
hasResource, hasAttribute := false, false
for _, item := range fieldKeysForName {
for _, item := range logicalFields {
switch item.FieldContext {
case telemetrytypes.FieldContextResource:
hasResource = true
@@ -49,18 +50,40 @@ func ResolveKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeysForName []*te
// when there is both resource and attribute context, default to resource only
if hasResource && hasAttribute {
filteredKeys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fieldKeysForName))
for _, item := range fieldKeysForName {
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
for _, item := range logicalFields {
if item.FieldContext == telemetrytypes.FieldContextResource {
filteredKeys = append(filteredKeys, item)
filtered = append(filtered, item)
}
}
fieldKeysForName = filteredKeys
logicalFields = filtered
warning += " " + "Using `resource` context by default. To query attributes explicitly, " +
fmt.Sprintf("use the fully qualified name (e.g., 'attribute.%s')", field.Name)
}
return fieldKeysForName, warning
return logicalFields, warning
}
// WrapAsLogicalFields wraps physical keys (candidate or synthesized) as
// single-member logical fields addressed by the requested spelling.
func WrapAsLogicalFields(requestedName string, keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
fields := make([]*telemetrytypes.LogicalField, 0, len(keys))
for _, key := range keys {
fields = append(fields, telemetrytypes.SingleLogicalField(requestedName, key))
}
return fields
}
// SingleKeys flattens logical fields to their single members. It is the adapter
// for signals whose fields are single-member by construction (every signal
// without family support); a condition builder that uses it compiles per
// physical key exactly as before.
func SingleKeys(fields []*telemetrytypes.LogicalField) []*telemetrytypes.TelemetryFieldKey {
keys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fields))
for _, field := range fields {
keys = append(keys, field.Single())
}
return keys
}
// NewKeyNotFoundError builds the error a condition builder returns when a filter term

View File

@@ -0,0 +1,94 @@
package querybuilder
import (
"context"
"fmt"
"strings"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
// The two functions below are the only place family expressions are built.
// They compose exclusively from the mapper's per-key primitives (FieldFor,
// ExistsFor), so every member honors its own storage: materialized columns,
// evolutions, and JSON plans ride the member keys, and a signal supports
// families the moment its primitives are correct.
// LogicalValueExpr returns the value expression for a resolved logical field:
// the member's own expression for a single-member field, and a current-first
// merge across the members' expressions for a family.
func LogicalValueExpr(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,
fm qbtypes.FieldMapper,
logical *telemetrytypes.LogicalField,
) (string, error) {
if !logical.IsFamily() {
return fm.FieldFor(ctx, orgID, tsStart, tsEnd, logical.Single())
}
memberExprs := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
expr, err := fm.FieldFor(ctx, orgID, tsStart, tsEnd, member)
if err != nil {
return "", err
}
memberExprs = append(memberExprs, expr)
}
if logical.FieldDataType == telemetrytypes.FieldDataTypeString {
// The trailing '' keeps single-key semantics for rows without any
// member: string maps read '' for an absent key, and negative
// operators must keep including such rows (see AddDefaultExistsFilter).
values := make([]string, 0, len(memberExprs))
for _, expr := range memberExprs {
values = append(values, fmt.Sprintf("NULLIF(%s, '')", expr))
}
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
}
// Numeric and boolean maps return zero for an absent key. If a family of
// either type is enabled, this tail must become zero too.
branches := make([]string, 0, len(logical.Members)*2)
for i, member := range logical.Members {
guard, err := fm.ExistsFor(ctx, orgID, tsStart, tsEnd, member, true)
if err != nil {
return "", err
}
branches = append(branches, guard, memberExprs[i])
}
return "multiIf(" + strings.Join(branches, ", ") + ", NULL)", nil
}
// LogicalExistsExpr returns the existence predicate for a resolved logical
// field: the member's own predicate for a single-member field, presence of
// any member for a family.
func LogicalExistsExpr(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,
fm qbtypes.FieldMapper,
logical *telemetrytypes.LogicalField,
exists bool,
) (string, error) {
if !logical.IsFamily() {
return fm.ExistsFor(ctx, orgID, tsStart, tsEnd, logical.Single(), exists)
}
guards := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
guard, err := fm.ExistsFor(ctx, orgID, tsStart, tsEnd, member, true)
if err != nil {
return "", err
}
guards = append(guards, guard)
}
combined := "(" + strings.Join(guards, " OR ") + ")"
if exists {
return combined, nil
}
return "NOT " + combined, nil
}

View File

@@ -0,0 +1,95 @@
package querybuilder
import (
"context"
"testing"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// stubFieldMapper provides just the two per-key primitives the shared
// composition builds on; the remaining FieldMapper methods are unused here.
type stubFieldMapper struct{}
func (stubFieldMapper) FieldFor(_ context.Context, _ valuer.UUID, _, _ uint64, key *telemetrytypes.TelemetryFieldKey) (string, error) {
return "value(" + key.Name + ")", nil
}
func (stubFieldMapper) ExistsFor(_ context.Context, _ valuer.UUID, _, _ uint64, key *telemetrytypes.TelemetryFieldKey, exists bool) (string, error) {
if exists {
return "has(" + key.Name + ")", nil
}
return "NOT has(" + key.Name + ")", nil
}
func (stubFieldMapper) ColumnFor(context.Context, valuer.UUID, uint64, uint64, *telemetrytypes.TelemetryFieldKey) ([]*schema.Column, error) {
return nil, qbtypes.ErrColumnNotFound
}
func (stubFieldMapper) ColumnExpressionFor(context.Context, valuer.UUID, uint64, uint64, *telemetrytypes.TelemetryFieldKey, telemetrytypes.FieldDataType, map[string][]*telemetrytypes.TelemetryFieldKey) (string, error) {
return "", qbtypes.ErrColumnNotFound
}
func (stubFieldMapper) CandidateKeys(context.Context, valuer.UUID, *telemetrytypes.TelemetryFieldKey, any, map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
return nil
}
func stringFamily(names ...string) *telemetrytypes.LogicalField {
members := make([]*telemetrytypes.TelemetryFieldKey, 0, len(names))
for _, name := range names {
members = append(members, &telemetrytypes.TelemetryFieldKey{Name: name, FieldDataType: telemetrytypes.FieldDataTypeString})
}
return &telemetrytypes.LogicalField{Name: names[0], FieldDataType: telemetrytypes.FieldDataTypeString, Members: members}
}
func TestLogicalValueExprSingleMemberDelegatesToFieldFor(t *testing.T) {
logical := telemetrytypes.SingleLogicalField("a", &telemetrytypes.TelemetryFieldKey{Name: "a"})
expr, err := LogicalValueExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, logical)
require.NoError(t, err)
assert.Equal(t, "value(a)", expr)
}
func TestLogicalValueExprStringFamilyMergesCurrentFirst(t *testing.T) {
expr, err := LogicalValueExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, stringFamily("current", "old"))
require.NoError(t, err)
// The trailing '' preserves keyless-row semantics for negative operators.
assert.Equal(t, "COALESCE(NULLIF(value(current), ''), NULLIF(value(old), ''), '')", expr)
}
func TestLogicalValueExprNumericFamilyGuardsEveryMember(t *testing.T) {
logical := &telemetrytypes.LogicalField{
Name: "current",
FieldDataType: telemetrytypes.FieldDataTypeNumber,
Members: []*telemetrytypes.TelemetryFieldKey{
{Name: "current", FieldDataType: telemetrytypes.FieldDataTypeNumber},
{Name: "old", FieldDataType: telemetrytypes.FieldDataTypeNumber},
},
}
expr, err := LogicalValueExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, logical)
require.NoError(t, err)
assert.Equal(t, "multiIf(has(current), value(current), has(old), value(old), NULL)", expr)
}
func TestLogicalExistsExprSingleMemberDelegatesToExistsFor(t *testing.T) {
logical := telemetrytypes.SingleLogicalField("a", &telemetrytypes.TelemetryFieldKey{Name: "a"})
expr, err := LogicalExistsExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, logical, false)
require.NoError(t, err)
assert.Equal(t, "NOT has(a)", expr)
}
func TestLogicalExistsExprFamilyIsAnyMemberPresence(t *testing.T) {
family := stringFamily("current", "old")
expr, err := LogicalExistsExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, family, true)
require.NoError(t, err)
assert.Equal(t, "(has(current) OR has(old))", expr)
expr, err = LogicalExistsExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, family, false)
require.NoError(t, err)
assert.Equal(t, "NOT (has(current) OR has(old))", expr)
}

View File

@@ -0,0 +1,165 @@
package querybuilder
import (
"testing"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func traceKey(name string, ctx telemetrytypes.FieldContext) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: ctx,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
func memberNames(logical *telemetrytypes.LogicalField) []string {
names := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
names = append(names, member.Name)
}
return names
}
// The deployment.environment(.name) family (enabled in pkg/semconv) drives the
// grouping tests below.
func TestMatchingLogicalFieldsGroupsFamilyMembers(t *testing.T) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {traceKey("deployment.environment.name", telemetrytypes.FieldContextResource)},
"deployment.environment": {traceKey("deployment.environment", telemetrytypes.FieldContextResource)},
}
for _, requested := range []string{"deployment.environment.name", "deployment.environment"} {
fields := MatchingLogicalFields(&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")
assert.Equal(t, telemetrytypes.FieldContextResource, logical.FieldContext)
assert.True(t, logical.IsFamily())
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, memberNames(logical),
"members are current-first regardless of the requested spelling")
}
}
// Member precedence is the family's current-first order, not lookup arrival
// order: a current-name key found only under its context-prefixed spelling
// arrives in the second lookup pass yet must still sort first.
func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment": {traceKey("deployment.environment", telemetrytypes.FieldContextResource)},
"resource.deployment.environment.name": {traceKey("resource.deployment.environment.name", telemetrytypes.FieldContextResource)},
}
fields := MatchingLogicalFields(&telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
}, fieldKeys)
require.Len(t, fields, 1)
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,
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")},
}
fields := MatchingLogicalFields(&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]))
}
// A family and a genuine same-name collision stack cleanly: the family stays
// one logical field, the collision adds another, and resource preference keeps
// the family as a unit.
func TestResolveLogicalFieldsKeepsFamilyThroughAmbiguity(t *testing.T) {
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {
traceKey("deployment.environment.name", telemetrytypes.FieldContextResource),
traceKey("deployment.environment.name", telemetrytypes.FieldContextAttribute),
},
"deployment.environment": {traceKey("deployment.environment", telemetrytypes.FieldContextResource)},
}
requested := &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}
fields := MatchingLogicalFields(requested, fieldKeys)
require.Len(t, fields, 2, "resource family + attribute collision")
resolved, warning := ResolveLogicalFields(requested, fields)
assert.NotEmpty(t, warning)
require.Len(t, resolved, 1)
assert.Equal(t, telemetrytypes.FieldContextResource, resolved[0].FieldContext)
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, memberNames(resolved[0]))
}
// Members of a family with different data types never merge: the identity
// (signal, context, data type) separates them into distinct logical fields.
func TestMatchingLogicalFieldsNeverMergesAcrossDataTypes(t *testing.T) {
numberKey := traceKey("deployment.environment", telemetrytypes.FieldContextResource)
numberKey.FieldDataType = telemetrytypes.FieldDataTypeNumber
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {traceKey("deployment.environment.name", telemetrytypes.FieldContextResource)},
"deployment.environment": {numberKey},
}
fields := MatchingLogicalFields(&telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 2)
for _, logical := range fields {
assert.False(t, logical.IsFamily())
}
}
func TestExpandKeySelectorsForFamilies(t *testing.T) {
selectors := []*telemetrytypes.FieldKeySelector{
{Name: "deployment.environment.name", Signal: telemetrytypes.SignalTraces, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact},
{Name: "service.name", Signal: telemetrytypes.SignalTraces, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact},
{Name: "deployment.environment.name", Signal: telemetrytypes.SignalLogs, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact},
}
expanded := ExpandKeySelectorsForFamilies(selectors)
names := make([]string, 0, len(expanded))
for _, selector := range expanded {
names = append(names, selector.Name)
}
assert.Equal(t, []string{
"deployment.environment.name",
"service.name",
"deployment.environment.name",
"deployment.environment",
}, names, "one sibling selector for the trace family member; logs and non-family names untouched")
sibling := expanded[len(expanded)-1]
assert.Equal(t, telemetrytypes.SignalTraces, sibling.Signal)
assert.Equal(t, telemetrytypes.FieldSelectorMatchTypeExact, sibling.SelectorMatchType)
}
func TestExpandKeySelectorsForFamiliesDeduplicatesAndSkipsFuzzy(t *testing.T) {
both := []*telemetrytypes.FieldKeySelector{
{Name: "deployment.environment.name", Signal: telemetrytypes.SignalTraces, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact},
{Name: "deployment.environment", Signal: telemetrytypes.SignalTraces, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact},
}
assert.Len(t, ExpandKeySelectorsForFamilies(both), 2, "both spellings already referenced")
fuzzy := []*telemetrytypes.FieldKeySelector{
{Name: "deployment.environment.name", Signal: telemetrytypes.SignalTraces, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeFuzzy},
}
assert.Len(t, ExpandKeySelectorsForFamilies(fuzzy), 1, "fuzzy (search-style) selectors never expand")
}

View File

@@ -56,17 +56,6 @@ func QueryStringToKeysSelectors(query string) []*telemetrytypes.FieldKeySelector
FieldDataType: key.FieldDataType,
})
}
// todo(tushar): consider reverting changes done to this method in below PR to avoid scope specific checks
// https://github.com/SigNoz/signoz/issues/11374
if key.FieldContext == telemetrytypes.FieldContextScope {
keys = append(keys, &telemetrytypes.FieldKeySelector{
Name: key.FieldContext.StringValue() + "." + key.Name,
Signal: key.Signal,
FieldContext: telemetrytypes.FieldContextUnspecified, // this allows 'scope.' prefix for keys with other context as well
FieldDataType: key.FieldDataType,
})
}
}
}

View File

@@ -72,23 +72,6 @@ func TestQueryToKeys(t *testing.T) {
},
},
},
{
query: `scope.version = '1.0.0'`,
expectedKeys: []telemetrytypes.FieldKeySelector{
{
Name: "version",
Signal: telemetrytypes.SignalUnspecified,
FieldContext: telemetrytypes.FieldContextScope,
FieldDataType: telemetrytypes.FieldDataTypeUnspecified,
},
{
Name: "scope.version",
Signal: telemetrytypes.SignalUnspecified,
FieldContext: telemetrytypes.FieldContextUnspecified,
FieldDataType: telemetrytypes.FieldDataTypeUnspecified,
},
},
},
}
for _, testCase := range testCases {

View File

@@ -10,6 +10,7 @@ import (
"github.com/SigNoz/signoz/pkg/errors"
grammar "github.com/SigNoz/signoz/pkg/parser/filterquery/grammar"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -360,7 +361,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
return ErrorConditionLiteral
}
}
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
if !ok {
return ErrorConditionLiteral
}
@@ -379,7 +380,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
// VisitComparison handles all comparison operators.
func (v *filterExpressionVisitor) VisitComparison(ctx *grammar.ComparisonContext) any {
key := v.Visit(ctx.Key()).(*telemetrytypes.TelemetryFieldKey)
matching := MatchingFieldKeys(key, v.fieldKeys)
matching := MatchingLogicalFields(key, v.fieldKeys)
// Handle EXISTS specially
if ctx.EXISTS() != nil {
@@ -675,7 +676,7 @@ func (v *filterExpressionVisitor) VisitFullText(ctx *grammar.FullTextContext) an
v.errors = append(v.errors, "full text search is not supported")
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
if !ok {
return ErrorConditionLiteral
}
@@ -730,7 +731,7 @@ func (v *filterExpressionVisitor) VisitFunctionCall(ctx *grammar.FunctionCallCon
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(key, MatchingFieldKeys(key, v.fieldKeys), operator, value)
conds, ok := v.buildConditions(key, MatchingLogicalFields(key, v.fieldKeys), operator, value)
if !ok {
return ErrorConditionLiteral
}
@@ -922,7 +923,7 @@ func (v *filterExpressionVisitor) VisitKey(ctx *grammar.KeyContext) any {
// buildConditions invokes the condition builder for a filter term, folding its
// warnings/errors into visitor state; returns false if an error was recorded.
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) ([]string, bool) {
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.LogicalField, op qbtypes.FilterOperator, value any) ([]string, bool) {
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter}, op, value, v.builder)
if err != nil {
_, _, _, _, errURL, _ := errors.Unwrapb(err)
@@ -979,30 +980,123 @@ func assignIfEmpty(s *string, value string) {
}
}
// MatchingFieldKeys returns the field keys from the map that match the given key,
// honoring any context/data type the user specified.
func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
fieldKeysForName := []*telemetrytypes.TelemetryFieldKey{}
// familyMemberNames returns the physical spellings to look up for the
// referenced key: the semantic-convention family members (current-first) when
// the key can resolve to traces, else just the requested name. Only trace
// field mappers understand families today; logs and metrics keep the
// requested spelling until theirs land.
func familyMemberNames(field *telemetrytypes.TelemetryFieldKey) []string {
if field.Signal != telemetrytypes.SignalUnspecified && field.Signal != telemetrytypes.SignalTraces {
return []string{field.Name}
}
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: field.FieldContext,
})
}
// match by name; keep items whose context and data type match (unspecified matches any)
for _, item := range fieldKeys[field.Name] {
if (field.FieldContext == telemetrytypes.FieldContextUnspecified || field.FieldContext == item.FieldContext) &&
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || field.FieldDataType == item.FieldDataType) {
fieldKeysForName = append(fieldKeysForName, item)
}
// 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 a single 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 therefore the
// length of the returned slice, and a family is never ambiguous with itself.
// Members alias the metadata map entries; nothing is copied or mutated.
func MatchingLogicalFields(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
members := familyMemberNames(field)
memberRank := make(map[string]int, len(members))
for i, member := range members {
memberRank[member] = i
}
// A context may have been split off a name that legitimately contained it (e.g.
// `attribute.key`); also look up the context-prefixed name so both readings resolve.
if field.FieldContext != telemetrytypes.FieldContextUnspecified {
contextPrefixedFieldName := fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), field.Name)
for _, item := range fieldKeys[contextPrefixedFieldName] {
// Context already matched via the lookup key; only data type needs checking.
if field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || item.FieldDataType == field.FieldDataType {
fieldKeysForName = append(fieldKeysForName, item)
fields := make([]*telemetrytypes.LogicalField, 0)
indexByIdentity := make(map[string]int)
// rank of the family member each physical key matched under; the stored
// name of a context-prefixed match differs from the member name.
ranks := make(map[*telemetrytypes.TelemetryFieldKey]int)
appendMatches := func(lookupName string, memberName string, contextAlreadyMatched bool) {
for _, item := range fieldKeys[lookupName] {
if !contextAlreadyMatched && field.FieldContext != telemetrytypes.FieldContextUnspecified && field.FieldContext != item.FieldContext {
continue
}
if field.FieldDataType != telemetrytypes.FieldDataTypeUnspecified && field.FieldDataType != item.FieldDataType {
continue
}
// A member lookup may have found a same-named field in a scope where
// this family does not apply. Keep exact names, but reject cross-member
// matches outside the generated family scope.
traceFamilyMatch := len(members) > 1 && item.Signal == telemetrytypes.SignalTraces
if memberName != field.Name {
if !traceFamilyMatch {
continue
}
itemSelector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: item.FieldContext,
}
if !slices.Contains(semconv.Members(semconv.KindAttribute, itemSelector), memberName) {
continue
}
}
if !traceFamilyMatch {
fields = append(fields, telemetrytypes.SingleLogicalField(field.Name, item))
continue
}
identity := item.Signal.StringValue() + ";" + item.FieldContext.StringValue() + ";" + item.FieldDataType.StringValue()
index, found := indexByIdentity[identity]
if !found {
index = len(fields)
indexByIdentity[identity] = index
fields = append(fields, &telemetrytypes.LogicalField{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
FieldDataType: item.FieldDataType,
})
}
logical := fields[index]
duplicate := false
for _, existing := range logical.Members {
if existing.Name == item.Name {
duplicate = true
break
}
}
if !duplicate {
ranks[item] = memberRank[memberName]
logical.Members = append(logical.Members, item)
}
}
}
return fieldKeysForName
for _, member := range members {
appendMatches(member, member, false)
}
// A context may have been split off a name that legitimately contained it
// (e.g. `attribute.key`); preserve that alternate reading for every family
// member.
if field.FieldContext != telemetrytypes.FieldContextUnspecified {
for _, member := range members {
appendMatches(fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), member), member, true)
}
}
// Precedence is a property of the family, not of arrival order: members
// sort current-first no matter which lookup pass found them.
for _, logical := range fields {
slices.SortStableFunc(logical.Members, func(a, b *telemetrytypes.TelemetryFieldKey) int {
return ranks[a] - ranks[b]
})
}
return fields
}

View File

@@ -588,9 +588,11 @@ func TestVisitKey(t *testing.T) {
// VisitKey only parses; the condition builder matches, resolves ambiguity
// and decides not-found handling. Replay that here against the generic
// builder behavior (error unless the key is ignored).
matching := MatchingFieldKeys(key, tt.fieldKeys)
keys, warning := ResolveKeys(key, matching)
// 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(key, tt.fieldKeys)
resolved, warning := ResolveLogicalFields(key, matching)
keys := SingleKeys(resolved)
var gotErrors []string
var gotMainErrURL, gotMainWrnURL string
@@ -766,7 +768,8 @@ func (b *resourceConditionBuilder) ConditionFor(
return nil, nil, nil
}
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
keys := SingleKeys(resolved)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -808,7 +811,8 @@ func (b *conditionBuilder) ConditionFor(
return []string{fmt.Sprintf("%s_cond", key.Name)}, nil, nil
}
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
keys := SingleKeys(resolved)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -44,6 +44,70 @@ func keyIndexFilter(key *telemetrytypes.TelemetryFieldKey) any {
return fmt.Sprintf(`%%%s%%`, key.Name)
}
// The three helpers below take the members of one logical field. With a single
// member they render exactly the pre-family shapes; a family widens key/value
// index hints to any-member and presence to any-member (all-absent when negated).
func keyIndexCondition(sb *sqlbuilder.SelectBuilder, column string, members []*telemetrytypes.TelemetryFieldKey) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
conditions = append(conditions, sb.Like(column, keyIndexFilter(member)))
}
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
func valueIndexCondition(
sb *sqlbuilder.SelectBuilder,
column string,
members []*telemetrytypes.TelemetryFieldKey,
op qbtypes.FilterOperator,
value any,
caseInsensitive bool,
) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
patterns := valueForIndexFilter(op, member, value)
switch values := patterns.(type) {
case []string:
for _, pattern := range values {
conditions = append(conditions, sb.Like(column, pattern))
}
default:
if caseInsensitive {
conditions = append(conditions, sb.ILike(column, values))
} else {
conditions = append(conditions, sb.Like(column, values))
}
}
}
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
func memberPresenceCondition(sb *sqlbuilder.SelectBuilder, column string, members []*telemetrytypes.TelemetryFieldKey, exists bool) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
field := fmt.Sprintf("simpleJSONHas(%s, '%s')", column, member.Name)
if exists {
conditions = append(conditions, sb.E(field, true))
} else {
conditions = append(conditions, sb.NE(field, true))
}
}
if exists {
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
return sb.And(conditions...)
}
// SkipResourceFilter is not applicable here: the fingerprint table only stores resource attributes.
func (b *defaultConditionBuilder) ConditionFor(
ctx context.Context,
@@ -57,7 +121,7 @@ func (b *defaultConditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
// has/hasAny/hasAll/hasToken are logs-body-only functions; they never apply to the
// resource fingerprint table, so skip them (the main query still evaluates them).
@@ -65,21 +129,21 @@ func (b *defaultConditionBuilder) ConditionFor(
return nil, nil, nil
}
keys, warning := querybuilder.ResolveKeys(key, matches)
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
}
conds := make([]string, 0, len(keys))
for _, k := range keys {
// the resource fingerprint table only stores resource attributes; keys from
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
// the resource fingerprint table only stores resource attributes; fields from
// any other context contribute no condition and are omitted. An empty result
// (including an unknown key) lets the caller skip this filter entirely.
if k.FieldContext != telemetrytypes.FieldContextResource {
if logical.FieldContext != telemetrytypes.FieldContextResource {
continue
}
cond, err := b.conditionForKey(ctx, startNs, endNs, k, op, value, sb)
cond, err := b.conditionForLogicalField(ctx, startNs, endNs, logical, op, value, sb)
if err != nil {
return nil, nil, err
}
@@ -88,11 +152,11 @@ func (b *defaultConditionBuilder) ConditionFor(
return conds, warnings, nil
}
func (b *defaultConditionBuilder) conditionForKey(
func (b *defaultConditionBuilder) conditionForLogicalField(
ctx context.Context,
startNs uint64,
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
logical *telemetrytypes.LogicalField,
op qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -102,7 +166,7 @@ func (b *defaultConditionBuilder) conditionForKey(
// as we store resource values as string
formattedValue := querybuilder.FormatValueForContains(value)
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, key)
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, logical.Single())
if err != nil {
return "", err
}
@@ -115,10 +179,12 @@ func (b *defaultConditionBuilder) conditionForKey(
// as we have not changed the resource column in the resource fingerprint table.
column := columns[0]
keyIdxFilter := sb.Like(column.Name, keyIndexFilter(key))
valueForIndexFilter := valueForIndexFilter(op, key, value)
members := logical.Members
isFamily := logical.IsFamily()
keyIdxFilter := keyIndexCondition(sb, column.Name, members)
singleValueIndexFilter := valueForIndexFilter(op, members[0], value)
fieldName, err := b.fm.FieldFor(ctx, valuer.UUID{}, startNs, endNs, key)
fieldName, err := querybuilder.LogicalValueExpr(ctx, valuer.UUID{}, startNs, endNs, b.fm, logical)
if err != nil {
return "", err
}
@@ -128,12 +194,17 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.E(fieldName, formattedValue),
keyIdxFilter,
sb.Like(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, members, op, value, false),
), nil
case qbtypes.FilterOperatorNotEqual:
if isFamily {
// A negated value-index hint would drop rows where another member
// holds the value; the fingerprint scan is small enough without it.
return sb.NE(fieldName, formattedValue), nil
}
return sb.And(
sb.NE(fieldName, formattedValue),
sb.NotLike(column.Name, valueForIndexFilter),
sb.NotLike(column.Name, singleValueIndexFilter),
), nil
case qbtypes.FilterOperatorGreaterThan:
return sb.And(sb.GT(fieldName, formattedValue), keyIdxFilter), nil
@@ -148,7 +219,7 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.ILike(fieldName, formattedValue),
keyIdxFilter,
sb.ILike(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, members, op, value, true),
), nil
case qbtypes.FilterOperatorNotLike, qbtypes.FilterOperatorNotILike:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
@@ -185,13 +256,11 @@ func (b *defaultConditionBuilder) conditionForKey(
inConditions = append(inConditions, sb.E(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.Or(inConditions...)
valConditions := make([]string, 0, len(values))
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.Like(column.Name, v))
}
}
mainCondition = sb.And(mainCondition, keyIdxFilter, sb.Or(valConditions...))
mainCondition = sb.And(
mainCondition,
keyIdxFilter,
valueIndexCondition(sb, column.Name, members, op, value, false),
)
return mainCondition, nil
case qbtypes.FilterOperatorNotIn:
@@ -204,8 +273,13 @@ func (b *defaultConditionBuilder) conditionForKey(
notInConditions = append(notInConditions, sb.NE(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.And(notInConditions...)
if isFamily {
// A negated value-index hint would drop rows where another member
// holds the value; the fingerprint scan is small enough without it.
return mainCondition, nil
}
valConditions := make([]string, 0, len(values))
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
if valuesForIndexFilter, ok := singleValueIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.NotLike(column.Name, v))
}
@@ -215,13 +289,11 @@ func (b *defaultConditionBuilder) conditionForKey(
case qbtypes.FilterOperatorExists:
return sb.And(
sb.E(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
memberPresenceCondition(sb, column.Name, members, true),
keyIdxFilter,
), nil
case qbtypes.FilterOperatorNotExists:
return sb.And(
sb.NE(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
), nil
return memberPresenceCondition(sb, column.Name, members, false), nil
case qbtypes.FilterOperatorRegexp:
return sb.And(
@@ -237,7 +309,7 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.ILike(fieldName, fmt.Sprintf(`%%%s%%`, formattedValue)),
keyIdxFilter,
sb.ILike(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, members, op, value, true),
), nil
case qbtypes.FilterOperatorNotContains:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else

View File

@@ -0,0 +1,90 @@
package resourcefilter
import (
"context"
"testing"
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"
)
func familyFieldKeys() map[string][]*telemetrytypes.TelemetryFieldKey {
newKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
return map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {newKey("deployment.environment.name")},
"deployment.environment": {newKey("deployment.environment")},
}
}
func familyConditionSQL(t *testing.T, op qbtypes.FilterOperator, value any) (string, []any) {
t.Helper()
cb := NewConditionBuilder(NewFieldMapper())
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := cb.ConditionFor(context.Background(), valuer.UUID{}, 0, 0,
&telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"},
familyFieldKeys(), qbtypes.ConditionBuilderOptions{}, op, value, sb)
require.NoError(t, err)
require.Len(t, conds, 1)
sb.Where(conds...)
return sb.BuildWithFlavor(sqlbuilder.ClickHouse)
}
const familyValueExpr = "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')"
func TestFamilyEqualWidensIndexHintsToAnyMember(t *testing.T) {
sql, args := familyConditionSQL(t, qbtypes.FilterOperatorEqual, "production")
assert.Contains(t, sql, familyValueExpr+" = ?")
// key hint: either member name may appear in the labels JSON
assert.Contains(t, sql, "(labels LIKE ? OR labels LIKE ?)")
assert.Contains(t, args, "%deployment.environment.name%")
assert.Contains(t, args, "%deployment.environment%")
assert.Contains(t, args, `%deployment.environment.name":"production%`)
assert.Contains(t, args, `%deployment.environment":"production%`)
}
func TestFamilyNotEqualDropsNegatedValueHint(t *testing.T) {
sql, args := familyConditionSQL(t, qbtypes.FilterOperatorNotEqual, "production")
assert.Contains(t, sql, familyValueExpr+" <> ?")
// A negated per-member value hint would drop rows where the other member
// holds the value, so the family form carries no index hints at all.
assert.NotContains(t, sql, "NOT LIKE")
assert.Equal(t, []any{"production"}, args)
}
func TestFamilyExistsIsAnyMemberPresence(t *testing.T) {
sql, _ := familyConditionSQL(t, qbtypes.FilterOperatorExists, nil)
assert.Contains(t, sql, "(simpleJSONHas(labels, 'deployment.environment.name') = ? OR simpleJSONHas(labels, 'deployment.environment') = ?)")
sql, _ = familyConditionSQL(t, qbtypes.FilterOperatorNotExists, nil)
assert.Contains(t, sql, "(simpleJSONHas(labels, 'deployment.environment.name') <> ? AND simpleJSONHas(labels, 'deployment.environment') <> ?)")
}
// With only one member in metadata the SQL keeps the exact pre-family shape,
// including the negated value hint on !=.
func TestSingleMemberShapesUnchanged(t *testing.T) {
cb := NewConditionBuilder(NewFieldMapper())
soloKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": familyFieldKeys()["deployment.environment.name"],
}
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := cb.ConditionFor(context.Background(), valuer.UUID{}, 0, 0,
&telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"},
soloKeys, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotEqual, "production", sb)
require.NoError(t, err)
sb.Where(conds...)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "simpleJSONExtractString(labels, 'deployment.environment.name') <> ?")
assert.Contains(t, sql, "labels NOT LIKE ?")
assert.NotContains(t, sql, "COALESCE")
}

View File

@@ -71,6 +71,33 @@ func (m *defaultFieldMapper) FieldFor(
return columns[0].Name, nil
}
// ExistsFor reports key presence in the fingerprint labels JSON. Only resource
// context keys have a presence notion here; anything else is a real column and
// always present.
func (m *defaultFieldMapper) ExistsFor(
ctx context.Context,
_ valuer.UUID,
tsStart, tsEnd uint64,
key *telemetrytypes.TelemetryFieldKey,
exists bool,
) (string, error) {
columns, err := m.getColumn(ctx, tsStart, tsEnd, key)
if err != nil {
return "", err
}
if key.FieldContext != telemetrytypes.FieldContextResource {
if exists {
return "true", nil
}
return "false", nil
}
pred := fmt.Sprintf("simpleJSONHas(%s, '%s')", columns[0].Name, key.Name)
if exists {
return pred, nil
}
return "NOT " + pred, nil
}
func (m *defaultFieldMapper) ColumnExpressionFor(
ctx context.Context,
orgID valuer.UUID,

View File

@@ -99,7 +99,7 @@ func (b *resourceFilterStatementBuilder[T]) Build(
q.Select("fingerprint")
q.From(fmt.Sprintf("%s.%s", b.dbName, b.tableName))
keySelectors := b.getKeySelectors(query)
keySelectors := querybuilder.ExpandKeySelectorsForFamilies(b.getKeySelectors(query))
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, err

View File

@@ -266,7 +266,7 @@ func (b *scopedTraceStatementBuilder) fetchKeys(ctx context.Context, orgID value
SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact,
})
}
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, selectors)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, querybuilder.ExpandKeySelectorsForFamilies(selectors))
return keys, err
}
@@ -442,7 +442,7 @@ func (b *scopedTraceStatementBuilder) resolveSpanPredicate(ctx context.Context,
for i := range selectors {
selectors[i].Signal = telemetrytypes.SignalTraces
}
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, selectors)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, querybuilder.ExpandKeySelectorsForFamilies(selectors))
if err != nil {
return "", nil, "", err
}

View File

@@ -119,7 +119,7 @@ func (b *traceQueryStatementBuilder) Build(
// We modify SelectFields above (injecting default fields), and those default
// fields can carry keys that need evolutions, so fetch keys after that.
keySelectors := getKeySelectors(query)
keySelectors := querybuilder.ExpandKeySelectorsForFamilies(getKeySelectors(query))
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {

View File

@@ -75,6 +75,27 @@ func TestStatementBuilder(t *testing.T) {
},
expectedErr: nil,
},
{
name: "family filter merges both spellings",
requestType: qbtypes.RequestTypeScalar,
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'",
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_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_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 <= ? ORDER BY __result_0 DESC",
Args: []any{"production", "%deployment.environment.name%", "%deployment.environment%", "%deployment.environment.name\":\"production%", "%deployment.environment\":\"production%", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
expectedErr: nil,
},
{
name: "OR b/w resource attr and attribute",
requestType: qbtypes.RequestTypeTimeSeries,
@@ -373,94 +394,6 @@ func TestStatementBuilder(t *testing.T) {
},
expectedErr: nil,
},
{
name: "scope.name filter and group by",
requestType: qbtypes.RequestTypeTimeSeries,
query: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.TraceAggregation{
{
Expression: "count()",
},
},
Filter: &qbtypes.Filter{
Expression: "scope.name = 'opentelemetry-io'",
},
Limit: 10,
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "scope.name",
FieldContext: telemetrytypes.FieldContextScope,
},
},
},
},
expected: qbtypes.Statement{
Query: "WITH __limit_cte AS (SELECT toString(multiIf(scope.name::String IS NOT NULL, scope.name::String, NULL)) AS `scope.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.name::String = ? AND scope.name::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `scope.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(scope.name::String IS NOT NULL, scope.name::String, NULL)) AS `scope.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.name::String = ? AND scope.name::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`scope.name`) GLOBAL IN (SELECT `scope.name` FROM __limit_cte) GROUP BY ts, `scope.name`",
Args: []any{"opentelemetry-io", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10, "opentelemetry-io", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
expectedErr: nil,
},
{
name: "scope.version filter with scope.name group by",
requestType: qbtypes.RequestTypeTimeSeries,
query: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.TraceAggregation{
{
Expression: "count()",
},
},
Filter: &qbtypes.Filter{
Expression: "scope.version = '1.0.0'",
},
Limit: 10,
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "scope.name",
FieldContext: telemetrytypes.FieldContextScope,
},
},
},
},
expected: qbtypes.Statement{
Query: "WITH __limit_cte AS (SELECT toString(multiIf(scope.name::String IS NOT NULL, scope.name::String, NULL)) AS `scope.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.version::String = ? AND scope.version::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `scope.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(scope.name::String IS NOT NULL, scope.name::String, NULL)) AS `scope.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.version::String = ? AND scope.version::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`scope.name`) GLOBAL IN (SELECT `scope.name` FROM __limit_cte) GROUP BY ts, `scope.name`",
Args: []any{"1.0.0", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10, "1.0.0", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
expectedErr: nil,
},
{
name: "scope.version filter only (no scope field in group by)",
requestType: qbtypes.RequestTypeTimeSeries,
query: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.TraceAggregation{
{
Expression: "count()",
},
},
Filter: &qbtypes.Filter{
Expression: "scope.version = '1.0.0'",
},
Limit: 10,
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "service.name",
},
},
},
},
expected: qbtypes.Statement{
Query: "WITH __limit_cte AS (SELECT toString(multiIf(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) IS NOT NULL, multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL), NULL)) AS `service.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.version::String = ? AND scope.version::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `service.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) IS NOT NULL, multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL), NULL)) AS `service.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.version::String = ? AND scope.version::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`service.name`) GLOBAL IN (SELECT `service.name` FROM __limit_cte) GROUP BY ts, `service.name`",
Args: []any{"1.0.0", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10, "1.0.0", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
},
}
fl := flaggertest.New(t)
@@ -887,52 +820,6 @@ func TestStatementBuilderListQueryWithCorruptData(t *testing.T) {
},
expectedErr: nil,
},
{
name: "List query with scope filter only (no scope in select or group by)",
requestType: qbtypes.RequestTypeRaw,
keysMap: map[string][]*telemetrytypes.TelemetryFieldKey{
"scope.version": {
{
Name: "scope.version",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextScope,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
},
},
query: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Filter: &qbtypes.Filter{
Expression: "scope.version = '1.0.0'",
},
Limit: 10,
},
expected: qbtypes.Statement{
Query: "SELECT timestamp AS `timestamp`, trace_id AS `trace_id`, span_id AS `span_id`, trace_state AS `trace_state`, parent_span_id AS `parent_span_id`, flags AS `flags`, name AS `name`, kind AS `kind`, kind_string AS `kind_string`, duration_nano AS `duration_nano`, status_code AS `status_code`, status_message AS `status_message`, status_code_string AS `status_code_string`, events AS `events`, links AS `links`, response_status_code AS `response_status_code`, external_http_url AS `external_http_url`, http_url AS `http_url`, external_http_method AS `external_http_method`, http_method AS `http_method`, http_host AS `http_host`, db_name AS `db_name`, db_operation AS `db_operation`, has_error AS `has_error`, is_remote AS `is_remote`, attributes_string, attributes_number, attributes_bool, resources_string FROM signoz_traces.distributed_signoz_index_v3 WHERE (scope.version::String = ? AND scope.version::String IS NOT NULL) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"1.0.0", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
},
{
// Regression test: scope.version in selectFields with no metadata (isColumn=true filters it out)
// must still produce scope.version::String, not scope.attributes.version::String
name: "scope.version in selectFields only, no metadata (intrinsic field fallback)",
requestType: qbtypes.RequestTypeRaw,
keysMap: map[string][]*telemetrytypes.TelemetryFieldKey{},
query: qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Filter: &qbtypes.Filter{},
SelectFields: []telemetrytypes.TelemetryFieldKey{
{Name: "scope.version", FieldContext: telemetrytypes.FieldContextUnspecified},
},
Limit: 10,
},
expected: qbtypes.Statement{
Query: "SELECT timestamp AS `timestamp`, trace_id AS `trace_id`, span_id AS `span_id`, scope.version::String AS `scope.version` FROM signoz_traces.distributed_signoz_index_v3 WHERE timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
},
}
for _, c := range cases {

View File

@@ -212,7 +212,7 @@ func (b *traceOperatorCTEBuilder) buildQueryCTE(ctx context.Context, queryName s
return cteName, nil
}
keySelectors := getKeySelectors(*query)
keySelectors := querybuilder.ExpandKeySelectorsForFamilies(getKeySelectors(*query))
b.stmtBuilder.logger.DebugContext(ctx, "Key selectors for query", slog.String("query_name", queryName), slog.Any("key_selectors", keySelectors))
keys, _, err := b.stmtBuilder.metadataStore.GetKeysMulti(ctx, b.orgID, keySelectors)
if err != nil {
@@ -442,7 +442,7 @@ func (b *traceOperatorCTEBuilder) buildFinalQuery(ctx context.Context, selectFro
}
}
keySelectors := b.getKeySelectors()
keySelectors := querybuilder.ExpandKeySelectorsForFamilies(b.getKeySelectors())
keys, _, err := b.stmtBuilder.metadataStore.GetKeysMulti(ctx, b.orgID, keySelectors)
if err != nil {
return nil, err

View File

@@ -38,8 +38,11 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
// an unknown key simply yields no condition rather than an error.
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
// an unknown key simply yields no condition rather than an error. Metadata
// fields have no family support, so every logical field is single-member
// and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -58,6 +58,19 @@ func (m *fieldMapper) ColumnFor(ctx context.Context, _ valuer.UUID, tsStart, tsE
return columns, nil
}
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
func (m *fieldMapper) ExistsFor(ctx context.Context, _ valuer.UUID, tsStart, tsEnd uint64, key *telemetrytypes.TelemetryFieldKey, exists bool) (string, error) {
columns, err := m.getColumn(ctx, tsStart, tsEnd, key)
if err != nil {
return "", err
}
pred := fmt.Sprintf("mapContains(%s, '%s')", columns[0].Name, key.Name)
if exists {
return pred, nil
}
return "NOT " + pred, nil
}
func (m *fieldMapper) FieldFor(ctx context.Context, _ valuer.UUID, startNs, endNs uint64, key *telemetrytypes.TelemetryFieldKey) (string, error) {
columns, err := m.getColumn(ctx, startNs, endNs, key)
if err != nil {

View File

@@ -180,7 +180,7 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
`CASE
// WHEN tagType = 'spanfield' THEN 1
WHEN tagType = 'resource' THEN 2
WHEN tagType = 'scope' THEN 3
// WHEN tagType = 'scope' THEN 3
WHEN tagType = 'tag' THEN 4
ELSE 5
END as priority`,

View File

@@ -139,7 +139,10 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
// Audit fields have no family support, so every logical field is
// single-member and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -97,6 +97,19 @@ func (m *fieldMapper) ColumnFor(ctx context.Context, _ valuer.UUID, _, _ uint64,
return m.getColumn(ctx, key)
}
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
func (m *fieldMapper) ExistsFor(ctx context.Context, orgID valuer.UUID, tsStart, tsEnd uint64, key *telemetrytypes.TelemetryFieldKey, exists bool) (string, error) {
fieldExpression, err := m.FieldFor(ctx, orgID, tsStart, tsEnd, key)
if err != nil {
return "", err
}
columns, err := m.getColumn(ctx, key)
if err != nil {
return "", err
}
return querybuilder.ExistsExpression(columns, key, tsStart, tsEnd, fieldExpression, exists)
}
func (m *fieldMapper) ColumnExpressionFor(
ctx context.Context,
orgID valuer.UUID,

View File

@@ -452,7 +452,7 @@ func (c *conditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
// search() resolves its own (optional) scope; handle it before key resolution.
@@ -460,7 +460,10 @@ func (c *conditionBuilder) ConditionFor(
return c.conditionForSearch(ctx, orgID, key, value, sb)
}
keys, warning := querybuilder.ResolveKeys(key, matches)
// Logs fields have no family support yet, so every logical field is
// single-member and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, matches)
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -279,7 +279,7 @@ func (m *fieldMapper) ColumnExpressionFor(
}
var stmts []string
for _, key := range candidates {
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, key, true)
guard, err := m.ExistsFor(ctx, orgID, tsStart, tsEnd, key, true)
if err != nil {
return "", err
}
@@ -308,7 +308,7 @@ func (m *fieldMapper) ColumnExpressionFor(
if !m.membershipGuarded(ctx, orgID, tsStart, tsEnd, candidates[0]) {
return m.FieldFor(ctx, orgID, tsStart, tsEnd, candidates[0])
}
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, candidates[0], true)
guard, err := m.ExistsFor(ctx, orgID, tsStart, tsEnd, candidates[0], true)
if err != nil {
return "", err
}
@@ -326,7 +326,7 @@ func (m *fieldMapper) ColumnExpressionFor(
var stmts []string
for _, key := range candidates {
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, key, true)
guard, err := m.ExistsFor(ctx, orgID, tsStart, tsEnd, key, true)
if err != nil {
return "", err
}
@@ -569,7 +569,8 @@ func (m *fieldMapper) membershipGuarded(ctx context.Context, orgID valuer.UUID,
return columnType == schema.ColumnTypeEnumMap || columnType == schema.ColumnTypeEnumJSON
}
func (m *fieldMapper) existsExpressionFor(
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
func (m *fieldMapper) ExistsFor(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,

View File

@@ -162,7 +162,9 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
keys := querybuilder.MatchingFieldKeys(key, fieldKeys)
// Metric labels have no family support, so every logical field is
// single-member and flattens losslessly to its physical key.
keys := querybuilder.SingleKeys(querybuilder.MatchingLogicalFields(key, fieldKeys))
var warnings []string
if len(keys) == 0 {
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {

View File

@@ -97,6 +97,18 @@ func (m *fieldMapper) ColumnFor(ctx context.Context, _ valuer.UUID, tsStart, tsE
return m.getColumn(ctx, tsStart, tsEnd, key)
}
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
// Intrinsic fields always exist; labels are checked for key membership.
func (m *fieldMapper) ExistsFor(_ context.Context, _ valuer.UUID, _, _ uint64, key *telemetrytypes.TelemetryFieldKey, exists bool) (string, error) {
if slices.Contains(IntrinsicFields, key.Name) {
return "true", nil
}
if exists {
return fmt.Sprintf("has(JSONExtractKeys(labels), '%s')", key.Name), nil
}
return fmt.Sprintf("not has(JSONExtractKeys(labels), '%s')", key.Name), nil
}
func (m *fieldMapper) ColumnExpressionFor(
ctx context.Context,
orgID valuer.UUID,

View File

@@ -32,7 +32,7 @@ func (c *conditionBuilder) conditionFor(
orgID valuer.UUID,
startNs uint64,
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
logical *telemetrytypes.LogicalField,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -42,13 +42,13 @@ func (c *conditionBuilder) conditionFor(
value = querybuilder.FormatValueForContains(value)
}
fieldExpression, err := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
fieldExpression, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, c.fm, logical)
if err != nil {
return "", err
}
// TODO(srikanthccv): maybe extend this to every possible attribute
if key.Name == "duration_nano" || key.Name == "durationNano" { // QoL improvement
if logical.Name == "duration_nano" || logical.Name == "durationNano" { // QoL improvement
switch v := value.(type) {
case string:
if duration, err := time.ParseDuration(v); err == nil {
@@ -65,7 +65,7 @@ func (c *conditionBuilder) conditionFor(
}
}
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(key, value, fieldExpression, operator)
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), value, fieldExpression, operator)
// regular operators
switch operator {
@@ -154,11 +154,7 @@ func (c *conditionBuilder) conditionFor(
// in the query builder, `exists` and `not exists` are used for
// key membership checks, so depending on the column type, the condition changes
case qbtypes.FilterOperatorExists, qbtypes.FilterOperatorNotExists:
columns, err := c.fm.ColumnFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err
}
pred, err := querybuilder.ExistsExpression(columns, key, startNs, endNs, fieldExpression, operator == qbtypes.FilterOperatorExists)
pred, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, c.fm, logical, operator == qbtypes.FilterOperatorExists)
if err != nil {
return "", err
}
@@ -210,10 +206,10 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
keys, warning := querybuilder.ResolveKeys(key, matches)
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -221,10 +217,10 @@ func (c *conditionBuilder) ConditionFor(
// A bare key that names a real column filters on the column too — first. When metadata
// only knows the name under other contexts, prepend the column and keep metadata matches
// only where their type is consistent with it (a corrupt entry can't degrade the column).
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(keys) > 0 {
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(logicalFields) > 0 {
hasColumn := false
for _, k := range keys {
if k.FieldContext == telemetrytypes.FieldContextSpan {
for _, logical := range logicalFields {
if logical.FieldContext == telemetrytypes.FieldContextSpan {
hasColumn = true
break
}
@@ -232,49 +228,49 @@ func (c *conditionBuilder) ConditionFor(
if !hasColumn {
probe := telemetrytypes.NewTelemetryFieldKey(key.Name, telemetrytypes.FieldContextSpan, key.FieldDataType)
if cols, colErr := c.fm.ColumnFor(ctx, orgID, startNs, endNs, probe); colErr == nil && len(cols) > 0 {
combined := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys)+1)
combined = append(combined, probe)
for _, k := range keys {
if columnMatchesDataType(cols[0], k.FieldDataType) {
combined = append(combined, k)
combined := make([]*telemetrytypes.LogicalField, 0, len(logicalFields)+1)
combined = append(combined, telemetrytypes.SingleLogicalField(key.Name, probe))
for _, logical := range logicalFields {
if columnMatchesDataType(cols[0], logical.FieldDataType) {
combined = append(combined, logical)
}
}
keys = combined
logicalFields = combined
}
}
}
synthesized := false
if len(keys) == 0 {
if len(logicalFields) == 0 {
// Not in metadata. CandidateKeys resolves it: fold contexts (span/trace) get the
// metadata map so it can honor a real column, correct to a stripped-name metadata
// match, or synthesize; strict contexts pass nil and keep their synthesize path.
keys = c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys))
if len(keys) == 0 {
logicalFields = querybuilder.WrapAsLogicalFields(key.Name, c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys)))
if len(logicalFields) == 0 {
return nil, warnings, querybuilder.NewKeyNotFoundError(key.Name)
}
synthesized = true
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
// When a resource sub-query already covers the term, drop resource keys from the main
// When a resource sub-query already covers the term, drop resource fields from the main
// query. Synthesized keys are exempt: the sub-query skips keys absent from metadata.
if skipResourceFilter && !synthesized {
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
for _, k := range keys {
if k.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, k)
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
for _, logical := range logicalFields {
if logical.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, logical)
}
}
if len(filtered) == 0 {
return nil, warnings, nil
}
keys = filtered
logicalFields = filtered
}
conds := make([]string, 0, len(keys))
for _, k := range keys {
cond, err := c.conditionForKey(ctx, orgID, startNs, endNs, k, operator, value, sb)
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
cond, err := c.conditionForLogicalField(ctx, orgID, startNs, endNs, logical, operator, value, sb)
if err != nil {
return nil, nil, err
}
@@ -283,28 +279,28 @@ func (c *conditionBuilder) ConditionFor(
return conds, warnings, nil
}
func (c *conditionBuilder) conditionForKey(
func (c *conditionBuilder) conditionForLogicalField(
ctx context.Context,
orgID valuer.UUID,
startNs uint64,
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
logical *telemetrytypes.LogicalField,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
) (string, error) {
if c.isSpanScopeField(key.Name) {
return c.buildSpanScopeCondition(key, operator, value, startNs)
if c.isSpanScopeField(logical.Name) {
return c.buildSpanScopeCondition(logical.Single(), operator, value, startNs)
}
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, operator, value, sb)
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, operator, value, sb)
if err != nil {
return "", err
}
if operator.AddDefaultExistsFilter() {
// skip adding exists filter for intrinsic fields
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, logical.Single())
if slices.Contains(maps.Keys(IntrinsicFields), field) ||
slices.Contains(maps.Keys(IntrinsicFieldsDeprecated), field) ||
slices.Contains(maps.Keys(CalculatedFields), field) ||
@@ -312,7 +308,7 @@ func (c *conditionBuilder) conditionForKey(
return condition, nil
}
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, qbtypes.FilterOperatorExists, nil, sb)
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, qbtypes.FilterOperatorExists, nil, sb)
if err != nil {
return "", err
}

View File

@@ -121,20 +121,6 @@ var (
FieldContext: telemetrytypes.FieldContextSpan,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
"scope.name": {
Name: "scope.name",
Description: "Instrumentation scope name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextScope,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
"scope.version": {
Name: "scope.version",
Description: "Instrumentation scope version",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextScope,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
}
IntrinsicFieldsDeprecated = map[string]telemetrytypes.TelemetryFieldKey{
"traceID": {

View File

@@ -0,0 +1,148 @@
package tracestelemetryschema
import (
"context"
"fmt"
"testing"
"time"
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"
)
// familyFixture returns the deployment.environment(.name) family keys as trace
// resource attributes (with the canonical evolution timeline), a metadata map
// holding them, and a time range inside the JSON-column window.
func familyFixture() (current, old *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey, startNs, endNs uint64) {
releaseTime := time.Date(2025, 5, 22, 22, 0, 0, 0, time.UTC)
newKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
Evolutions: MockEvolutionData(releaseTime),
}
}
current = newKey("deployment.environment.name")
old = newKey("deployment.environment")
fieldKeys = map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
}
return current, old, fieldKeys, uint64(1747947419000000000), uint64(1747983448000000000)
}
// memberValueExprs returns each member's own FieldFor output; family
// expressions must be exactly the composition of these.
func memberValueExprs(t *testing.T, fm qbtypes.FieldMapper, startNs, endNs uint64, members ...*telemetrytypes.TelemetryFieldKey) []string {
t.Helper()
exprs := make([]string, 0, len(members))
for _, member := range members {
expr, err := fm.FieldFor(context.Background(), valuer.UUID{}, startNs, endNs, member)
require.NoError(t, err)
exprs = append(exprs, expr)
}
return exprs
}
func TestConditionForFamilyMergesMembersCurrentFirst(t *testing.T) {
current, old, fieldKeys, startNs, endNs := familyFixture()
fm := NewFieldMapper()
cb := NewConditionBuilder(fm)
// The requested spelling is the old name; precedence must still be
// current-first.
requested := &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment"}
sb := sqlbuilder.NewSelectBuilder()
conds, warnings, err := cb.ConditionFor(context.Background(), valuer.UUID{}, startNs, endNs, requested, fieldKeys, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb)
require.NoError(t, err)
assert.Empty(t, warnings, "a family is one logical field, never ambiguous with itself")
require.Len(t, conds, 1)
exprs := memberValueExprs(t, fm, startNs, endNs, current, old)
family := fmt.Sprintf("COALESCE(NULLIF(%s, ''), NULLIF(%s, ''), '')", exprs[0], exprs[1])
sb.Where(conds...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, family+" = ?")
// Equal adds the default exists filter: presence of any member.
assert.Contains(t, sql, fmt.Sprintf("(%s IS NOT NULL OR %s IS NOT NULL)", exprs[0], exprs[1]))
assert.Equal(t, []any{"production"}, args)
}
func TestConditionForFamilyNegativeKeepsKeylessRows(t *testing.T) {
current, old, fieldKeys, startNs, endNs := familyFixture()
fm := NewFieldMapper()
cb := NewConditionBuilder(fm)
requested := &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := cb.ConditionFor(context.Background(), valuer.UUID{}, startNs, endNs, requested, fieldKeys, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotEqual, "production", sb)
require.NoError(t, err)
require.Len(t, conds, 1)
exprs := memberValueExprs(t, fm, startNs, endNs, current, old)
sb.Where(conds...)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
// The trailing '' makes rows without any member read '' (single-key map
// semantics), so `!=` keeps including them; no exists filter is added.
assert.Contains(t, sql, fmt.Sprintf("COALESCE(NULLIF(%s, ''), NULLIF(%s, ''), '') <> ?", exprs[0], exprs[1]))
assert.NotContains(t, sql, "IS NOT NULL OR")
}
func TestConditionForFamilyExists(t *testing.T) {
current, old, fieldKeys, startNs, endNs := familyFixture()
fm := NewFieldMapper()
cb := NewConditionBuilder(fm)
requested := &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := cb.ConditionFor(context.Background(), valuer.UUID{}, startNs, endNs, requested, fieldKeys, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotExists, nil, sb)
require.NoError(t, err)
require.Len(t, conds, 1)
exprs := memberValueExprs(t, fm, startNs, endNs, current, old)
sb.Where(conds...)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, fmt.Sprintf("NOT (%s IS NOT NULL OR %s IS NOT NULL)", exprs[0], exprs[1]))
}
// A single-member key's condition is byte-identical to the pre-family shape:
// composition only appears when metadata proves a second member.
func TestConditionForSingleMemberIsUnchanged(t *testing.T) {
current, _, _, startNs, endNs := familyFixture()
fm := NewFieldMapper()
cb := NewConditionBuilder(fm)
soloKeys := map[string][]*telemetrytypes.TelemetryFieldKey{current.Name: {current}}
requested := &telemetrytypes.TelemetryFieldKey{Name: current.Name}
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := cb.ConditionFor(context.Background(), valuer.UUID{}, startNs, endNs, requested, soloKeys, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb)
require.NoError(t, err)
require.Len(t, conds, 1)
exprs := memberValueExprs(t, fm, startNs, endNs, current)
sb.Where(conds...)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, exprs[0]+" = ?")
assert.NotContains(t, sql, "COALESCE")
}
func TestColumnExpressionForFamilyGroupBy(t *testing.T) {
current, old, fieldKeys, startNs, endNs := familyFixture()
fm := NewFieldMapper()
expr, err := fm.ColumnExpressionFor(context.Background(), valuer.UUID{}, startNs, endNs,
&telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, telemetrytypes.FieldDataTypeString, fieldKeys)
require.NoError(t, err)
exprs := memberValueExprs(t, fm, startNs, endNs, current, old)
family := fmt.Sprintf("COALESCE(NULLIF(%s, ''), NULLIF(%s, ''), '')", exprs[0], exprs[1])
guard := fmt.Sprintf("(%s IS NOT NULL OR %s IS NOT NULL)", exprs[0], exprs[1])
assert.Equal(t, fmt.Sprintf("multiIf(%s, %s, NULL)", guard, family), expr)
}

View File

@@ -52,7 +52,6 @@ var (
ValueType: schema.ColumnTypeString,
}},
"resource": {Name: "resource", Type: schema.JSONColumnType{}},
"scope": {Name: "scope", Type: schema.JSONColumnType{}},
"events": {Name: "events", Type: schema.ArrayColumnType{
ElementType: schema.ColumnTypeString,
@@ -177,7 +176,7 @@ func (m *fieldMapper) getColumn(
case telemetrytypes.FieldContextResource:
return []*schema.Column{indexV3Columns["resource"], indexV3Columns["resources_string"]}, nil
case telemetrytypes.FieldContextScope:
return []*schema.Column{indexV3Columns["scope"]}, nil
return []*schema.Column{}, qbtypes.ErrColumnNotFound
case telemetrytypes.FieldContextAttribute:
switch key.FieldDataType {
case telemetrytypes.FieldDataTypeString:
@@ -288,24 +287,14 @@ func (m *fieldMapper) resolveColumnExprs(
switch column.Type.GetType() {
case schema.ColumnTypeEnumJSON:
// json is only supported for resource context as of now
if key.FieldContext != telemetrytypes.FieldContextResource {
return nil, nil, nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "only resource context fields are supported for json columns, got %s", key.FieldContext.String)
}
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
// once clickHouse dependency is updated, we need to check if we can remove it.
switch key.FieldContext {
case telemetrytypes.FieldContextResource:
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
case telemetrytypes.FieldContextScope:
switch key.Name {
case "scope.name", "scope.version":
exprs = append(exprs, fmt.Sprintf("%s::String", key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s IS NOT NULL", key.Name))
default:
exprs = append(exprs, fmt.Sprintf("%s.attributes.`%s`::String", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s.attributes.`%s` IS NOT NULL", columnName, key.Name))
}
default:
return nil, nil, nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "only resource and scope context fields are supported for json columns, got %s", key.FieldContext.String)
}
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
case schema.ColumnTypeEnumString,
schema.ColumnTypeEnumUInt64,
schema.ColumnTypeEnumUInt32,
@@ -347,6 +336,68 @@ func (m *fieldMapper) resolveColumnExprs(
return exprs, existExprs, columns, nil
}
// logicalForResolvedColumn upgrades a directly-resolvable key (the FieldFor
// probe succeeded) to its family when the metadata map proves membership;
// otherwise the key stays a single-member logical field.
func logicalForResolvedColumn(field *telemetrytypes.TelemetryFieldKey, keys map[string][]*telemetrytypes.TelemetryFieldKey) *telemetrytypes.LogicalField {
for _, logical := range querybuilder.MatchingLogicalFields(field, keys) {
if logical.IsFamily() &&
logical.FieldContext == field.FieldContext &&
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || logical.FieldDataType == field.FieldDataType) {
return logical
}
}
return telemetrytypes.SingleLogicalField(field.Name, field)
}
// upgradeToFamilies swaps single-member candidates for their family when the
// metadata map proves membership. Candidate order and every non-family
// candidate stay exactly as the legacy flow produced them; sibling candidates
// of an already-emitted family are dropped rather than duplicated.
func upgradeToFamilies(field *telemetrytypes.TelemetryFieldKey, candidates []*telemetrytypes.LogicalField, keys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
var families []*telemetrytypes.LogicalField
for _, logical := range querybuilder.MatchingLogicalFields(field, keys) {
if logical.IsFamily() {
families = append(families, logical)
}
}
if len(families) == 0 {
return candidates
}
out := make([]*telemetrytypes.LogicalField, 0, len(candidates))
emitted := make(map[*telemetrytypes.LogicalField]bool)
for _, candidate := range candidates {
var family *telemetrytypes.LogicalField
for _, fam := range families {
if fam.FieldContext != candidate.FieldContext || fam.FieldDataType != candidate.FieldDataType {
continue
}
memberOfFamily := candidate.Single().Name == field.Name
for _, member := range fam.Members {
if member.Name == candidate.Single().Name {
memberOfFamily = true
break
}
}
if memberOfFamily {
family = fam
break
}
}
if family == nil {
out = append(out, candidate)
continue
}
if emitted[family] {
continue
}
emitted[family] = true
out = append(out, family)
}
return out
}
// ColumnExpressionFor returns the bare (unaliased) SQL expression for the field, resolving
// unknown keys via CandidateKeys and wrapping guardable columns with exists-guard multiIfs
// so an absent key yields NULL.
@@ -359,18 +410,23 @@ func (m *fieldMapper) ColumnExpressionFor(
keys map[string][]*telemetrytypes.TelemetryFieldKey,
) (string, error) {
// Resolve the candidate column(s).
var candidates []*telemetrytypes.TelemetryFieldKey
// Resolve the candidate logical field(s).
var candidates []*telemetrytypes.LogicalField
switch _, err := m.FieldFor(ctx, orgID, startNs, endNs, field); {
case err == nil:
candidates = []*telemetrytypes.TelemetryFieldKey{field}
// A directly-resolvable key upgrades to its family when the metadata
// map proves membership; otherwise it stays single-member.
candidates = []*telemetrytypes.LogicalField{logicalForResolvedColumn(field, keys)}
case errors.Is(err, qbtypes.ErrColumnNotFound):
// column (when the bare name is one) plus metadata matches, else synthesized
// type-variant keys.
candidates = m.CandidateKeys(ctx, orgID, field, nil, keys)
if len(candidates) == 0 {
// The legacy candidate flow, unchanged: column (when the bare name is
// one) plus metadata matches, else synthesized type-variant keys. The
// family step below only swaps candidates for their family; it never
// changes candidate order or non-family behavior.
raw := m.CandidateKeys(ctx, orgID, field, nil, keys)
if len(raw) == 0 {
return "", errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "field `%s` not found", field.Name).WithSuggestions(errors.NewSuggestionsOnLevenshteinDistance(field.Name, errors.NounKeys, maps.Keys(keys))...)
}
candidates = upgradeToFamilies(field, querybuilder.WrapAsLogicalFields(field.Name, raw), keys)
default:
return "", err
}
@@ -384,21 +440,21 @@ func (m *fieldMapper) ColumnExpressionFor(
dummyValue = 0.0
}
stmts := make([]string, 0, len(candidates)*2)
for _, key := range candidates {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
for _, logical := range candidates {
value, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logical)
if err != nil {
return "", err
}
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
guard, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, m, logical, true)
if err != nil {
return "", err
}
coerced := value
// a time column keeps its native type; coercing it would yield seconds
if temporal, err := m.columnIsTemporal(ctx, startNs, endNs, key); err != nil {
if temporal, err := m.logicalIsTemporal(ctx, startNs, endNs, logical); err != nil {
return "", err
} else if !temporal {
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(key, dummyValue, value, qbtypes.FilterOperatorUnknown)
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), dummyValue, value, qbtypes.FilterOperatorUnknown)
}
stmts = append(stmts, guard, coerced)
}
@@ -406,13 +462,14 @@ func (m *fieldMapper) ColumnExpressionFor(
}
if len(candidates) == 1 {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, candidates[0])
logical := candidates[0]
value, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logical)
if err != nil {
return "", err
}
exprs, existExprs, _, _ := m.resolveColumnExprs(ctx, startNs, endNs, candidates[0])
if len(exprs) == 1 && len(existExprs) == 1 {
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, candidates[0], true)
exprs, existExprs, _, _ := m.resolveColumnExprs(ctx, startNs, endNs, logical.Single())
if !logical.IsFamily() && len(exprs) == 1 && len(existExprs) == 1 {
guard, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, m, logical, true)
if err != nil {
return "", err
}
@@ -424,12 +481,12 @@ func (m *fieldMapper) ColumnExpressionFor(
// Multiple candidates (collision / synth): multiIf picks the first that exists,
// stringified so branches share a type.
args := make([]string, 0, len(candidates))
for _, key := range candidates {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
for _, logical := range candidates {
value, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logical)
if err != nil {
return "", err
}
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
guard, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, m, logical, true)
if err != nil {
return "", err
}
@@ -438,6 +495,15 @@ func (m *fieldMapper) ColumnExpressionFor(
return fmt.Sprintf("multiIf(%s, NULL)", strings.Join(args, ", ")), nil
}
// logicalIsTemporal reports whether the logical field resolves to a single time
// column. A family is attribute-backed and never temporal.
func (m *fieldMapper) logicalIsTemporal(ctx context.Context, startNs, endNs uint64, logical *telemetrytypes.LogicalField) (bool, error) {
if logical.IsFamily() {
return false, nil
}
return m.columnIsTemporal(ctx, startNs, endNs, logical.Single())
}
// columnIsTemporal reports whether key resolves to a single time column, after evolution
// selection. Multiple columns mean an attribute-map union, which is never temporal.
func (m *fieldMapper) columnIsTemporal(ctx context.Context, startNs, endNs uint64, key *telemetrytypes.TelemetryFieldKey) (bool, error) {
@@ -533,7 +599,8 @@ func (m *fieldMapper) CandidateKeys(ctx context.Context, _ valuer.UUID, field *t
return nil
}
func (m *fieldMapper) existsExpressionFor(
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
func (m *fieldMapper) ExistsFor(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,

View File

@@ -83,33 +83,6 @@ func TestGetFieldKeyName(t *testing.T) {
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
expectedError: nil,
},
{
name: "Scope field - scope.name",
key: telemetrytypes.TelemetryFieldKey{
Name: "scope.name",
FieldContext: telemetrytypes.FieldContextScope,
},
expectedResult: "scope.name::String",
expectedError: nil,
},
{
name: "Scope field - scope.version",
key: telemetrytypes.TelemetryFieldKey{
Name: "scope.version",
FieldContext: telemetrytypes.FieldContextScope,
},
expectedResult: "scope.version::String",
expectedError: nil,
},
{
name: "Scope field - custom attribute",
key: telemetrytypes.TelemetryFieldKey{
Name: "custom.attr",
FieldContext: telemetrytypes.FieldContextScope,
},
expectedResult: "scope.attributes.`custom.attr`::String",
expectedError: nil,
},
{
// Query like `attribute.attribute_string:string` should resolve to `attributes_string['attribute_string']`.
name: "Attribute key whose name collides with contextual map column resolves as a map lookup",

View File

@@ -113,17 +113,18 @@ func BuildCompleteFieldKeyMap(releaseTime time.Time) map[string][]*telemetrytype
FieldDataType: telemetrytypes.FieldDataTypeBool,
},
},
"scope.name": {
// both spellings of an enabled semantic-convention family
"deployment.environment.name": {
{
Name: "scope.name",
FieldContext: telemetrytypes.FieldContextScope,
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
},
"scope.version": {
"deployment.environment": {
{
Name: "scope.version",
FieldContext: telemetrytypes.FieldContextScope,
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
},

View File

@@ -34,6 +34,11 @@ type FieldMapper interface {
// the name (or `{context}.{name}`) first, else synthesized type-variant keys for sources
// that support it, else nil (caller errors). value is the filter operand, nil otherwise.
CandidateKeys(ctx context.Context, orgID valuer.UUID, field *telemetrytypes.TelemetryFieldKey, value any, keys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey
// ExistsFor returns the existence predicate for a single physical key (negated when
// exists is false), self-contained and arg-free so it can guard column expressions.
// It is the per-member primitive querybuilder.LogicalExistsExpr and the numeric branch
// of querybuilder.LogicalValueExpr compose family expressions from.
ExistsFor(ctx context.Context, orgID valuer.UUID, tsStart, tsEnd uint64, key *telemetrytypes.TelemetryFieldKey, exists bool) (string, error)
}
// ConditionBuilder builds the conditions for a filter term. The builder owns key resolution:

View File

@@ -0,0 +1,69 @@
package telemetrytypes
import "strings"
// LogicalField is resolution output: one queryable field, addressed by the
// spelling the request used, backed by the physical member keys that store it.
//
// The resolver expresses ambiguity ("possibly different fields sharing a
// name") as a []*LogicalField — never inside one LogicalField. Within one
// LogicalField, members are alternate physical spellings of the same field
// (a semantic-convention family), ordered current-first; compilers merge
// them into one expression with current-wins precedence. Across the slice,
// compilers build one condition per LogicalField and combine per the
// operator, exactly as they combine ambiguous keys.
//
// Members always has at least one entry. A non-family field has exactly
// one. Members alias the metadata map entries and must not be mutated.
type LogicalField struct {
// Name is the requested spelling. It is the response identity: aliases,
// series labels, and warnings use it, so responses echo the request.
Name string
// The physical identity every member shares. Members with a different
// signal, field context, or data type belong to different logical
// fields by definition.
Signal Signal
FieldContext FieldContext
FieldDataType FieldDataType
// Members are the physical keys that store this field, ordered
// current-first. Each member carries its own physical facts
// (Materialized, Evolutions, JSONPlan, ...), so per-member accessors
// need no sibling information.
Members []*TelemetryFieldKey
}
// SingleLogicalField wraps one physical key as its own logical field.
func SingleLogicalField(name string, key *TelemetryFieldKey) *LogicalField {
return &LogicalField{
Name: name,
Signal: key.Signal,
FieldContext: key.FieldContext,
FieldDataType: key.FieldDataType,
Members: []*TelemetryFieldKey{key},
}
}
// Single returns the only member. It is the accessor for signals whose
// logical fields are always single-member (everything except traces today).
func (l *LogicalField) Single() *TelemetryFieldKey {
return l.Members[0]
}
// IsFamily reports whether the field has more than one physical member.
func (l *LogicalField) IsFamily() bool {
return len(l.Members) > 1
}
// String implements fmt.Stringer for warning messages.
func (l *LogicalField) String() string {
if len(l.Members) == 1 {
return l.Members[0].String()
}
names := make([]string, 0, len(l.Members))
for _, member := range l.Members {
names = append(names, member.Name)
}
return l.Name + "(" + l.FieldContext.StringValue() + ", " + l.FieldDataType.StringValue() + ", members: " + strings.Join(names, ", ") + ")"
}

View File

@@ -24,8 +24,6 @@ pytest_plugins = [
"fixtures.browser",
"fixtures.keycloak",
"fixtures.idp",
"fixtures.googleidp",
"fixtures.tls",
"fixtures.notification_channel",
"fixtures.maildev",
"fixtures.alerts",

View File

@@ -1,231 +0,0 @@
import functools
import time
from collections.abc import Callable
from http import HTTPStatus
from pathlib import Path
from urllib.parse import urlparse
import docker
import docker.errors
import pytest
import requests
from jwcrypto import jwk, jwt
from testcontainers.core.container import Network
from wiremock.resources.mappings import HttpMethods, Mapping, MappingRequest, MappingResponse
from wiremock.testing.testcontainer import WireMockContainer
from fixtures import reuse, types
from fixtures.logger import setup_logger
from fixtures.tls import CA_ID_LABEL, KEYSTORE_PASSWORD, ca_id, issue_server_keystore
logger = setup_logger(__name__)
# The google callback authn hardcodes Google's issuer, so the mock must be
# reachable as accounts.google.com over TLS from the signoz container: the
# wiremock container joins the network under that alias and serves HTTPS on 443
# with a certificate issued by the integration CA that signoz trusts.
ISSUER = "https://accounts.google.com"
ISSUER_HOST = "accounts.google.com"
GOOGLE_DOMAIN = "google.integration.test"
# One signing key for the whole session: the token and JWKS stubs are always
# installed together, so per-call keys would only add RSA keygen latency.
@functools.cache
def signing_key() -> jwk.JWK:
return jwk.JWK.generate(kty="RSA", size=2048, kid="googleidp-integration", use="sig", alg="RS256")
def perform_google_login(
signoz: types.SigNoz,
googleidp: types.TestContainerDocker,
get_session_context: Callable[[str], dict],
email: str,
) -> str:
"""Drive the google login flow for email and return the final redirect URL.
The authorize URL points at https://accounts.google.com (resolvable only
inside the docker network), so it is rewritten to the mock's host-mapped
port, mirroring how the oidc suite rewrites keycloak URLs.
"""
session_context = get_session_context(email)
assert len(session_context["orgs"]) == 1
assert len(session_context["orgs"][0]["authNSupport"]["callback"]) == 1
url = session_context["orgs"][0]["authNSupport"]["callback"][0]["url"]
assert url.startswith(f"{ISSUER}/")
parsed_url = urlparse(url)
authorize_url = googleidp.host_configs["8080"].get(f"{parsed_url.path}?{parsed_url.query}")
response = requests.get(authorize_url, allow_redirects=False, timeout=5)
assert response.status_code == HTTPStatus.FOUND
callback_url = response.headers["Location"]
assert "/api/v1/complete/google" in callback_url
response = requests.get(callback_url, allow_redirects=False, timeout=30)
assert response.status_code == HTTPStatus.SEE_OTHER
return response.headers["Location"]
def get_google_domain(signoz: types.SigNoz, admin_token: str) -> dict:
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/domains"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
return next(
(domain for domain in response.json()["data"] if domain["name"] == GOOGLE_DOMAIN),
None,
)
def google_oidc_mappings(email: str, name: str, hd: str, audience: str, email_verified: bool = True) -> list[Mapping]:
"""Wiremock mappings for one Google OIDC login: discovery, an auto-approving
authorize redirect, a token response with an RS256 id_token for the given
identity, and the JWKS the signoz container verifies it against."""
now = int(time.time())
token = jwt.JWT(
header={"alg": "RS256", "kid": signing_key()["kid"], "typ": "JWT"},
claims={
"iss": ISSUER,
"aud": audience,
"sub": f"google-oauth2|{email}",
"email": email,
"email_verified": email_verified,
"name": name,
"hd": hd,
"iat": now,
"exp": now + 3600,
},
)
token.make_signed_token(signing_key())
id_token = token.serialize()
return [
Mapping(
request=MappingRequest(method=HttpMethods.GET, url_path="/.well-known/openid-configuration"),
response=MappingResponse(
status=200,
json_body={
"issuer": ISSUER,
"authorization_endpoint": f"{ISSUER}/o/oauth2/v2/auth",
"token_endpoint": f"{ISSUER}/token",
"jwks_uri": f"{ISSUER}/jwks",
"response_types_supported": ["code"],
"subject_types_supported": ["public"],
"id_token_signing_alg_values_supported": ["RS256"],
"scopes_supported": ["openid", "email", "profile"],
"token_endpoint_auth_methods_supported": ["client_secret_basic", "client_secret_post"],
},
),
),
Mapping(
request=MappingRequest(method=HttpMethods.GET, url_path="/o/oauth2/v2/auth"),
response=MappingResponse(
status=302,
headers={
# Triple-stache: redirect_uri and state are URLs; handlebars
# would otherwise HTML-escape their special characters.
# request.query values arrive URL-decoded, so the state is
# re-encoded into the redirect exactly as google does.
"Location": "{{{request.query.redirect_uri}}}?code=integration-test-code&state={{{urlEncode request.query.state}}}",
},
transformers=["response-template"],
),
),
Mapping(
request=MappingRequest(method=HttpMethods.POST, url_path="/token"),
response=MappingResponse(
status=200,
json_body={
"access_token": "integration-test-access-token",
"token_type": "Bearer",
"expires_in": 3600,
"id_token": id_token,
},
),
),
Mapping(
request=MappingRequest(method=HttpMethods.GET, url_path="/jwks"),
response=MappingResponse(
status=200,
json_body={"keys": [signing_key().export_public(as_dict=True)]},
),
),
]
@pytest.fixture(name="googleidp", scope="package")
def googleidp( # pylint: disable=too-many-arguments,too-many-positional-arguments
network: Network,
tls: types.TLS,
tmpfs: Callable[[str], Path],
request: pytest.FixtureRequest,
pytestconfig: pytest.Config,
) -> types.TestContainerDocker:
"""Wiremock impersonating Google's OIDC provider. Stubs are installed per
test via make_http_mocks with google_oidc_mappings; port 8080 serves the
admin API and the authorize redirect to the test process."""
def create() -> types.TestContainerDocker:
keystore_path = issue_server_keystore(tls, tmpfs("googleidp-certs"), ISSUER_HOST)
container = WireMockContainer(image="wiremock/wiremock:2.35.1-1", secure=False)
container.with_volume_mapping(str(keystore_path.parent), "/certs", "ro")
container.with_network(network)
container.with_network_aliases(ISSUER_HOST)
container.with_kwargs(labels={CA_ID_LABEL: ca_id(tls)})
try:
container.start(f"--port 8080 --https-port 443 --https-keystore /certs/keystore.p12 --keystore-type PKCS12 --keystore-password {KEYSTORE_PASSWORD} --local-response-templating")
except Exception:
# Ryuk is disabled: a started-but-unready container would survive
# and keep squatting on the accounts.google.com alias, poisoning
# DNS for any replacement on the shared network.
container.stop()
raise
return types.TestContainerDocker(
id=container.get_wrapped_container().id,
host_configs={
"8080": types.TestContainerUrlConfig("http", container.get_container_host_ip(), container.get_exposed_port(8080)),
},
container_configs={
"443": types.TestContainerUrlConfig("https", ISSUER_HOST, 443),
},
)
def delete(container: types.TestContainerDocker) -> None:
client = docker.from_env()
try:
client.containers.get(container_id=container.id).stop()
client.containers.get(container_id=container.id).remove(v=True)
except docker.errors.NotFound:
logger.info("googleidp container %s already gone", container.id)
def restore(cache: dict) -> types.TestContainerDocker:
return types.TestContainerDocker.from_cache(cache)
def stale(container: types.TestContainerDocker) -> bool:
client = docker.from_env()
try:
labels = client.containers.get(container_id=container.id).attrs["Config"]["Labels"]
except docker.errors.NotFound:
return True
return labels.get(CA_ID_LABEL) != ca_id(tls)
return reuse.wrap(
request,
pytestconfig,
"googleidp",
lambda: types.TestContainerDocker(id="", host_configs={}, container_configs={}),
create,
delete,
restore,
stale=stale,
)

View File

@@ -999,8 +999,6 @@ def generate_traces_with_corrupt_metadata() -> list[Traces]:
"cloud.provider": "integration",
"cloud.account.id": "000",
"trace_id": "corrupt_data",
"scope_name": "corrupt_data",
"scope.scope.name": "corrupt_data",
},
attributes={
"net.transport": "IP.TCP",
@@ -1009,10 +1007,7 @@ def generate_traces_with_corrupt_metadata() -> list[Traces]:
"http.request.method": "POST",
"http.response.status_code": "200",
"timestamp": "corrupt_data",
"version": "1.0.0",
"scope.scope.version": "1.0.0",
},
scope={"name": "io.signoz.http.server", "version": "2.0.0"},
),
Traces(
timestamp=now - timedelta(seconds=3.5),
@@ -1032,24 +1027,12 @@ def generate_traces_with_corrupt_metadata() -> list[Traces]:
"cloud.provider": "integration",
"cloud.account.id": "000",
"timestamp": "corrupt_data",
"scope.attributes.name": "corrupt_data",
},
attributes={
"db.name": "integration",
"db.operation": "SELECT",
"db.statement": "SELECT * FROM integration",
"trace_d": "corrupt_data",
"scope.attributes.version": "corrupt_data",
},
scope={
"name": "io.opentelemetry.contrib.http",
"version": "1.0.0",
"attributes": {
"telemetry.sdk.language": "cpp",
"name": "not-the-real-name",
"version": "not-the-real-version",
"attributes": "literally-a-key-named-attributes",
},
},
),
Traces(
@@ -1070,15 +1053,12 @@ def generate_traces_with_corrupt_metadata() -> list[Traces]:
"cloud.provider": "integration",
"cloud.account.id": "000",
"duration_nano": "corrupt_data",
"scope.scope.attributes.version": "corrupt_data",
},
attributes={
"http.request.method": "PATCH",
"http.status_code": "404",
"id": "1",
"scope.scope.version": "corrupt_data",
},
scope={"name": "io.signoz.http.client", "version": "2.0.0"},
),
Traces(
timestamp=now - timedelta(seconds=1),
@@ -1097,7 +1077,6 @@ def generate_traces_with_corrupt_metadata() -> list[Traces]:
"host.name": "linux-001",
"cloud.provider": "integration",
"cloud.account.id": "001",
"scope.scope.version": "corrupt_data",
},
attributes={
"message.type": "SENT",
@@ -1105,10 +1084,7 @@ def generate_traces_with_corrupt_metadata() -> list[Traces]:
"messaging.message.id": "001",
"duration_nano": "corrupt_data",
"id": 1,
"scope": "corrupt_data",
"scope.attributes.name": "corrupt_data",
},
scope={"name": "io.signoz.messaging", "version": "3.0.0"},
),
]

View File

@@ -39,7 +39,6 @@ def wrap( # pylint: disable=too-many-arguments,too-many-positional-arguments
delete: Callable[[T], None],
restore: Callable[[dict], T],
rebuild: bool = False,
stale: Callable[[T], bool] | None = None,
) -> T:
"""
Wraps a resource creation and cleanup process with reuse and teardown options.
@@ -51,7 +50,6 @@ def wrap( # pylint: disable=too-many-arguments,too-many-positional-arguments
- delete: function to delete the resource
- restore: function to restore resource from cache
- rebuild: under --reuse, delete the cached resource and recreate it instead of restoring it
- stale: under --reuse, decides whether a restored resource is still usable; a stale resource is deleted and recreated
"""
resource = empty()
@@ -64,14 +62,8 @@ def wrap( # pylint: disable=too-many-arguments,too-many-positional-arguments
delete(restore(existing_resource))
pytestconfig.cache.set(key, None)
else:
restored = restore(existing_resource)
if stale is not None and stale(restored):
logger.info("Recreating stale %s(%s)", key, existing_resource)
delete(restored)
pytestconfig.cache.set(key, None)
else:
logger.info("Reusing existing %s(%s)", key, existing_resource)
return restored
logger.info("Reusing existing %s(%s)", key, existing_resource)
return restore(existing_resource)
if not teardown(request):
resource = create()
@@ -96,23 +88,15 @@ def wrap( # pylint: disable=too-many-arguments,too-many-positional-arguments
return
resource = restore(existing_resource)
logger.info(
"Removing %s",
resource.__log__() if hasattr(resource, "__log__") else resource,
)
delete(resource)
pytestconfig.cache.set(key, None)
return
# A run without --reuse owns only what it created this session: the
# cache key (and whatever a parked --reuse stack has under it) is left
# untouched.
logger.info(
"Removing %s",
resource.__log__() if hasattr(resource, "__log__") else resource,
)
delete(resource)
pytestconfig.cache.set(key, None)
request.addfinalizer(finalizer)
if reuse(request):

View File

@@ -13,7 +13,6 @@ from testcontainers.core.container import DockerContainer, Network
from fixtures import reuse, types
from fixtures.logger import setup_logger
from fixtures.tls import CA_CONTAINER_PATH, CA_ID_LABEL, ca_id
logger = setup_logger(__name__)
@@ -28,12 +27,10 @@ def create_signoz(
pytestconfig: pytest.Config,
cache_key: str = "signoz",
env_overrides: dict | None = None,
tls: types.TLS | None = None,
) -> types.SigNoz:
"""
Factory function for creating a SigNoz container.
Accepts optional env_overrides to customize the container environment, and
an optional integration CA (tls) to trust in addition to the system roots.
Accepts optional env_overrides to customize the container environment.
"""
def create() -> types.SigNoz:
@@ -118,13 +115,6 @@ def create_signoz(
"rw",
)
# The CA lands in the directory Go scans for system roots, so tests can
# stand in for real TLS hosts (e.g. the fake accounts.google.com) while
# the bundled roots keep working for everything else.
if tls:
container.with_volume_mapping(tls.ca_cert_path, CA_CONTAINER_PATH, "ro")
container.with_kwargs(labels={CA_ID_LABEL: ca_id(tls)})
container.start()
def ready(container: DockerContainer) -> None:
@@ -203,16 +193,6 @@ def create_signoz(
gateway=gateway,
)
def stale(container: types.SigNoz) -> bool:
if not tls:
return False
client = docker.from_env()
try:
labels = client.containers.get(container_id=container.self.id).attrs["Config"]["Labels"]
except docker.errors.NotFound:
return True
return labels.get(CA_ID_LABEL) != ca_id(tls)
return reuse.wrap(
request,
pytestconfig,
@@ -232,7 +212,6 @@ def create_signoz(
delete=delete,
restore=restore,
rebuild=pytestconfig.getoption("--rebuild"),
stale=stale,
)
@@ -243,7 +222,6 @@ def signoz( # pylint: disable=too-many-arguments,too-many-positional-arguments
gateway: types.TestContainerDocker,
sqlstore: types.TestContainerSQL,
clickhouse: types.TestContainerClickhouse,
tls: types.TLS,
request: pytest.FixtureRequest,
pytestconfig: pytest.Config,
) -> types.SigNoz:
@@ -255,5 +233,4 @@ def signoz( # pylint: disable=too-many-arguments,too-many-positional-arguments
clickhouse=clickhouse,
request=request,
pytestconfig=pytestconfig,
tls=tls,
)

143
tests/fixtures/tls.py vendored
View File

@@ -1,143 +0,0 @@
import contextlib
import datetime
import hashlib
import os
import uuid
from pathlib import Path
import pytest
from cryptography import x509
from cryptography.hazmat.primitives import hashes, serialization
from cryptography.hazmat.primitives.asymmetric import rsa
from cryptography.hazmat.primitives.serialization import pkcs12
from cryptography.x509.oid import NameOID
from fixtures import reuse, types
from fixtures.logger import setup_logger
logger = setup_logger(__name__)
# The integration CA is mounted into the directory Go scans for system roots,
# so signoz containers trust it in addition to the bundled Debian roots (not
# instead of them, which is what SSL_CERT_FILE would do) and mocks that must be
# reached over TLS under a real hostname (e.g. accounts.google.com) can serve
# certificates issued by it via issue_server_keystore.
CA_CONTAINER_PATH = "/etc/ssl/certs/signoz-integration-ca.pem"
KEYSTORE_PASSWORD = "password" # noqa: S105
# Containers chained to the CA carry its id as this label; comparing it against
# the current CA spots reused containers built against a rotated CA (or before
# the CA existed) so they can be recreated instead of failing TLS opaquely.
CA_ID_LABEL = "signoz.integration.ca"
def ca_id(tls: types.TLS) -> str:
return hashlib.sha256(Path(tls.ca_cert_path).read_bytes()).hexdigest()[:12]
@pytest.fixture(name="tls", scope="package")
def tls(
request: pytest.FixtureRequest,
pytestconfig: pytest.Config,
) -> types.TLS:
"""The integration CA. Server certificates for mocks are issued from it
with issue_server_keystore.
The CA cannot live in tmpfs: pytest wipes basetemp at every session start,
while reused containers (and the keystores issued for them in later
sessions) must keep chaining to the same CA. The pytest cache directory is
the cross-session store, like the reuse metadata itself."""
def create() -> types.TLS:
# Each CA gets a fresh directory: a run without --reuse must never
# rotate the files that a parked stack's containers bind-mount.
ca_dir = pytestconfig.cache.mkdir("tls") / uuid.uuid4().hex
ca_dir.mkdir()
now = datetime.datetime.now(datetime.UTC)
ca_key = rsa.generate_private_key(public_exponent=65537, key_size=2048)
ca_name = x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, "signoz-integration-ca")])
ca_cert = (
x509.CertificateBuilder()
.subject_name(ca_name)
.issuer_name(ca_name)
.public_key(ca_key.public_key())
.serial_number(x509.random_serial_number())
.not_valid_before(now - datetime.timedelta(days=1))
.not_valid_after(now + datetime.timedelta(days=3650))
.add_extension(x509.BasicConstraints(ca=True, path_length=None), critical=True)
.sign(ca_key, hashes.SHA256())
)
(ca_dir / "ca.pem").write_bytes(ca_cert.public_bytes(serialization.Encoding.PEM))
(ca_dir / "ca.key").write_bytes(
ca_key.private_bytes(
encoding=serialization.Encoding.PEM,
format=serialization.PrivateFormat.PKCS8,
encryption_algorithm=serialization.NoEncryption(),
)
)
return types.TLS(ca_cert_path=str(ca_dir / "ca.pem"), ca_key_path=str(ca_dir / "ca.key"))
def delete(tls: types.TLS) -> None:
for path in (tls.ca_cert_path, tls.ca_key_path):
try:
os.remove(path)
except FileNotFoundError:
logger.info("CA file %s already gone", path)
with contextlib.suppress(OSError):
os.rmdir(Path(tls.ca_cert_path).parent)
def restore(cache: dict) -> types.TLS:
return types.TLS.from_cache(cache)
def stale(tls: types.TLS) -> bool:
return not (Path(tls.ca_cert_path).is_file() and Path(tls.ca_key_path).is_file())
return reuse.wrap(
request,
pytestconfig,
"tls",
lambda: types.TLS(ca_cert_path="", ca_key_path=""),
create,
delete,
restore,
stale=stale,
)
def issue_server_keystore(tls: types.TLS, directory: Path, hostname: str) -> Path:
"""Write a PKCS12 keystore (keystore.p12, password KEYSTORE_PASSWORD) into
directory, holding a certificate for hostname issued by the integration CA.
Mount it into a mock container that must serve TLS as hostname."""
ca_cert = x509.load_pem_x509_certificate(Path(tls.ca_cert_path).read_bytes())
ca_key = serialization.load_pem_private_key(Path(tls.ca_key_path).read_bytes(), password=None)
now = datetime.datetime.now(datetime.UTC)
leaf_key = rsa.generate_private_key(public_exponent=65537, key_size=2048)
leaf_cert = (
x509.CertificateBuilder()
.subject_name(x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, hostname)]))
.issuer_name(ca_cert.subject)
.public_key(leaf_key.public_key())
.serial_number(x509.random_serial_number())
.not_valid_before(now - datetime.timedelta(days=1))
.not_valid_after(now + datetime.timedelta(days=3650))
.add_extension(x509.SubjectAlternativeName([x509.DNSName(hostname)]), critical=False)
.add_extension(x509.ExtendedKeyUsage([x509.oid.ExtendedKeyUsageOID.SERVER_AUTH]), critical=False)
.sign(ca_key, hashes.SHA256())
)
keystore_path = directory / "keystore.p12"
keystore_path.write_bytes(
pkcs12.serialize_key_and_certificates(
name=hostname.encode(),
key=leaf_key,
cert=leaf_cert,
cas=[ca_cert],
encryption_algorithm=serialization.BestAvailableEncryption(KEYSTORE_PASSWORD.encode()),
)
)
return keystore_path

View File

@@ -302,7 +302,6 @@ class Traces(ABC):
db_operation: str
has_error: bool
is_remote: str
scope_json: dict[str, Any]
resource: list[TracesResource]
tag_attributes: list[TracesTagAttributes]
@@ -328,7 +327,6 @@ class Traces(ABC):
links: list[TracesLink] = [],
trace_state: str = "",
flags: np.uint32 = 0,
scope: dict[str, Any] = {},
resource_write_mode: Literal["legacy_only", "dual_write"] = "dual_write",
) -> None:
if timestamp is None:
@@ -410,33 +408,6 @@ class Traces(ABC):
# Calculate resource fingerprint
self.resource_fingerprint = LogsOrTracesFingerprint(self.resources_string).calculate()
# Process scope mirroring the InstrumentationScope on the OTLP span.
scope_name = scope.get("name", "")
scope_version = scope.get("version", "")
scope_string = {k: str(v) for k, v in scope.get("attributes", {}).items()}
self.scope_json = {
"name": scope_name,
"version": scope_version,
"attributes": scope_string,
}
scope_keys = {"scope.name": scope_name, "scope.version": scope_version}
scope_keys.update(scope_string)
for k, v in scope_keys.items():
if v == "":
continue
self.tag_attributes.append(
TracesTagAttributes(
timestamp=timestamp,
tag_key=k,
tag_type="scope",
tag_data_type="string",
string_value=v,
number_value=None,
)
)
self.attribute_keys.append(TracesResourceOrAttributeKeys(name=k, datatype="string", tag_type="scope"))
# Process attributes by type and populate custom fields
self.attribute_string = {}
self.attributes_number = {}
@@ -688,7 +659,6 @@ class Traces(ABC):
self.has_error,
self.is_remote,
self.resource_json,
self.scope_json,
],
dtype=object,
)
@@ -719,7 +689,6 @@ class Traces(ABC):
attributes=data.get("attributes", {}),
trace_state=data.get("trace_state", ""),
flags=data.get("flags", 0),
scope=data.get("scope", {}),
)
@classmethod
@@ -859,7 +828,6 @@ def insert_traces_to_clickhouse(conn, traces: list[Traces]) -> None:
"has_error",
"is_remote",
"resource",
"scope",
],
data=[trace.np_arr() for trace in traces],
)

View File

@@ -158,26 +158,6 @@ class Network:
return f"Network(id={self.id}, name={self.name})"
@dataclass
class TLS:
__test__ = False
ca_cert_path: str
ca_key_path: str
@staticmethod
def from_cache(cache: dict) -> "TLS":
return TLS(ca_cert_path=cache["ca_cert_path"], ca_key_path=cache["ca_key_path"])
def __cache__(self) -> dict:
return {
"ca_cert_path": self.ca_cert_path,
"ca_key_path": self.ca_key_path,
}
def __log__(self) -> str:
return f"TLS(ca_cert_path={self.ca_cert_path}, ca_key_path={self.ca_key_path})"
# Alerts related types

View File

@@ -37,7 +37,6 @@ def test_telemetry_databases_exist(signoz: types.SigNoz) -> None:
def test_teardown(
signoz: types.SigNoz, # pylint: disable=unused-argument
idp: types.TestContainerIDP, # pylint: disable=unused-argument
googleidp: types.TestContainerDocker, # pylint: disable=unused-argument
create_user_admin: types.Operation, # pylint: disable=unused-argument
migrator: types.Operation, # pylint: disable=unused-argument
maildev: types.TestContainerDocker, # pylint: disable=unused-argument

View File

@@ -1,185 +0,0 @@
from collections.abc import Callable
from http import HTTPStatus
import requests
from wiremock.resources.mappings import Mapping
from fixtures import types
from fixtures.auth import (
USER_ADMIN_EMAIL,
USER_ADMIN_PASSWORD,
assert_user_has_role,
find_user_with_roles_by_email,
)
from fixtures.googleidp import GOOGLE_DOMAIN, get_google_domain, google_oidc_mappings, perform_google_login
from fixtures.types import Operation, SigNoz
GOOGLE_CLIENT_ID = "google-client-id.apps.googleusercontent.com"
GOOGLE_CLIENT_SECRET = "google-client-secret"
def test_create_auth_domain(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
# Reruns against a reused stack find the domain from the previous run;
# drop it so creation always starts from a clean slate.
domain = get_google_domain(signoz, admin_token)
if domain:
response = requests.delete(
signoz.self.host_configs["8080"].get(f"/api/v1/domains/{domain['id']}"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v1/domains"),
json={
"name": GOOGLE_DOMAIN,
"config": {
"ssoEnabled": True,
"ssoType": "google_auth",
"googleAuthConfig": {
"clientId": GOOGLE_CLIENT_ID,
"clientSecret": GOOGLE_CLIENT_SECRET,
},
},
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.CREATED
def test_google_authn(
signoz: SigNoz,
googleidp: types.TestContainerDocker,
make_http_mocks: Callable[[types.TestContainerDocker, list[Mapping]], None],
get_token: Callable[[str, str], str],
get_session_context: Callable[[str], dict],
) -> None:
email = "viewer@google.integration.test"
make_http_mocks(googleidp, google_oidc_mappings(email=email, name="Google Viewer", hd=GOOGLE_DOMAIN, audience=GOOGLE_CLIENT_ID))
redirect_url = perform_google_login(signoz, googleidp, get_session_context, email)
assert "accessToken=" in redirect_url
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
found_user = find_user_with_roles_by_email(signoz, admin_token, email)
assert found_user["displayName"] == "Google Viewer"
assert_user_has_role(found_user, "signoz-viewer")
def test_google_authn_hd_mismatch(
signoz: SigNoz,
googleidp: types.TestContainerDocker,
make_http_mocks: Callable[[types.TestContainerDocker, list[Mapping]], None],
get_token: Callable[[str, str], str],
get_session_context: Callable[[str], dict],
) -> None:
# The id_token carries a hosted-domain claim for a different workspace than
# the auth domain; the callback must reject it and provision no user.
email = "intruder@google.integration.test"
make_http_mocks(googleidp, google_oidc_mappings(email=email, name="Intruder", hd="other.workspace.test", audience=GOOGLE_CLIENT_ID))
redirect_url = perform_google_login(signoz, googleidp, get_session_context, email)
assert "callbackauthnerr" in redirect_url
assert "accessToken=" not in redirect_url
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK
assert not any(user["email"] == email for user in response.json()["data"])
def test_google_authn_unverified_email(
signoz: SigNoz,
googleidp: types.TestContainerDocker,
make_http_mocks: Callable[[types.TestContainerDocker, list[Mapping]], None],
get_token: Callable[[str, str], str],
get_session_context: Callable[[str], dict],
) -> None:
email = "unverified@google.integration.test"
make_http_mocks(googleidp, google_oidc_mappings(email=email, name="Unverified", hd=GOOGLE_DOMAIN, audience=GOOGLE_CLIENT_ID, email_verified=False))
redirect_url = perform_google_login(signoz, googleidp, get_session_context, email)
assert "callbackauthnerr" in redirect_url
# Opting the domain into insecureSkipEmailVerified must let the same
# unverified identity through.
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
domain = get_google_domain(signoz, admin_token)
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/domains/{domain['id']}"),
json={
"config": {
"ssoEnabled": True,
"ssoType": "google_auth",
"googleAuthConfig": {
"clientId": GOOGLE_CLIENT_ID,
"clientSecret": GOOGLE_CLIENT_SECRET,
"insecureSkipEmailVerified": True,
},
},
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT
redirect_url = perform_google_login(signoz, googleidp, get_session_context, email)
assert "accessToken=" in redirect_url
found_user = find_user_with_roles_by_email(signoz, admin_token, email)
assert_user_has_role(found_user, "signoz-viewer")
def test_google_role_mapping_default_role(
signoz: SigNoz,
googleidp: types.TestContainerDocker,
make_http_mocks: Callable[[types.TestContainerDocker, list[Mapping]], None],
get_token: Callable[[str, str], str],
get_session_context: Callable[[str], dict],
) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
domain = get_google_domain(signoz, admin_token)
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/domains/{domain['id']}"),
json={
"config": {
"ssoEnabled": True,
"ssoType": "google_auth",
"googleAuthConfig": {
"clientId": GOOGLE_CLIENT_ID,
"clientSecret": GOOGLE_CLIENT_SECRET,
},
"roleMapping": {
"defaultRole": "EDITOR",
},
},
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT
email = "editor@google.integration.test"
make_http_mocks(googleidp, google_oidc_mappings(email=email, name="Google Editor", hd=GOOGLE_DOMAIN, audience=GOOGLE_CLIENT_ID))
redirect_url = perform_google_login(signoz, googleidp, get_session_context, email)
assert "accessToken=" in redirect_url
found_user = find_user_with_roles_by_email(signoz, admin_token, email)
assert_user_has_role(found_user, "signoz-editor")

View File

@@ -130,45 +130,6 @@ def test_reset_password(signoz: types.SigNoz, get_token: Callable[[str, str], st
assert token is not None
def test_reset_password_v2(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
found_user = find_user_by_email(signoz, admin_token, PASSWORD_USER_EMAIL)
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v2/users/{found_user['id']}/reset_password_tokens"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.CREATED, response.text
token = response.json()["data"]["token"]
# A password failing the strength policy is rejected without consuming the token
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/factor_password/reset"),
json={"password": "password", "token": token},
timeout=2,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/factor_password/reset"),
json={"password": "resetV2Password123Z$", "token": token},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
assert get_token(PASSWORD_USER_EMAIL, "resetV2Password123Z$") is not None
# The token is single use, so replaying it no longer resolves
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/factor_password/reset"),
json={"password": "resetV2Password456Z$", "token": token},
timeout=2,
)
assert response.status_code == HTTPStatus.NOT_FOUND, response.text
def test_reset_password_with_no_password(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)

View File

@@ -1240,13 +1240,6 @@ def test_traces_list_span_scope(
lambda x: {"duration_nano": int(x[1].duration_nano), "span_id": x[1].span_id, "timestamp": format_timestamp(x[1].timestamp), "trace_id": x[1].trace_id},
id="select_attribute_duration_order_intrinsic",
),
# Case 9: filter on the intrinsic scope.version. Only x[1] should match.
pytest.param(
BuilderQuery(signal="traces", name="A", select_fields=[TelemetryFieldKey("timestamp")], filter_expression="scope.version = '1.0.0'", limit=1),
HTTPStatus.OK,
lambda x: {"span_id": x[1].span_id, "timestamp": format_timestamp(x[1].timestamp), "trace_id": x[1].trace_id},
id="filter_scope_version",
),
],
)
def test_traces_list_with_corrupt_data(
@@ -1290,156 +1283,6 @@ def test_traces_list_with_corrupt_data(
assert get_rows(response)[0]["data"] == expected(traces)
@pytest.mark.parametrize(
"filter_expression,expected_indices",
[
# Intrinsic scope.name / scope.version resolve to the JSON sub-columns.
pytest.param("scope.name = 'io.signoz.payment'", [1], id="intrinsic_scope_name"),
pytest.param("scope.version = '2.3.1'", [0], id="intrinsic_scope_version"),
# A scope attribute resolves against the scope JSON column's attributes.
pytest.param("scope.telemetry.sdk.language = 'python'", [1], id="scope_attribute"),
# `env.tier` is a span attribute on span 0 and a scope attribute on
# span 1. Unprefixed -> no explicit context, so it is checked in every
# applicable context (attribute OR scope) and both spans match.
pytest.param("env.tier = 'gold'", [0, 1], id="bare_cross_context"),
# The explicit `scope.` prefix forces scope context only, so span 0's
# span attribute is ignored — only span 1 matches.
pytest.param("scope.env.tier = 'gold'", [1], id="scope_prefixed_cross_context"),
# `scope.name` matches BOTH the intrinsic scope.name field (span 0) and a
# scope attribute literally named `name` (span 1's scope attribute
# name='io.signoz.checkout').
pytest.param("scope.name = 'io.signoz.checkout'", [0, 1], id="scope_name_collision"),
# `scope.name` also matches a span attribute literally named `scope.name`
# (attribute context) — span 2 carries attribute scope.name='attr-scope-name'.
pytest.param("scope.name = 'attr-scope-name'", [2], id="scope_name_attribute_collision"),
# An unprefixed `name` resolves to the intrinsic span `name` column and a
# `name` scope attribute, but NOT the scope.name field. Span 2's span
# name and span 1's scope attribute `name` both equal 'io.signoz.checkout';
# span 0's scope.name field equals it too but is NOT matched.
pytest.param("name = 'io.signoz.checkout'", [1, 2], id="bare_name_excludes_scope_name_field"),
# A value that no resolvable key holds (scope.name/scope.version field,
# a `name`/`version` scope attribute, or a same-named attribute/resource)
# returns nothing.
pytest.param("scope.version = 'corrupt_data'", [], id="scope_version_no_match"),
pytest.param("scope.name = 'corrupt_data'", [], id="scope_name_no_match"),
],
)
def test_traces_list_with_scope_filter(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
filter_expression: str,
expected_indices: list[int],
) -> None:
"""
Setup three spans with different scope key resolution:
- x[0]: scope.name/version 'io.signoz.checkout'/'2.3.1'; span attribute
env.tier='gold'.
- x[1]: scope.name/version 'io.signoz.payment'/'4.5.6'; scope attributes
telemetry.sdk.language='python', env.tier='gold', and a `name` scope
attribute colliding with x[0]'s scope.name value.
- x[2]: span name 'io.signoz.checkout' (colliding with x[0]'s scope.name
value) and a span attribute literally named `scope.name`.
Tests:
- Filtering on scope.name / scope.version / a scope attribute.
- An unprefixed key is resolved across contexts (scope checked alongside
attribute / intrinsic), while a `scope.`-prefixed key is scope-only.
- `scope.name` hits the intrinsic field, a `name` scope attribute, and a
span attribute `scope.name` (cross-context), while a bare `name` hits
the span name column (and a `name` scope attribute) but never the
scope.name field.
"""
now = datetime.now(tz=UTC).replace(microsecond=0)
trace_id = TraceIdGenerator.trace_id()
span_ids = [TraceIdGenerator.span_id() for _ in range(3)]
traces = [
Traces(
timestamp=now - timedelta(seconds=4),
duration=timedelta(seconds=2),
trace_id=trace_id,
span_id=span_ids[0],
parent_span_id="",
name="GET /checkout",
kind=TracesKind.SPAN_KIND_SERVER,
status_code=TracesStatusCode.STATUS_CODE_OK,
resources={"service.name": "checkout"},
attributes={"http.request.method": "GET", "env.tier": "gold"},
scope={
"name": "io.signoz.checkout",
"version": "2.3.1",
"attributes": {"telemetry.sdk.language": "go"},
},
),
Traces(
timestamp=now - timedelta(seconds=2),
duration=timedelta(seconds=1),
trace_id=trace_id,
span_id=span_ids[1],
parent_span_id="",
name="POST /pay",
kind=TracesKind.SPAN_KIND_SERVER,
status_code=TracesStatusCode.STATUS_CODE_OK,
resources={"service.name": "payment"},
attributes={"http.request.method": "POST"},
# env.tier is a scope attribute here (cross-context with span 0);
# `name` is a scope attribute colliding with span 0's scope.name.
scope={
"name": "io.signoz.payment",
"version": "4.5.6",
"attributes": {
"telemetry.sdk.language": "python",
"env.tier": "gold",
"name": "io.signoz.checkout",
},
},
),
Traces(
timestamp=now - timedelta(seconds=1),
duration=timedelta(seconds=1),
trace_id=trace_id,
span_id=span_ids[2],
parent_span_id="",
# span name collides with span 0's scope.name value
name="io.signoz.checkout",
kind=TracesKind.SPAN_KIND_SERVER,
status_code=TracesStatusCode.STATUS_CODE_OK,
resources={"service.name": "probe"},
# a span attribute named `scope.name`
attributes={"scope.name": "attr-scope-name"},
scope={"name": "span-gamma", "version": "9.9.9"},
),
]
insert_traces(traces)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms, end_ms = _query_window(now)
response = make_query_request(
signoz,
token,
start_ms=start_ms,
end_ms=end_ms,
request_type=RequestType.RAW,
queries=[
BuilderQuery(
signal="traces",
name="A",
select_fields=[TelemetryFieldKey("timestamp")],
filter_expression=filter_expression,
limit=10,
).to_dict()
],
)
assert response.status_code == HTTPStatus.OK, response.text
got_span_ids = {row["data"]["span_id"] for row in get_rows(response)}
expected_span_ids = {traces[i].span_id for i in expected_indices}
assert got_span_ids == expected_span_ids
@pytest.mark.parametrize("surface", ["filter", "select", "order"])
def test_traces_list_unknown_span_context_synthesizes(
signoz: types.SigNoz,

View File

@@ -19,8 +19,6 @@ dependencies = [
"fastapi>=0.115",
"uvicorn[standard]>=0.34",
"py>=1.11",
"cryptography>=50.0.0",
"jwcrypto>=1.5.8",
]
[dependency-groups]

4
tests/uv.lock generated
View File

@@ -1070,10 +1070,8 @@ version = "0.1.0"
source = { virtual = "." }
dependencies = [
{ name = "clickhouse-connect" },
{ name = "cryptography" },
{ name = "fastapi" },
{ name = "isodate" },
{ name = "jwcrypto" },
{ name = "numpy" },
{ name = "psycopg2" },
{ name = "py" },
@@ -1095,10 +1093,8 @@ dev = [
[package.metadata]
requires-dist = [
{ name = "clickhouse-connect", specifier = ">=0.8.18" },
{ name = "cryptography", specifier = ">=50.0.0" },
{ name = "fastapi", specifier = ">=0.115" },
{ name = "isodate", specifier = ">=0.7.2" },
{ name = "jwcrypto", specifier = ">=1.5.8" },
{ name = "numpy", specifier = ">=2.3.2" },
{ name = "psycopg2", specifier = ">=2.9.10" },
{ name = "py", specifier = ">=1.11" },