Compare commits

..

34 Commits

Author SHA1 Message Date
Nikhil Soni
883e9492d6 Merge remote-tracking branch 'origin/main' into ns/scope 2026-08-10 15:01:20 +05:30
Pandey
cd8346a91a test(callbackauthn): cover the google authn flow end to end (#12486)
#### Description

- Adds end-to-end coverage for the google authn flow
(`callbackauthn/04_google.py`): happy-path login, hd-claim mismatch
rejection, unverified-email rejection + `insecureSkipEmailVerified`
opt-in, and roleMapping defaultRole.
- The google callback authn hardcodes `https://accounts.google.com` as
its issuer and fully verifies the RS256 id_token, so a wiremock
container impersonates Google: it joins the test network under the
`accounts.google.com` alias and serves HTTPS with a certificate issued
by a new integration CA (`tests/fixtures/tls.py`), which every signoz
container now trusts via `SSL_CERT_FILE`.
- Stubs (discovery, auto-approving authorize, token, JWKS) are installed
per test via the existing `make_http_mocks` fixture with a pre-signed
id_token for the identity under test; one session-scoped RSA key signs
all tokens.
2026-08-10 08:35:58 +00:00
Vikrant Gupta
db5f4b4cd5 feat(user): add v2 reset password endpoint (#12489)
#### Description

- `POST /api/v1/resetPassword` was the last password endpoint with no v2
equivalent. Adds `POST /api/v2/factor_password/reset`, next to the
existing `/factor_password/forgot`, so the whole recovery flow lives
under one namespace.
- v1 keeps working and is now marked deprecated. Both routes share the
same handler, so behaviour is identical.
- A malformed request body now returns a structured 400 instead of a
500. This applies to v1 too, since the handler is shared.

#### Issues closed by this PR

Contributes to SigNoz/platform-pod#2667

#### Additional Information


- The generated frontend client is included because CI re-runs `pnpm
generate:api` and fails on drift. The UI still calls v1; moving it to
the new `useResetPassword` hook is a separate PR to keep review
ownership split.
- Not fixed here: a reset doesn't revoke existing sessions, though a
voluntary password change does. Worth its own ticket.
2026-08-10 08:27:46 +00:00
Aditya Singh
3bd9b2a96a feat(log-details): update log details header (#12310)
## Pull Request

---

### 📄 Summary
> Why does this change exist?  
> What problem does it solve, and why is this the right approach?

Slice 1 of the log-details drawer revamp, a reworked drawer **header**,
gated behind the new `isLogDetailsV2` flag (ships off). Highlights and
the DataViewer land in
the following stacked PRs.

 **Change points**

- Added a new **Log Details Header** with:
  - A formatted timestamp based on the user's timezone
  - A ⋯ menu with **Copy log** and **Copy link to log**
  - Up/down arrows to move between logs
  - An optional **Open in Explorer** button
 
- Moved the log navigation logic into a separate `useLogNavigation` hook
- Updated the Log Details drawer:
- Shows the new header when the feature flag is on (otherwise keeps the
old title)
- Replaced the WARN/ERROR divider with `LogStateIndicator`, which shows
colors for all log levels
- Moved the copy actions into the ⋯ menu and removed the old inline copy
button in V2
- Uses the new navigation hook for both keyboard shortcuts and header
arrows

- Updated the copy link handler:
- `onLogCopy` now accepts an optional click event, so it can be called
from the ⋯ menu
- Moved the "Copied to clipboard" toast from the top-right to the
bottom-right

- Everything is gated behind feature flag for now


#### Screenshots / Screen Recordings (if applicable)


https://github.com/user-attachments/assets/4a73e299-4715-4feb-81e3-762fdc3d0757




#### Issues closed by this PR
> Reference issues using `Closes #issue-number` to enable automatic
closure on merge.

---

###  Change Type
_Select all that apply_

- [ ]  Feature
- [ ] 🐛 Bug fix
- [ ] ♻️ Refactor
- [ ] 🛠️ Infra / Tooling
- [ ] 🧪 Test-only

---

### 🐛 Bug Context
> Required if this PR fixes a bug

#### Root Cause
> What caused the issue?  
> Regression, faulty assumption, edge case, refactor, etc.

#### Fix Strategy
> How does this PR address the root cause?

---

### 🧪 Testing Strategy
> How was this change validated?

- Tests added/updated:
- Manual verification:
- Edge cases covered:

---

### ⚠️ Risk & Impact Assessment
> What could break? How do we recover?

- Blast radius:
- Potential regressions:
- Rollback plan:

---

### 📝 Changelog
> Fill only if this affects users, APIs, UI, or documented behavior  
> Use **N/A** for internal or non-user-facing changes

| Field | Value |
|------|-------|
| Deployment Type | Cloud / OSS / Enterprise |
| Change Type | Feature / Bug Fix / Maintenance |
| Description | User-facing summary |

---

### 📋 Checklist
- [ ] Tests added or explicitly not required
- [ ] Manually tested
- [ ] Breaking changes documented
- [ ] Backward compatibility considered

---

## 👀 Notes for Reviewers

<!-- Anything reviewers should keep in mind while reviewing -->

---
2026-08-10 06:54:54 +00:00
Nikhil Soni
ab3b88966e Merge remote-tracking branch 'origin/main' into ns/scope
# Conflicts:
#	pkg/telemetryschema/tracestelemetryschema/field_mapper.go
#	tests/conftest.py
#	tests/integration/tests/querier/04_traces.py
2026-08-05 16:44:41 +05:30
Nikhil Soni
08ebc37109 fix: handle fixed scope.name and scope.version fields to work without scope 2026-07-02 19:04:41 +05:30
Nikhil Soni
aa5a1c5e62 chore: update collector version 2026-07-02 14:46:07 +05:30
Nikhil Soni
7b34a47ac5 chore: run formatting on tests 2026-07-02 14:41:33 +05:30
Nikhil Soni
b4b2d7bb66 test: add test to show cross context matching 2026-06-24 18:34:38 +05:30
Nikhil Soni
e16416475b refactor: drop unused fields 2026-06-24 18:07:47 +05:30
Nikhil Soni
0ea7c1ae6e test: add more cases for scope name 2026-06-24 18:05:50 +05:30
Nikhil Soni
a023c8ed4a test: add integration test for scope fields 2026-06-24 15:25:17 +05:30
Nikhil Soni
a73ae62cd1 Merge remote-tracking branch 'origin/main' into ns/scope 2026-06-24 13:02:18 +05:30
Nikhil Soni
ec6fb58052 chore: add more tests 2026-06-24 12:37:42 +05:30
Nikhil Soni
d3d13eb7ff fix: remove handling of normalized properties for scope
Otherwise it will be impossible to query if scope attribute also
exists with same name - name and version
2026-05-21 11:19:22 +05:30
Nikhil Soni
782de2b210 fix: use correct error type for internal issues 2026-05-20 15:32:26 +05:30
Nikhil Soni
d3c38693f3 fix: allow 'scope.' prefix for keys with other context 2026-05-20 15:30:25 +05:30
Nikhil Soni
8791df3697 fix: avoid removing context prefix to support attr with prefix 2026-05-19 19:14:23 +05:30
Nikhil Soni
eb719c3d0d fix: use key selector with context prefix 2026-05-19 18:56:00 +05:30
Nikhil Soni
f10435c210 Merge remote-tracking branch 'origin' into ns/scope 2026-05-19 13:55:36 +05:30
Nikhil Soni
f3f1e9cb59 chore: add tests for denormalized field name as well 2026-05-19 13:55:25 +05:30
Nikhil Soni
d0370ce3ef fix: handle fields with included context for scope (select clause) 2026-05-14 17:02:56 +05:30
Nikhil Soni
d169761e65 Merge remote-tracking branch 'origin/main' into ns/scope 2026-05-14 11:50:33 +05:30
Nikhil Soni
87864ef5d4 chore: remove duplicates from .gitignore 2026-05-11 15:45:32 +05:30
Nikhil Soni
2e0bc8998e chore: use name as key name for scope instead of scope.name 2026-05-11 15:40:45 +05:30
Nikhil Soni
7e1f4aa50d Merge remote-tracking branch 'origin/main' into ns/scope 2026-05-11 14:27:19 +05:30
Nikhil Soni
35da39247c Merge branch 'main' into ns/scope 2026-05-07 17:41:11 +05:30
Nikhil Soni
ceccc47a34 fix: fix test for case without resource filter 2026-05-07 16:04:03 +05:30
Nikhil Soni
23da5e22ec Merge branch 'main' into ns/scope 2026-05-07 13:27:34 +05:30
Nikhil Soni
4c1b479149 chore: add tests for scope fields 2026-04-28 20:27:10 +05:30
Nikhil Soni
f72204a8b2 refactor: simplify field mapper for scope 2026-04-28 20:26:37 +05:30
Nikhil Soni
deb3f385fa chore: remove underscore version of scope fields 2026-04-23 10:26:55 +05:30
Nikhil Soni
77ce5f86b1 fix: use scope as json field instead with name and version 2026-04-23 01:15:02 +05:30
Nikhil Soni
ff211de441 feat: add support for scope fields in traces 2026-04-14 10:45:08 +05:30
63 changed files with 1960 additions and 1424 deletions

View File

@@ -12198,9 +12198,9 @@ paths:
- dashboard
/api/v1/resetPassword:
post:
deprecated: false
deprecated: true
description: This endpoint resets the password by token
operationId: ResetPassword
operationId: ResetPasswordDeprecated
requestBody:
content:
application/json:
@@ -15567,6 +15567,41 @@ 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,9 +255,10 @@ export const useCreateInvite = <
};
/**
* This endpoint resets the password by token
* @deprecated
* @summary Reset password
*/
export const resetPassword = (
export const resetPasswordDeprecated = (
typesPostableResetPasswordDTO?: BodyType<TypesPostableResetPasswordDTO>,
signal?: AbortSignal,
) => {
@@ -270,23 +271,23 @@ export const resetPassword = (
});
};
export const getResetPasswordMutationOptions = <
export const getResetPasswordDeprecatedMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof resetPassword>>,
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof resetPassword>>,
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
> => {
const mutationKey = ['resetPassword'];
const mutationKey = ['resetPasswordDeprecated'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
@@ -296,45 +297,47 @@ export const getResetPasswordMutationOptions = <
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof resetPassword>>,
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
{ data?: BodyType<TypesPostableResetPasswordDTO> }
> = (props) => {
const { data } = props ?? {};
return resetPassword(data);
return resetPasswordDeprecated(data);
};
return { mutationFn, ...mutationOptions };
};
export type ResetPasswordMutationResult = NonNullable<
Awaited<ReturnType<typeof resetPassword>>
export type ResetPasswordDeprecatedMutationResult = NonNullable<
Awaited<ReturnType<typeof resetPasswordDeprecated>>
>;
export type ResetPasswordMutationBody =
export type ResetPasswordDeprecatedMutationBody =
| BodyType<TypesPostableResetPasswordDTO>
| undefined;
export type ResetPasswordMutationError = ErrorType<RenderErrorResponseDTO>;
export type ResetPasswordDeprecatedMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @deprecated
* @summary Reset password
*/
export const useResetPassword = <
export const useResetPasswordDeprecated = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof resetPassword>>,
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof resetPassword>>,
Awaited<ReturnType<typeof resetPasswordDeprecated>>,
TError,
{ data?: BodyType<TypesPostableResetPasswordDTO> },
TContext
> => {
return useMutation(getResetPasswordMutationOptions(options));
return useMutation(getResetPasswordDeprecatedMutationOptions(options));
};
/**
* This endpoint lists all users
@@ -593,6 +596,89 @@ 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

@@ -0,0 +1,49 @@
.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

@@ -0,0 +1,153 @@
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

@@ -0,0 +1,52 @@
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

@@ -0,0 +1,161 @@
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,3 +1,10 @@
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,7 +8,9 @@ import { ToggleGroupSimple } from '@signozhq/ui/toggle-group';
import { Divider } from '@signozhq/ui/divider';
import { Typography } from '@signozhq/ui/typography';
import cx from 'classnames';
import { LogType } from 'components/Logs/LogStateIndicator/LogStateIndicator';
import LogStateIndicator, {
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';
@@ -23,6 +25,7 @@ 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';
@@ -48,8 +51,10 @@ import { ILogBody } from 'types/api/logs/log';
import { Query, TagFilter } from 'types/api/queryBuilder/queryBuilderData';
import { DataSource, StringOperators } from 'types/common/queryBuilder';
import { RESOURCE_KEYS, VIEW_TYPES, VIEWS } from './constants';
import { isLogDetailsV2, 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';
@@ -96,7 +101,8 @@ 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('.query-status-popover') ||
target.closest('[data-radix-popper-content-wrapper]')
) {
return;
}
@@ -112,49 +118,30 @@ function LogDetailInner({
};
}, [onClose]);
// Keyboard navigation - handle up/down arrow keys
// Only listen when in OVERVIEW tab
// eslint-disable-next-line sonarjs/cognitive-complexity
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.
useEffect(() => {
if (
!logs ||
!onNavigateLog ||
logs.length === 0 ||
selectedView !== VIEW_TYPES.OVERVIEW
) {
return;
if (selectedView !== VIEW_TYPES.OVERVIEW) {
return undefined;
}
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();
// 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);
}
}
goToPrev();
} else if (e.key === 'ArrowDown') {
e.preventDefault();
e.stopPropagation();
// 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);
}
}
goToNext();
}
};
@@ -162,7 +149,7 @@ function LogDetailInner({
return (): void => {
document.removeEventListener('keydown', handleKeyDown);
};
}, [log.id, logs, onNavigateLog, onScrollToLog, selectedView]);
}, [selectedView, goToPrev, goToNext]);
const listQuery = useMemo(() => {
if (!stagedQuery || stagedQuery.builder.queryData.length < 1) {
@@ -303,33 +290,6 @@ 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
@@ -338,57 +298,69 @@ function LogDetailInner({
maskClosable={false}
getContainer={getContainer}
title={
<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>
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>
{handleOpenInExplorer && (
<div>
<Button
variant="outlined"
color="secondary"
prefix={<Compass size={16} />}
className="open-in-explorer-btn"
onClick={handleOpenInExplorer}
<div className="log-detail-drawer__title-right">
<div className="log-arrows">
<Tooltip
title={isPrevDisabled ? '' : 'Move to previous log'}
placement="top"
mouseLeaveDelay={0}
>
Open in Explorer
</Button>
<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>
</div>
)
}
placement="right"
onClose={drawerCloseHandler}
@@ -407,7 +379,15 @@ function LogDetailInner({
data-testid="log-detail-drawer"
>
<div className="log-detail-drawer__log">
<Divider type="vertical" className={cx('log-type-indicator', logType)} />
{isLogDetailsV2 ? (
<LogStateIndicator
severityText={log.severity_text}
severityNumber={log.severity_number}
fontSize={options?.fontSize ?? FontSize.MEDIUM}
/>
) : (
<Divider type="vertical" className={cx('log-type-indicator', logType)} />
)}
<Tooltip
title={removeEscapeCharacters(logBody)}
placement="left"
@@ -483,22 +463,25 @@ function LogDetailInner({
</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>
{/* 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>
)}
</div>
</div>
{isFilterVisible && contextQuery?.builder.queryData[0] && (

View File

@@ -13,6 +13,7 @@ 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 { MouseEventHandler } from 'react';
import { MouseEvent } 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: MouseEventHandler<HTMLElement>;
onLogCopy: (event?: MouseEvent<HTMLElement>) => void;
onClearActiveLog: () => void;
};

View File

@@ -1,10 +1,4 @@
import {
MouseEventHandler,
useCallback,
useEffect,
useMemo,
useState,
} from 'react';
import { MouseEvent, useCallback, useEffect, useMemo, useState } from 'react';
// eslint-disable-next-line no-restricted-imports
import { useSelector } from 'react-redux';
import { useLocation } from 'react-router-dom';
@@ -46,14 +40,14 @@ export const useCopyLogLink = (logId?: string): UseCopyLogLink => {
[pathname],
);
const onLogCopy: MouseEventHandler<HTMLElement> = useCallback(
(event) => {
const onLogCopy = useCallback(
(event?: MouseEvent<HTMLElement>): void => {
if (!logId) {
return;
}
event.preventDefault();
event.stopPropagation();
event?.preventDefault();
event?.stopPropagation();
urlQuery.delete(QueryParams.activeLogId);
urlQuery.delete(QueryParams.relativeTime);
@@ -66,7 +60,7 @@ export const useCopyLogLink = (logId?: string): UseCopyLogLink => {
setCopy(link);
toast.success('Copied to clipboard', { position: 'top-right' });
toast.success('Copied to clipboard', { position: 'bottom-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: "ResetPassword",
ID: "ResetPasswordDeprecated",
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: false,
Deprecated: true,
SecuritySchemes: []handler.OpenAPISecurityScheme{},
})).Methods(http.MethodPost).GetError(); err != nil {
return err
@@ -299,6 +299,23 @@ 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,10 +40,7 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
// 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)
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -64,23 +64,6 @@ 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 := json.NewDecoder(r.Body).Decode(req); err != nil {
if err := binding.JSON.BindBody(r.Body, req); err != nil {
render.Error(w, err)
return
}

View File

@@ -1,43 +0,0 @@
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,25 +21,24 @@ const (
hasTokenFunctionDocURL = "https://signoz.io/docs/userguide/functions-reference/#hastoken-function"
)
// 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, ""
// 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, ""
}
warning := fmt.Sprintf(
"Key `%s` is ambiguous, found %d different combinations of field context / data type: %v.",
field.Name,
len(logicalFields),
logicalFields,
len(fieldKeysForName),
fieldKeysForName,
)
hasResource, hasAttribute := false, false
for _, item := range logicalFields {
for _, item := range fieldKeysForName {
switch item.FieldContext {
case telemetrytypes.FieldContextResource:
hasResource = true
@@ -50,40 +49,18 @@ func ResolveLogicalFields(field *telemetrytypes.TelemetryFieldKey, logicalFields
// when there is both resource and attribute context, default to resource only
if hasResource && hasAttribute {
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
for _, item := range logicalFields {
filteredKeys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fieldKeysForName))
for _, item := range fieldKeysForName {
if item.FieldContext == telemetrytypes.FieldContextResource {
filtered = append(filtered, item)
filteredKeys = append(filteredKeys, item)
}
}
logicalFields = filtered
fieldKeysForName = filteredKeys
warning += " " + "Using `resource` context by default. To query attributes explicitly, " +
fmt.Sprintf("use the fully qualified name (e.g., 'attribute.%s')", field.Name)
}
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
return fieldKeysForName, warning
}
// NewKeyNotFoundError builds the error a condition builder returns when a filter term

View File

@@ -1,94 +0,0 @@
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

@@ -1,95 +0,0 @@
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

@@ -1,165 +0,0 @@
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,6 +56,17 @@ 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,6 +72,23 @@ 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,7 +10,6 @@ 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"
@@ -361,7 +360,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
return ErrorConditionLiteral
}
}
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
if !ok {
return ErrorConditionLiteral
}
@@ -380,7 +379,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 := MatchingLogicalFields(key, v.fieldKeys)
matching := MatchingFieldKeys(key, v.fieldKeys)
// Handle EXISTS specially
if ctx.EXISTS() != nil {
@@ -676,7 +675,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.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
if !ok {
return ErrorConditionLiteral
}
@@ -731,7 +730,7 @@ func (v *filterExpressionVisitor) VisitFunctionCall(ctx *grammar.FunctionCallCon
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(key, MatchingLogicalFields(key, v.fieldKeys), operator, value)
conds, ok := v.buildConditions(key, MatchingFieldKeys(key, v.fieldKeys), operator, value)
if !ok {
return ErrorConditionLiteral
}
@@ -923,7 +922,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.LogicalField, op qbtypes.FilterOperator, value any) ([]string, bool) {
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) ([]string, bool) {
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter}, op, value, v.builder)
if err != nil {
_, _, _, _, errURL, _ := errors.Unwrapb(err)
@@ -980,123 +979,30 @@ func assignIfEmpty(s *string, value string) {
}
}
// familyMemberNames returns the physical spellings to look up for the
// referenced key: the semantic-convention family members (current-first) when
// 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,
})
}
// 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{}
// 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
}
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)
}
// 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)
}
}
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.
// 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 {
for _, member := range members {
appendMatches(fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), member), member, true)
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)
}
}
}
// 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
return fieldKeysForName
}

View File

@@ -588,11 +588,9 @@ 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). 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)
// builder behavior (error unless the key is ignored).
matching := MatchingFieldKeys(key, tt.fieldKeys)
keys, warning := ResolveKeys(key, matching)
var gotErrors []string
var gotMainErrURL, gotMainWrnURL string
@@ -768,8 +766,7 @@ func (b *resourceConditionBuilder) ConditionFor(
return nil, nil, nil
}
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
keys := SingleKeys(resolved)
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -811,8 +808,7 @@ func (b *conditionBuilder) ConditionFor(
return []string{fmt.Sprintf("%s_cond", key.Name)}, nil, nil
}
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
keys := SingleKeys(resolved)
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -44,70 +44,6 @@ 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,
@@ -121,7 +57,7 @@ func (b *defaultConditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
matches := querybuilder.MatchingFieldKeys(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).
@@ -129,21 +65,21 @@ func (b *defaultConditionBuilder) ConditionFor(
return nil, nil, nil
}
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
keys, warning := querybuilder.ResolveKeys(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
}
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
// the resource fingerprint table only stores resource attributes; fields from
conds := make([]string, 0, len(keys))
for _, k := range keys {
// the resource fingerprint table only stores resource attributes; keys 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 logical.FieldContext != telemetrytypes.FieldContextResource {
if k.FieldContext != telemetrytypes.FieldContextResource {
continue
}
cond, err := b.conditionForLogicalField(ctx, startNs, endNs, logical, op, value, sb)
cond, err := b.conditionForKey(ctx, startNs, endNs, k, op, value, sb)
if err != nil {
return nil, nil, err
}
@@ -152,11 +88,11 @@ func (b *defaultConditionBuilder) ConditionFor(
return conds, warnings, nil
}
func (b *defaultConditionBuilder) conditionForLogicalField(
func (b *defaultConditionBuilder) conditionForKey(
ctx context.Context,
startNs uint64,
endNs uint64,
logical *telemetrytypes.LogicalField,
key *telemetrytypes.TelemetryFieldKey,
op qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -166,7 +102,7 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
// as we store resource values as string
formattedValue := querybuilder.FormatValueForContains(value)
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, logical.Single())
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, key)
if err != nil {
return "", err
}
@@ -179,12 +115,10 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
// as we have not changed the resource column in the resource fingerprint table.
column := columns[0]
members := logical.Members
isFamily := logical.IsFamily()
keyIdxFilter := keyIndexCondition(sb, column.Name, members)
singleValueIndexFilter := valueForIndexFilter(op, members[0], value)
keyIdxFilter := sb.Like(column.Name, keyIndexFilter(key))
valueForIndexFilter := valueForIndexFilter(op, key, value)
fieldName, err := querybuilder.LogicalValueExpr(ctx, valuer.UUID{}, startNs, endNs, b.fm, logical)
fieldName, err := b.fm.FieldFor(ctx, valuer.UUID{}, startNs, endNs, key)
if err != nil {
return "", err
}
@@ -194,17 +128,12 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
return sb.And(
sb.E(fieldName, formattedValue),
keyIdxFilter,
valueIndexCondition(sb, column.Name, members, op, value, false),
sb.Like(column.Name, valueForIndexFilter),
), 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, singleValueIndexFilter),
sb.NotLike(column.Name, valueForIndexFilter),
), nil
case qbtypes.FilterOperatorGreaterThan:
return sb.And(sb.GT(fieldName, formattedValue), keyIdxFilter), nil
@@ -219,7 +148,7 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
return sb.And(
sb.ILike(fieldName, formattedValue),
keyIdxFilter,
valueIndexCondition(sb, column.Name, members, op, value, true),
sb.ILike(column.Name, valueForIndexFilter),
), nil
case qbtypes.FilterOperatorNotLike, qbtypes.FilterOperatorNotILike:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
@@ -256,11 +185,13 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
inConditions = append(inConditions, sb.E(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.Or(inConditions...)
mainCondition = sb.And(
mainCondition,
keyIdxFilter,
valueIndexCondition(sb, column.Name, members, op, value, false),
)
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...))
return mainCondition, nil
case qbtypes.FilterOperatorNotIn:
@@ -273,13 +204,8 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
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 := singleValueIndexFilter.([]string); ok {
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.NotLike(column.Name, v))
}
@@ -289,11 +215,13 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
case qbtypes.FilterOperatorExists:
return sb.And(
memberPresenceCondition(sb, column.Name, members, true),
sb.E(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
keyIdxFilter,
), nil
case qbtypes.FilterOperatorNotExists:
return memberPresenceCondition(sb, column.Name, members, false), nil
return sb.And(
sb.NE(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
), nil
case qbtypes.FilterOperatorRegexp:
return sb.And(
@@ -309,7 +237,7 @@ func (b *defaultConditionBuilder) conditionForLogicalField(
return sb.And(
sb.ILike(fieldName, fmt.Sprintf(`%%%s%%`, formattedValue)),
keyIdxFilter,
valueIndexCondition(sb, column.Name, members, op, value, true),
sb.ILike(column.Name, valueForIndexFilter),
), nil
case qbtypes.FilterOperatorNotContains:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else

View File

@@ -1,90 +0,0 @@
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,33 +71,6 @@ 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 := querybuilder.ExpandKeySelectorsForFamilies(b.getKeySelectors(query))
keySelectors := 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, querybuilder.ExpandKeySelectorsForFamilies(selectors))
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, 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, querybuilder.ExpandKeySelectorsForFamilies(selectors))
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, 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 := querybuilder.ExpandKeySelectorsForFamilies(getKeySelectors(query))
keySelectors := getKeySelectors(query)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {

View File

@@ -75,27 +75,6 @@ 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,
@@ -394,6 +373,94 @@ 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)
@@ -820,6 +887,52 @@ 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 := querybuilder.ExpandKeySelectorsForFamilies(getKeySelectors(*query))
keySelectors := 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 := querybuilder.ExpandKeySelectorsForFamilies(b.getKeySelectors())
keySelectors := b.getKeySelectors()
keys, _, err := b.stmtBuilder.metadataStore.GetKeysMulti(ctx, b.orgID, keySelectors)
if err != nil {
return nil, err

View File

@@ -38,11 +38,8 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
// 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)
// an unknown key simply yields no condition rather than an error.
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -58,19 +58,6 @@ 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,10 +139,7 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
// 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)
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -97,19 +97,6 @@ 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.MatchingLogicalFields(key, fieldKeys)
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
// search() resolves its own (optional) scope; handle it before key resolution.
@@ -460,10 +460,7 @@ func (c *conditionBuilder) ConditionFor(
return c.conditionForSearch(ctx, orgID, key, value, sb)
}
// 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)
keys, warning := querybuilder.ResolveKeys(key, matches)
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.ExistsFor(ctx, orgID, tsStart, tsEnd, key, true)
guard, err := m.existsExpressionFor(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.ExistsFor(ctx, orgID, tsStart, tsEnd, candidates[0], true)
guard, err := m.existsExpressionFor(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.ExistsFor(ctx, orgID, tsStart, tsEnd, key, true)
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, key, true)
if err != nil {
return "", err
}
@@ -569,8 +569,7 @@ func (m *fieldMapper) membershipGuarded(ctx context.Context, orgID valuer.UUID,
return columnType == schema.ColumnTypeEnumMap || columnType == schema.ColumnTypeEnumJSON
}
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
func (m *fieldMapper) ExistsFor(
func (m *fieldMapper) existsExpressionFor(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,

View File

@@ -162,9 +162,7 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
// 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))
keys := querybuilder.MatchingFieldKeys(key, fieldKeys)
var warnings []string
if len(keys) == 0 {
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {

View File

@@ -97,18 +97,6 @@ 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,
logical *telemetrytypes.LogicalField,
key *telemetrytypes.TelemetryFieldKey,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -42,13 +42,13 @@ func (c *conditionBuilder) conditionFor(
value = querybuilder.FormatValueForContains(value)
}
fieldExpression, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, c.fm, logical)
fieldExpression, err := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err
}
// TODO(srikanthccv): maybe extend this to every possible attribute
if logical.Name == "duration_nano" || logical.Name == "durationNano" { // QoL improvement
if key.Name == "duration_nano" || key.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(logical.Single(), value, fieldExpression, operator)
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(key, value, fieldExpression, operator)
// regular operators
switch operator {
@@ -154,7 +154,11 @@ 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:
pred, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, c.fm, logical, operator == qbtypes.FilterOperatorExists)
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)
if err != nil {
return "", err
}
@@ -206,10 +210,10 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
keys, warning := querybuilder.ResolveKeys(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -217,10 +221,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(logicalFields) > 0 {
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(keys) > 0 {
hasColumn := false
for _, logical := range logicalFields {
if logical.FieldContext == telemetrytypes.FieldContextSpan {
for _, k := range keys {
if k.FieldContext == telemetrytypes.FieldContextSpan {
hasColumn = true
break
}
@@ -228,49 +232,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.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)
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)
}
}
logicalFields = combined
keys = combined
}
}
}
synthesized := false
if len(logicalFields) == 0 {
if len(keys) == 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.
logicalFields = querybuilder.WrapAsLogicalFields(key.Name, c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys)))
if len(logicalFields) == 0 {
keys = c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys))
if len(keys) == 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 fields from the main
// When a resource sub-query already covers the term, drop resource keys from the main
// query. Synthesized keys are exempt: the sub-query skips keys absent from metadata.
if skipResourceFilter && !synthesized {
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
for _, logical := range logicalFields {
if logical.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, logical)
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
for _, k := range keys {
if k.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, k)
}
}
if len(filtered) == 0 {
return nil, warnings, nil
}
logicalFields = filtered
keys = filtered
}
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
cond, err := c.conditionForLogicalField(ctx, orgID, startNs, endNs, logical, operator, value, sb)
conds := make([]string, 0, len(keys))
for _, k := range keys {
cond, err := c.conditionForKey(ctx, orgID, startNs, endNs, k, operator, value, sb)
if err != nil {
return nil, nil, err
}
@@ -279,28 +283,28 @@ func (c *conditionBuilder) ConditionFor(
return conds, warnings, nil
}
func (c *conditionBuilder) conditionForLogicalField(
func (c *conditionBuilder) conditionForKey(
ctx context.Context,
orgID valuer.UUID,
startNs uint64,
endNs uint64,
logical *telemetrytypes.LogicalField,
key *telemetrytypes.TelemetryFieldKey,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
) (string, error) {
if c.isSpanScopeField(logical.Name) {
return c.buildSpanScopeCondition(logical.Single(), operator, value, startNs)
if c.isSpanScopeField(key.Name) {
return c.buildSpanScopeCondition(key, operator, value, startNs)
}
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, operator, value, sb)
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, 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, logical.Single())
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
if slices.Contains(maps.Keys(IntrinsicFields), field) ||
slices.Contains(maps.Keys(IntrinsicFieldsDeprecated), field) ||
slices.Contains(maps.Keys(CalculatedFields), field) ||
@@ -308,7 +312,7 @@ func (c *conditionBuilder) conditionForLogicalField(
return condition, nil
}
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, qbtypes.FilterOperatorExists, nil, sb)
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, qbtypes.FilterOperatorExists, nil, sb)
if err != nil {
return "", err
}

View File

@@ -121,6 +121,20 @@ 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

@@ -1,148 +0,0 @@
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,6 +52,7 @@ 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,
@@ -176,7 +177,7 @@ func (m *fieldMapper) getColumn(
case telemetrytypes.FieldContextResource:
return []*schema.Column{indexV3Columns["resource"], indexV3Columns["resources_string"]}, nil
case telemetrytypes.FieldContextScope:
return []*schema.Column{}, qbtypes.ErrColumnNotFound
return []*schema.Column{indexV3Columns["scope"]}, nil
case telemetrytypes.FieldContextAttribute:
switch key.FieldDataType {
case telemetrytypes.FieldDataTypeString:
@@ -287,14 +288,24 @@ 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.
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
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)
}
case schema.ColumnTypeEnumString,
schema.ColumnTypeEnumUInt64,
schema.ColumnTypeEnumUInt32,
@@ -336,68 +347,6 @@ 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.
@@ -410,23 +359,18 @@ func (m *fieldMapper) ColumnExpressionFor(
keys map[string][]*telemetrytypes.TelemetryFieldKey,
) (string, error) {
// Resolve the candidate logical field(s).
var candidates []*telemetrytypes.LogicalField
// Resolve the candidate column(s).
var candidates []*telemetrytypes.TelemetryFieldKey
switch _, err := m.FieldFor(ctx, orgID, startNs, endNs, field); {
case err == nil:
// 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)}
candidates = []*telemetrytypes.TelemetryFieldKey{field}
case errors.Is(err, qbtypes.ErrColumnNotFound):
// 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 {
// 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 {
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
}
@@ -440,21 +384,21 @@ func (m *fieldMapper) ColumnExpressionFor(
dummyValue = 0.0
}
stmts := make([]string, 0, len(candidates)*2)
for _, logical := range candidates {
value, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logical)
for _, key := range candidates {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err
}
guard, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, m, logical, true)
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
if err != nil {
return "", err
}
coerced := value
// a time column keeps its native type; coercing it would yield seconds
if temporal, err := m.logicalIsTemporal(ctx, startNs, endNs, logical); err != nil {
if temporal, err := m.columnIsTemporal(ctx, startNs, endNs, key); err != nil {
return "", err
} else if !temporal {
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), dummyValue, value, qbtypes.FilterOperatorUnknown)
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(key, dummyValue, value, qbtypes.FilterOperatorUnknown)
}
stmts = append(stmts, guard, coerced)
}
@@ -462,14 +406,13 @@ func (m *fieldMapper) ColumnExpressionFor(
}
if len(candidates) == 1 {
logical := candidates[0]
value, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logical)
value, err := m.FieldFor(ctx, orgID, startNs, endNs, candidates[0])
if err != nil {
return "", err
}
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)
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)
if err != nil {
return "", err
}
@@ -481,12 +424,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 _, logical := range candidates {
value, err := querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logical)
for _, key := range candidates {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err
}
guard, err := querybuilder.LogicalExistsExpr(ctx, orgID, startNs, endNs, m, logical, true)
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
if err != nil {
return "", err
}
@@ -495,15 +438,6 @@ 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) {
@@ -599,8 +533,7 @@ func (m *fieldMapper) CandidateKeys(ctx context.Context, _ valuer.UUID, field *t
return nil
}
// ExistsFor implements the per-key existence primitive of qbtypes.FieldMapper.
func (m *fieldMapper) ExistsFor(
func (m *fieldMapper) existsExpressionFor(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,

View File

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

View File

@@ -34,11 +34,6 @@ 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

@@ -1,69 +0,0 @@
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,6 +24,8 @@ pytest_plugins = [
"fixtures.browser",
"fixtures.keycloak",
"fixtures.idp",
"fixtures.googleidp",
"fixtures.tls",
"fixtures.notification_channel",
"fixtures.maildev",
"fixtures.alerts",

231
tests/fixtures/googleidp.py vendored Normal file
View File

@@ -0,0 +1,231 @@
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,6 +999,8 @@ 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",
@@ -1007,7 +1009,10 @@ 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),
@@ -1027,12 +1032,24 @@ 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(
@@ -1053,12 +1070,15 @@ 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),
@@ -1077,6 +1097,7 @@ 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",
@@ -1084,7 +1105,10 @@ 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,6 +39,7 @@ 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.
@@ -50,6 +51,7 @@ 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()
@@ -62,8 +64,14 @@ def wrap( # pylint: disable=too-many-arguments,too-many-positional-arguments
delete(restore(existing_resource))
pytestconfig.cache.set(key, None)
else:
logger.info("Reusing existing %s(%s)", key, existing_resource)
return restore(existing_resource)
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
if not teardown(request):
resource = create()
@@ -88,15 +96,23 @@ 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,6 +13,7 @@ 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__)
@@ -27,10 +28,12 @@ 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.
Accepts optional env_overrides to customize the container environment, and
an optional integration CA (tls) to trust in addition to the system roots.
"""
def create() -> types.SigNoz:
@@ -115,6 +118,13 @@ 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:
@@ -193,6 +203,16 @@ 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,
@@ -212,6 +232,7 @@ def create_signoz(
delete=delete,
restore=restore,
rebuild=pytestconfig.getoption("--rebuild"),
stale=stale,
)
@@ -222,6 +243,7 @@ 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:
@@ -233,4 +255,5 @@ 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 Normal file
View File

@@ -0,0 +1,143 @@
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,6 +302,7 @@ class Traces(ABC):
db_operation: str
has_error: bool
is_remote: str
scope_json: dict[str, Any]
resource: list[TracesResource]
tag_attributes: list[TracesTagAttributes]
@@ -327,6 +328,7 @@ 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:
@@ -408,6 +410,33 @@ 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 = {}
@@ -659,6 +688,7 @@ class Traces(ABC):
self.has_error,
self.is_remote,
self.resource_json,
self.scope_json,
],
dtype=object,
)
@@ -689,6 +719,7 @@ class Traces(ABC):
attributes=data.get("attributes", {}),
trace_state=data.get("trace_state", ""),
flags=data.get("flags", 0),
scope=data.get("scope", {}),
)
@classmethod
@@ -828,6 +859,7 @@ 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,6 +158,26 @@ 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,6 +37,7 @@ 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

@@ -0,0 +1,185 @@
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,6 +130,45 @@ 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,6 +1240,13 @@ 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(
@@ -1283,6 +1290,156 @@ 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,6 +19,8 @@ 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,8 +1070,10 @@ version = "0.1.0"
source = { virtual = "." }
dependencies = [
{ name = "clickhouse-connect" },
{ name = "cryptography" },
{ name = "fastapi" },
{ name = "isodate" },
{ name = "jwcrypto" },
{ name = "numpy" },
{ name = "psycopg2" },
{ name = "py" },
@@ -1093,8 +1095,10 @@ 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" },