Compare commits

...

45 Commits

Author SHA1 Message Date
Nikhil Soni
a09ace42fb chore: regenerate openapi spec and api clients
Assisted-by: Claude Opus 5.5
2026-10-06 22:11:15 +05:30
Nikhil Soni
62f1e0cebc chore(promote): trim comments and the materialized route description
Assisted-by: Claude Opus 5.5
2026-10-06 22:11:14 +05:30
Nikhil Soni
120c953caa refactor(promote): check the materialized path type in the module
Assisted-by: Claude Opus 5.5
2026-10-06 22:03:38 +05:30
Nikhil Soni
c689834ba5 refactor(promote): skip default and indexed materialized paths in the module
The default index reuses the /api/v2/traces/fields constants so both stay in sync.

Assisted-by: Claude Opus 5.5
2026-10-06 21:53:15 +05:30
Nikhil Soni
fd6e2e6061 chore: regenerate openapi spec and api clients
Assisted-by: Claude Opus 5.5
2026-10-06 19:52:02 +05:30
Nikhil Soni
a10586906f refactor(promote): index materialized attributes with the default bloom filter
The materialize field API only builds bloom_filter indexes, so column skip indexes are no longer
read. Resource and scope targets are dropped until needed.

Assisted-by: Claude Opus 5.5
2026-10-06 19:51:55 +05:30
Nikhil Soni
d1ed7426bc chore: regenerate openapi spec and api clients
Assisted-by: Claude Opus 5.5
2026-10-06 18:58:57 +05:30
Nikhil Soni
7ee745d4ce feat(promote): index materialized paths on their json sub-columns
Only string paths are indexed: number reads are dynamicType-gated casts
that no sub-column index matches. Resource and scope are index-only domains.

Assisted-by: Claude Opus 5.5
2026-10-06 18:58:51 +05:30
Nikhil Soni
82a4e395cb chore: regenerate openapi spec and api clients
Assisted-by: Claude Opus 5.5
2026-10-06 14:17:13 +05:30
Nikhil Soni
9eaea370e6 refactor(promote): group imports and document per-signal scopes
Assisted-by: Claude Opus 5.5
2026-10-06 14:16:31 +05:30
Nikhil Soni
75316dbd52 refactor(promote): move the field resource onto the signal
Assisted-by: Claude Opus 5.5
2026-10-05 21:57:47 +05:30
Nikhil Soni
b14d5663a3 chore(promote): renumber the field tuples migration after main's
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-05 21:57:47 +05:30
Nikhil Soni
77155f79c1 refactor(promote): inline the field scopes and drop types from the migration
Assisted-by: Claude Opus 5.5
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-05 21:57:47 +05:30
Nikhil Soni
dbcfbf8012 chore: regenerate openapi spec
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-05 21:57:46 +05:30
Nikhil Soni
9e6b8f832d feat(promote): authorize promoted paths per signal
Reuse the existing logs-field and traces-field resources instead of a new kind.

Assisted-by: Claude Opus 5.5
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-05 21:57:46 +05:30
Nikhil Soni
4a30d464d2 Revert "feat(promote): support the bloom_filter index type"
This reverts commit 8c22b1386267beefcb09ddc04dedd3a5f9886784.
2026-10-05 21:56:43 +05:30
Nikhil Soni
974fe05ba3 feat(promote): support the bloom_filter index type
IndexNameType maps a requested index type to the bare name used in
index names; ValidateAndSetDefaults rejects unsupported types with it
and index creation reuses the same mapping.
2026-10-05 21:56:43 +05:30
Nikhil Soni
55b541a268 Revert "refactor(promote): pass resolved targets into PromotePaths"
This reverts commit df7d64ba0d94f3e40ea834a868d9d7728826b69d.
2026-10-05 21:56:43 +05:30
Nikhil Soni
9d3f1db01f refactor(promote): pass resolved targets into PromotePaths
The handler resolves each path's target for validation; PromotePaths
now takes the pairs as TargetedPath instead of re-resolving them.
2026-10-05 21:56:43 +05:30
Nikhil Soni
a69bc741fb refactor(promote): validate promote paths in the handler
The handler resolves each path's target and validates it, matching the
listing filters pattern; the module receives validated paths. The
exported RejectedPathPrefixes collapses into the unexported
reservedPathPrefix predicate the validation folds in.
2026-10-05 21:56:43 +05:30
Nikhil Soni
26ad334fdb refactor(promote): rename JSONIndexSource to JSONIndexLookup
Rename the metadata store parameter struct to say what it does, flatten
RejectedPathPrefixes (the signal switch only added the logs body search
prefix), and move the unexported expression helper to the end of
target.go.
2026-10-05 21:56:43 +05:30
Nikhil Soni
204d7384b3 refactor(promote): drop the IndexesSupported target flag 2026-10-05 21:56:43 +05:30
Nikhil Soni
e7c49bfce7 Revert "refactor(promote): move index creation into the metadata store"
This reverts commit b137ffd95b.
2026-10-05 21:56:43 +05:30
Nikhil Soni
94de5cf72b refactor(promote): move index creation into the metadata store 2026-10-05 21:56:43 +05:30
Nikhil Soni
6f24594cf7 feat(promote): per-path indexes for trace attributes and bare paths in every domain
- the traces attributes target now supports per-path skip indexes; their
  expression is a bare type cast, CAST(dynamicElement(col.path, 'T'), 'T'),
  with no lower/assumeNotNull folding. The cast also unwraps the Nullable
  that dynamicElement returns, which bloom filter indexes reject.
- the body. path prefix is dropped for logs body: the API URL already
  names the context, so paths are bare attribute names in every domain;
  prefixed paths are rejected with a guiding error.
- ListLogsJSONIndexes generalizes to ListJSONIndexes(source) driven by the
  target, and the index expression unfolding accepts both the folded and
  the bare cast forms.
2026-10-05 21:56:43 +05:30
Nikhil Soni
8e607aa6d7 refactor(promote): drop comments that restate the code 2026-10-05 21:56:40 +05:30
Nikhil Soni
e35ad7b1cd refactor(promote): validate the listing filters in the handler 2026-10-05 21:56:40 +05:30
Nikhil Soni
f59cbd85bb chore: regenerate openapi spec and api clients 2026-10-05 21:56:40 +05:30
Nikhil Soni
3498eb00b1 feat(promote): add filters to the promoted paths listing
GET /api/v1/promoted_path accepts signal, context, promoted and indexes
query parameters narrowing the listing. The signal, context and path
fields of PromotePath are now marked required in the API contract.
2026-10-05 21:56:40 +05:30
Nikhil Soni
33b6d629b2 chore: regenerate openapi spec and api clients 2026-10-05 21:56:40 +05:30
Nikhil Soni
fc1d0cc504 refactor(promote)!: move the promotion domain from the URL into the request
The promote paths API collapses to /api/v1/promoted_path. Each PromotePath
carries its signal and context, so a single request can span domains and
the list endpoint returns every domain's paths annotated with theirs.
2026-10-05 21:56:40 +05:30
Nikhil Soni
a6ab2f3583 refactor(promote)!: rename the API path to promoted_path
A noun and singular, like the other API URLs.
2026-10-05 21:56:40 +05:30
Nikhil Soni
1525a334c2 refactor(promote): group target constructors; drop redundant and unsupported-feature tests
- move NewTargetFromPath next to the other target constructors
- drop TestNewTargetFromPath: thin glue over SignalFromText/FieldContextFromText/TargetFor
- drop traces index rejection cases: per-path indexes are simply not supported for traces yet
2026-10-05 21:56:40 +05:30
Nikhil Soni
378717a3c0 test(promote): cover the per-path skip index creation of the logs body domain 2026-10-05 21:56:40 +05:30
Nikhil Soni
41427714fe refactor(promote): move the path resolution to types with a validate method, table-drive the tests 2026-10-05 21:56:40 +05:30
Nikhil Soni
a329c0fef9 chore: regenerate openapi spec and api clients 2026-10-05 21:56:40 +05:30
Nikhil Soni
cadad6d61d test: align the subtest names with the table format rule 2026-10-05 21:56:40 +05:30
Nikhil Soni
d0eb63073e fix(promote): rename the signal path variable to telemetry_signal
orval generates an AbortSignal parameter named signal for every client
method, so a {signal} path variable produced a duplicate identifier in
the generated client (tsc error). The URL itself is unchanged in
behavior: /api/v1/promote_paths/{telemetry_signal}/{context}.
2026-10-05 21:56:40 +05:30
Nikhil Soni
d9f2ebbf0e refactor(promote): inline the promote and list helpers into their sole callers 2026-10-05 21:56:40 +05:30
Nikhil Soni
cb683a82b6 refactor(promote)!: drop the legacy logs promote_paths routes
There are no consumers of /api/v1/logs/promote_paths, so no backward
compatibility is needed: the logs body domain is served by the generic
/api/v1/promote_paths/{signal}/{context} routes and the legacy routes
and handler methods are removed.
2026-10-05 21:56:40 +05:30
Nikhil Soni
66fc58054b refactor(promote): move Target into target.go, enum-style SignalFromText, rename handler method
- Target type definition moves from types.go to target.go alongside its
  constructors, with inline comments
- SignalFromText follows the codebase enum pattern (switch over the
  declared values + Enum method) instead of a string-to-signal map
- generic route handler method renamed HandlePromotePaths -> PromotePaths
2026-10-05 21:56:40 +05:30
Nikhil Soni
e482ad78b4 refactor(promote): centralize domain construction and generalize routes
- target construction moves to promotetypes: a generic NewTarget plus
  per-domain constructors (NewLogsBodyTarget, NewTracesAttributesTarget)
  and a TargetFor registry keyed by (signal, context); implpromote and
  telemetrymetadata no longer hand-roll domain literals
- routes generalize to /api/v1/promote_paths/{signal}/{context}: the
  legacy logs body route (/api/v1/logs/promote_paths) is kept for
  compatibility but the domain now travels in the path, so a future logs
  attribute domain does not collide with the logs body route; supersedes
  the /api/v1/traces/promote_paths routes
- add telemetrytypes.SignalFromText for parsing the signal path variable
2026-10-05 21:56:40 +05:30
Nikhil Soni
d26cc928bf refactor(promote): template the promotion record with EvolutionEntry
Target now carries an EvolutionEntry template (signal, promoted column
name and type, field context) instead of loose signal/context/column
fields, so the store write is exactly row template + field names +
release time and the hardcoded JSON() column type moves to the domain
definitions. DBName/LocalTableName stay on Target explicitly as index
DDL config, used only by targets with index support.
2026-10-05 21:56:40 +05:30
Nikhil Soni
17b8b6504b refactor(promote): collapse module interface to target-parameterized methods
The per-domain methods were pure delegates; the promotion domain now
travels as promotetypes.Target through Module.ListPromotedPaths /
Module.PromotePaths, with the handler methods (one per route) passing
their domain's target.
2026-10-05 21:56:40 +05:30
Nikhil Soni
bcd3937c4c feat(promote): add traces attributes promotion API
Refactor the promote module into a target-parameterized core so the logs
body_v2 flow and future promotion domains share one implementation, and
add the spans attributes JSON column (attributes -> attributes_promoted)
as a second domain behind POST/GET /api/v1/traces/promote_paths.

- promotetypes.Target describes a promotion domain: signal, field
  context, db/table, base/promoted columns, path prefix rule and whether
  per-path skip indexes are supported
- index support is optional per target; traces starts promotion-only
  since the traces query builder does not consume per-path skip indexes
- metadata store GetPromotedPaths/PromotePaths take (signal, column,
  context) instead of being hardcoded to the logs body column
- fix the list response never attaching indexes to promoted entries
  (aggregated by unprefixed name but looked up by prefixed path) and
  reporting indexed+promoted paths twice
2026-10-05 21:56:40 +05:30
31 changed files with 2903 additions and 491 deletions

View File

@@ -60,6 +60,7 @@ jobs:
- rawexportdata
- promqlconformance
- promapiconformance
- promote
- querierauthz
- role
- rootuser

View File

@@ -7661,8 +7661,19 @@ components:
- metric
- value
type: object
PromotetypesIndexMaterializedPathsParams:
properties:
context:
type: string
dryRun:
type: boolean
signal:
type: string
type: object
PromotetypesPromotePath:
properties:
context:
type: string
indexes:
items:
$ref: '#/components/schemas/PromotetypesWrappedIndex'
@@ -7671,6 +7682,12 @@ components:
type: string
promote:
type: boolean
signal:
type: string
required:
- signal
- context
- path
type: object
PromotetypesWrappedIndex:
properties:
@@ -13146,110 +13163,6 @@ paths:
tags:
- llmpricingrules
x-signoz-stability: alpha
/api/v1/logs/promote_paths:
get:
deprecated: false
description: This endpoints promotes and indexes paths
operationId: ListPromotedAndIndexedPaths
responses:
"200":
content:
application/json:
schema:
properties:
data:
items:
$ref: '#/components/schemas/PromotetypesPromotePath'
nullable: true
type: array
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: Promote and index paths
tags:
- logs
x-signoz-stability: alpha
post:
deprecated: false
description: This endpoints promotes and indexes paths
operationId: HandlePromoteAndIndexPaths
requestBody:
content:
application/json:
schema:
items:
$ref: '#/components/schemas/PromotetypesPromotePath'
nullable: true
type: array
responses:
"201":
description: Created
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Promote and index paths
tags:
- logs
x-signoz-stability: alpha
/api/v1/org/preferences:
get:
deprecated: false
@@ -13421,6 +13334,215 @@ paths:
tags:
- preferences
x-signoz-stability: alpha
/api/v1/promoted_path:
get:
deprecated: false
description: This endpoint lists the promoted paths of every JSON column, each
annotated with its signal and context. The signal, context, promoted and indexes
query parameters filter the listing. Requires the list scope of each signal
the filters select, or of every signal when none match.
operationId: ListPromotedPaths
parameters:
- in: query
name: signal
schema:
type: string
- in: query
name: context
schema:
type: string
- in: query
name: promoted
schema:
nullable: true
type: boolean
- in: query
name: indexes
schema:
nullable: true
type: boolean
responses:
"200":
content:
application/json:
schema:
properties:
data:
items:
$ref: '#/components/schemas/PromotetypesPromotePath'
nullable: true
type: array
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- logs-field:list
- traces-field:list
- tokenizer:
- logs-field:list
- traces-field:list
summary: List promoted paths
tags:
- promote
x-signoz-stability: alpha
post:
deprecated: false
description: This endpoint promotes paths of JSON columns to their promoted
columns. Each path names its promotion domain with its signal and context,
e.g. traces/attribute. Requires the update scope of each signal in the request.
operationId: PromotePaths
requestBody:
content:
application/json:
schema:
items:
$ref: '#/components/schemas/PromotetypesPromotePath'
nullable: true
type: array
responses:
"201":
description: Created
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- logs-field:update
- traces-field:update
- tokenizer:
- logs-field:update
- traces-field:update
summary: Promote paths
tags:
- promote
x-signoz-stability: alpha
/api/v1/promoted_path/materialized:
post:
deprecated: false
description: This endpoint indexes the JSON sub-columns of materialized paths
and returns them. Requires the update scope of each signal the filters select.
operationId: IndexMaterializedPaths
parameters:
- in: query
name: signal
schema:
type: string
- in: query
name: context
schema:
type: string
- in: query
name: dryRun
schema:
type: boolean
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/PromotetypesIndexMaterializedPathsParams'
responses:
"200":
content:
application/json:
schema:
properties:
data:
items:
$ref: '#/components/schemas/PromotetypesPromotePath'
nullable: true
type: array
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- logs-field:update
- traces-field:update
- tokenizer:
- logs-field:update
- traces-field:update
summary: Index materialized paths
tags:
- promote
x-signoz-stability: alpha
/api/v1/roles:
get:
deprecated: false

View File

@@ -4,23 +4,15 @@
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
import { useMutation, useQuery } from 'react-query';
import { useMutation } from 'react-query';
import type {
InvalidateOptions,
MutationFunction,
QueryClient,
QueryFunction,
QueryKey,
UseMutationOptions,
UseMutationResult,
UseQueryOptions,
UseQueryResult,
} from 'react-query';
import type {
HandleExportRawDataPOSTParams,
ListPromotedAndIndexedPaths200,
PromotetypesPromotePathDTO,
Querybuildertypesv5QueryRangeRequestDTO,
RenderErrorResponseDTO,
} from '../sigNoz.schemas';
@@ -28,26 +20,6 @@ import type {
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
const withQueryKey = <T extends object, K>(
query: T,
queryKey: K,
): T & { queryKey: K } => {
const result = { queryKey } as T & { queryKey: K };
for (const key of Object.keys(query)) {
// The explicit queryKey always wins, matching the previous
// `{ ...query, queryKey }` spread where it was set last.
if (key === 'queryKey') {
continue;
}
Object.defineProperty(result, key, {
enumerable: true,
configurable: true,
get: () => (query as Record<string, unknown>)[key],
});
}
return result;
};
/**
* This endpoints allows complex query exporting raw data for traces and logs
* @summary Export raw data
@@ -149,175 +121,3 @@ export const useHandleExportRawDataPOST = <
> => {
return useMutation(getHandleExportRawDataPOSTMutationOptions(options));
};
/**
* This endpoints promotes and indexes paths
* @summary Promote and index paths
*/
export const listPromotedAndIndexedPaths = (signal?: AbortSignal) => {
return GeneratedAPIInstance<ListPromotedAndIndexedPaths200>({
url: `/api/v1/logs/promote_paths`,
method: 'GET',
signal,
});
};
export const getListPromotedAndIndexedPathsQueryKey = () => {
return [`/api/v1/logs/promote_paths`] as const;
};
export const getListPromotedAndIndexedPathsQueryOptions = <
TData = Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
TError,
TData
>;
}) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getListPromotedAndIndexedPathsQueryKey();
const queryFn: QueryFunction<
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>
> = ({ signal }) => listPromotedAndIndexedPaths(signal);
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type ListPromotedAndIndexedPathsQueryResult = NonNullable<
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>
>;
export type ListPromotedAndIndexedPathsQueryError =
ErrorType<RenderErrorResponseDTO>;
/**
* @summary Promote and index paths
*/
export function useListPromotedAndIndexedPaths<
TData = Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
TError,
TData
>;
}): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getListPromotedAndIndexedPathsQueryOptions(options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return withQueryKey(query, queryOptions.queryKey);
}
/**
* @summary Promote and index paths
*/
export const invalidateListPromotedAndIndexedPaths = async (
queryClient: QueryClient,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getListPromotedAndIndexedPathsQueryKey() },
options,
);
return queryClient;
};
/**
* This endpoints promotes and indexes paths
* @summary Promote and index paths
*/
export const handlePromoteAndIndexPaths = (
promotetypesPromotePathDTONull?: BodyType<
PromotetypesPromotePathDTO[] | null
> | null,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v1/logs/promote_paths`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: promotetypesPromotePathDTONull,
signal,
});
};
export const getHandlePromoteAndIndexPathsMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
> => {
const mutationKey = ['handlePromoteAndIndexPaths'];
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 handlePromoteAndIndexPaths>>,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> }
> = (props) => {
const { data } = props ?? {};
return handlePromoteAndIndexPaths(data);
};
return { mutationFn, ...mutationOptions };
};
export type HandlePromoteAndIndexPathsMutationResult = NonNullable<
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>
>;
export type HandlePromoteAndIndexPathsMutationBody =
| BodyType<PromotetypesPromotePathDTO[] | null>
| undefined;
export type HandlePromoteAndIndexPathsMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @summary Promote and index paths
*/
export const useHandlePromoteAndIndexPaths = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
> => {
return useMutation(getHandlePromoteAndIndexPathsMutationOptions(options));
};

View File

@@ -0,0 +1,336 @@
/**
* ! Do not edit manually
* * The file has been auto-generated using Orval for SigNoz
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
import { useMutation, useQuery } from 'react-query';
import type {
InvalidateOptions,
MutationFunction,
QueryClient,
QueryFunction,
QueryKey,
UseMutationOptions,
UseMutationResult,
UseQueryOptions,
UseQueryResult,
} from 'react-query';
import type {
IndexMaterializedPaths200,
IndexMaterializedPathsParams,
ListPromotedPaths200,
ListPromotedPathsParams,
PromotetypesIndexMaterializedPathsParamsDTO,
PromotetypesPromotePathDTO,
RenderErrorResponseDTO,
} from '../sigNoz.schemas';
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
const withQueryKey = <T extends object, K>(
query: T,
queryKey: K,
): T & { queryKey: K } => {
const result = { queryKey } as T & { queryKey: K };
for (const key of Object.keys(query)) {
// The explicit queryKey always wins, matching the previous
// `{ ...query, queryKey }` spread where it was set last.
if (key === 'queryKey') {
continue;
}
Object.defineProperty(result, key, {
enumerable: true,
configurable: true,
get: () => (query as Record<string, unknown>)[key],
});
}
return result;
};
/**
* This endpoint lists the promoted paths of every JSON column, each annotated with its signal and context. The signal, context, promoted and indexes query parameters filter the listing. Requires the list scope of each signal the filters select, or of every signal when none match.
* @summary List promoted paths
*/
export const listPromotedPaths = (
params?: ListPromotedPathsParams,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<ListPromotedPaths200>({
url: `/api/v1/promoted_path`,
method: 'GET',
params,
signal,
});
};
export const getListPromotedPathsQueryKey = (
params?: ListPromotedPathsParams,
) => {
return [`/api/v1/promoted_path`, ...(params ? [params] : [])] as const;
};
export const getListPromotedPathsQueryOptions = <
TData = Awaited<ReturnType<typeof listPromotedPaths>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
params?: ListPromotedPathsParams,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listPromotedPaths>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getListPromotedPathsQueryKey(params);
const queryFn: QueryFunction<
Awaited<ReturnType<typeof listPromotedPaths>>
> = ({ signal }) => listPromotedPaths(params, signal);
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
Awaited<ReturnType<typeof listPromotedPaths>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type ListPromotedPathsQueryResult = NonNullable<
Awaited<ReturnType<typeof listPromotedPaths>>
>;
export type ListPromotedPathsQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary List promoted paths
*/
export function useListPromotedPaths<
TData = Awaited<ReturnType<typeof listPromotedPaths>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
params?: ListPromotedPathsParams,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listPromotedPaths>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getListPromotedPathsQueryOptions(params, options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return withQueryKey(query, queryOptions.queryKey);
}
/**
* @summary List promoted paths
*/
export const invalidateListPromotedPaths = async (
queryClient: QueryClient,
params?: ListPromotedPathsParams,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getListPromotedPathsQueryKey(params) },
options,
);
return queryClient;
};
/**
* This endpoint promotes paths of JSON columns to their promoted columns. Each path names its promotion domain with its signal and context, e.g. traces/attribute. Requires the update scope of each signal in the request.
* @summary Promote paths
*/
export const promotePaths = (
promotetypesPromotePathDTONull?: BodyType<
PromotetypesPromotePathDTO[] | null
> | null,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v1/promoted_path`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: promotetypesPromotePathDTONull,
signal,
});
};
export const getPromotePathsMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof promotePaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof promotePaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
> => {
const mutationKey = ['promotePaths'];
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 promotePaths>>,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> }
> = (props) => {
const { data } = props ?? {};
return promotePaths(data);
};
return { mutationFn, ...mutationOptions };
};
export type PromotePathsMutationResult = NonNullable<
Awaited<ReturnType<typeof promotePaths>>
>;
export type PromotePathsMutationBody =
| BodyType<PromotetypesPromotePathDTO[] | null>
| undefined;
export type PromotePathsMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Promote paths
*/
export const usePromotePaths = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof promotePaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof promotePaths>>,
TError,
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
TContext
> => {
return useMutation(getPromotePathsMutationOptions(options));
};
/**
* This endpoint indexes the JSON sub-columns of materialized paths and returns them. Requires the update scope of each signal the filters select.
* @summary Index materialized paths
*/
export const indexMaterializedPaths = (
promotetypesIndexMaterializedPathsParamsDTO?: BodyType<PromotetypesIndexMaterializedPathsParamsDTO>,
params?: IndexMaterializedPathsParams,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<IndexMaterializedPaths200>({
url: `/api/v1/promoted_path/materialized`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: promotetypesIndexMaterializedPathsParamsDTO,
params,
signal,
});
};
export const getIndexMaterializedPathsMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof indexMaterializedPaths>>,
TError,
{
data?: BodyType<PromotetypesIndexMaterializedPathsParamsDTO>;
params?: IndexMaterializedPathsParams;
},
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof indexMaterializedPaths>>,
TError,
{
data?: BodyType<PromotetypesIndexMaterializedPathsParamsDTO>;
params?: IndexMaterializedPathsParams;
},
TContext
> => {
const mutationKey = ['indexMaterializedPaths'];
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 indexMaterializedPaths>>,
{
data?: BodyType<PromotetypesIndexMaterializedPathsParamsDTO>;
params?: IndexMaterializedPathsParams;
}
> = (props) => {
const { data, params } = props ?? {};
return indexMaterializedPaths(data, params);
};
return { mutationFn, ...mutationOptions };
};
export type IndexMaterializedPathsMutationResult = NonNullable<
Awaited<ReturnType<typeof indexMaterializedPaths>>
>;
export type IndexMaterializedPathsMutationBody =
| BodyType<PromotetypesIndexMaterializedPathsParamsDTO>
| undefined;
export type IndexMaterializedPathsMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @summary Index materialized paths
*/
export const useIndexMaterializedPaths = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof indexMaterializedPaths>>,
TError,
{
data?: BodyType<PromotetypesIndexMaterializedPathsParamsDTO>;
params?: IndexMaterializedPathsParams;
},
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof indexMaterializedPaths>>,
TError,
{
data?: BodyType<PromotetypesIndexMaterializedPathsParamsDTO>;
params?: IndexMaterializedPathsParams;
},
TContext
> => {
return useMutation(getIndexMaterializedPathsMutationOptions(options));
};

View File

@@ -9423,6 +9423,21 @@ export interface PrometheusSuccessResponseSchemaDTO {
warnings?: string[];
}
export interface PromotetypesIndexMaterializedPathsParamsDTO {
/**
* @type string
*/
context?: string;
/**
* @type boolean
*/
dryRun?: boolean;
/**
* @type string
*/
signal?: string;
}
export interface PromotetypesWrappedIndexDTO {
fieldDataType?: TelemetrytypesFieldDataTypeDTO;
/**
@@ -9436,6 +9451,10 @@ export interface PromotetypesWrappedIndexDTO {
}
export interface PromotetypesPromotePathDTO {
/**
* @type string
*/
context: string;
/**
* @type array
*/
@@ -9443,11 +9462,15 @@ export interface PromotetypesPromotePathDTO {
/**
* @type string
*/
path?: string;
path: string;
/**
* @type boolean
*/
promote?: boolean;
/**
* @type string
*/
signal: string;
}
export interface Querybuildertypesv5AggregationMetaDTO {
@@ -12537,17 +12560,6 @@ export type ListUnmappedLLMModels200 = {
status: string;
};
export type ListPromotedAndIndexedPaths200 = {
/**
* @type array,null
*/
data: PromotetypesPromotePathDTO[] | null;
/**
* @type string
*/
status: string;
};
export type ListOrgPreferences200 = {
/**
* @type array
@@ -12573,6 +12585,69 @@ export type GetOrgPreference200 = {
export type UpdateOrgPreferencePathParameters = {
name: string;
};
export type ListPromotedPathsParams = {
/**
* @type string
* @description undefined
*/
signal?: string;
/**
* @type string
* @description undefined
*/
context?: string;
/**
* @type boolean,null
* @description undefined
*/
promoted?: boolean | null;
/**
* @type boolean,null
* @description undefined
*/
indexes?: boolean | null;
};
export type ListPromotedPaths200 = {
/**
* @type array,null
*/
data: PromotetypesPromotePathDTO[] | null;
/**
* @type string
*/
status: string;
};
export type IndexMaterializedPathsParams = {
/**
* @type string
* @description undefined
*/
signal?: string;
/**
* @type string
* @description undefined
*/
context?: string;
/**
* @type boolean
* @description undefined
*/
dryRun?: boolean;
};
export type IndexMaterializedPaths200 = {
/**
* @type array,null
*/
data: PromotetypesPromotePathDTO[] | null;
/**
* @type string
*/
status: string;
};
export type ListRoles200 = {
/**
* @type array

View File

@@ -3,42 +3,77 @@ package signozapiserver
import (
"net/http"
"github.com/SigNoz/signoz/pkg/http/handler"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/types/promotetypes"
"github.com/gorilla/mux"
"github.com/SigNoz/signoz/pkg/http/handler"
"github.com/SigNoz/signoz/pkg/types/authtypes"
"github.com/SigNoz/signoz/pkg/types/coretypes"
"github.com/SigNoz/signoz/pkg/types/promotetypes"
)
func (provider *provider) addPromoteRoutes(router *mux.Router) error {
if err := router.Handle("/api/v1/logs/promote_paths", handler.New(provider.authzMiddleware.EditAccess(provider.promoteHandler.HandlePromoteAndIndexPaths), handler.OpenAPIDef{
ID: "HandlePromoteAndIndexPaths",
Tags: []string{"logs"},
Summary: "Promote and index paths",
Description: "This endpoints promotes and indexes paths",
if err := router.Handle("/api/v1/promoted_path", handler.New(provider.authzMiddleware.CheckResources(provider.promoteHandler.PromotePaths, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName), handler.OpenAPIDef{
ID: "PromotePaths",
Tags: []string{"promote"},
Summary: "Promote paths",
Description: "This endpoint promotes paths of JSON columns to their promoted columns. Each path names its promotion domain with its signal and context, e.g. traces/attribute. Requires the update scope of each signal in the request.",
Request: new([]*promotetypes.PromotePath),
RequestContentType: "application/json",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest},
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodPost).GetError(); err != nil {
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceLogsField.Scope(coretypes.VerbUpdate), coretypes.ResourceMetaResourceTracesField.Scope(coretypes.VerbUpdate)}),
}, handler.WithResourceDefs(handler.TelemetryResourceDef{
Verb: coretypes.VerbUpdate,
Category: coretypes.ActionCategoryConfigurationChange,
Selector: coretypes.WildcardSelector,
Resources: promotetypes.PromotePathsResources,
}))).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v1/logs/promote_paths", handler.New(provider.authzMiddleware.ViewAccess(provider.promoteHandler.ListPromotedAndIndexedPaths), handler.OpenAPIDef{
ID: "ListPromotedAndIndexedPaths",
Tags: []string{"logs"},
Summary: "Promote and index paths",
Description: "This endpoints promotes and indexes paths",
if err := router.Handle("/api/v1/promoted_path", handler.New(provider.authzMiddleware.CheckResources(provider.promoteHandler.ListPromotedPaths, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName), handler.OpenAPIDef{
ID: "ListPromotedPaths",
Tags: []string{"promote"},
Summary: "List promoted paths",
Description: "This endpoint lists the promoted paths of every JSON column, each annotated with its signal and context. The signal, context, promoted and indexes query parameters filter the listing. Requires the list scope of each signal the filters select, or of every signal when none match.",
Request: nil,
RequestQuery: new(promotetypes.ListPromotedPathsFilters),
RequestContentType: "",
Response: new([]*promotetypes.PromotePath),
ResponseContentType: "",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
})).Methods(http.MethodGet).GetError(); err != nil {
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceLogsField.Scope(coretypes.VerbList), coretypes.ResourceMetaResourceTracesField.Scope(coretypes.VerbList)}),
}, handler.WithResourceDefs(handler.TelemetryResourceDef{
Verb: coretypes.VerbList,
Category: coretypes.ActionCategoryDataAccess,
Selector: coretypes.WildcardSelector,
Resources: promotetypes.FilteredSignalsResources,
}))).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v1/promoted_path/materialized", handler.New(provider.authzMiddleware.CheckResources(provider.promoteHandler.IndexMaterializedPaths, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName), handler.OpenAPIDef{
ID: "IndexMaterializedPaths",
Tags: []string{"promote"},
Summary: "Index materialized paths",
Description: "This endpoint indexes the JSON sub-columns of materialized paths and returns them. Requires the update scope of each signal the filters select.",
Request: nil,
RequestQuery: new(promotetypes.IndexMaterializedPathsParams),
RequestContentType: "",
Response: new([]*promotetypes.PromotePath),
ResponseContentType: "",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest},
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceLogsField.Scope(coretypes.VerbUpdate), coretypes.ResourceMetaResourceTracesField.Scope(coretypes.VerbUpdate)}),
}, handler.WithResourceDefs(handler.TelemetryResourceDef{
Verb: coretypes.VerbUpdate,
Category: coretypes.ActionCategoryConfigurationChange,
Selector: coretypes.WildcardSelector,
Resources: promotetypes.FilteredSignalsResources,
}))).Methods(http.MethodPost).GetError(); err != nil {
return err
}

View File

@@ -19,7 +19,7 @@ func NewHandler(module promote.Module) promote.Handler {
return &handler{module: module}
}
func (h *handler) HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Request) {
func (h *handler) PromotePaths(w http.ResponseWriter, r *http.Request) {
// TODO(Nitya): Use in multi tenant setup
_, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
@@ -32,8 +32,23 @@ func (h *handler) HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Requ
render.Error(w, err)
return
}
if len(req) == 0 {
render.Error(w, errors.NewInvalidInputf(errors.CodeInvalidInput, "paths cannot be empty"))
return
}
for _, path := range req {
target, err := path.Target()
if err != nil {
render.Error(w, err)
return
}
if err := path.ValidateAndSetDefaults(target); err != nil {
render.Error(w, err)
return
}
}
err = h.module.PromoteAndIndexPaths(r.Context(), req...)
err = h.module.PromotePaths(r.Context(), req...)
if err != nil {
render.Error(w, err)
return
@@ -42,7 +57,7 @@ func (h *handler) HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Requ
render.Success(w, http.StatusCreated, nil)
}
func (h *handler) ListPromotedAndIndexedPaths(w http.ResponseWriter, r *http.Request) {
func (h *handler) ListPromotedPaths(w http.ResponseWriter, r *http.Request) {
// TODO(Nitya): Use in multi tenant setup
_, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
@@ -50,7 +65,38 @@ func (h *handler) ListPromotedAndIndexedPaths(w http.ResponseWriter, r *http.Req
return
}
paths, err := h.module.ListPromotedAndIndexedPaths(r.Context())
var filters promotetypes.ListPromotedPathsFilters
if err := binding.Query.BindQuery(r.URL.Query(), &filters); err != nil {
render.Error(w, err)
return
}
if err := filters.Validate(); err != nil {
render.Error(w, err)
return
}
paths, err := h.module.ListPromotedPaths(r.Context(), filters)
if err != nil {
render.Error(w, err)
return
}
render.Success(w, http.StatusOK, paths)
}
func (h *handler) IndexMaterializedPaths(w http.ResponseWriter, r *http.Request) {
var params promotetypes.IndexMaterializedPathsParams
if err := binding.Query.BindQuery(r.URL.Query(), &params); err != nil {
render.Error(w, err)
return
}
filters := params.Filters()
if err := filters.Validate(); err != nil {
render.Error(w, err)
return
}
paths, err := h.module.IndexMaterializedPaths(r.Context(), params)
if err != nil {
render.Error(w, err)
return

View File

@@ -2,14 +2,11 @@ package implpromote
import (
"context"
"maps"
"slices"
"strings"
schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/modules/promote"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
@@ -31,46 +28,72 @@ func NewModule(metadataStore telemetrytypes.MetadataStore, telemetrystore teleme
return &module{metadataStore: metadataStore, telemetryStore: telemetrystore}
}
func (m *module) ListPromotedAndIndexedPaths(ctx context.Context) ([]promotetypes.PromotePath, error) {
indexes, err := m.metadataStore.ListLogsJSONIndexes(ctx)
func (m *module) ListPromotedPaths(ctx context.Context, filters promotetypes.ListPromotedPathsFilters) ([]promotetypes.PromotePath, error) {
response := make([]promotetypes.PromotePath, 0)
for _, target := range promotetypes.Targets() {
if !filters.MatchesTarget(target) {
continue
}
paths, err := m.listPromotedPaths(ctx, target)
if err != nil {
return nil, err
}
for _, path := range paths {
if filters.MatchesPath(path) {
response = append(response, path)
}
}
}
return response, nil
}
func (m *module) listPromotedPaths(ctx context.Context, target promotetypes.Target) ([]promotetypes.PromotePath, error) {
promotedPaths, err := m.metadataStore.GetPromotedPaths(ctx, target.Entry)
if err != nil {
return nil, err
}
response := make([]promotetypes.PromotePath, 0, len(promotedPaths))
for path := range promotedPaths {
response = append(response, promotetypes.PromotePath{
Signal: target.Entry.Signal.StringValue(),
Context: target.Entry.FieldContext.StringValue(),
Path: path,
Promote: true,
})
}
indexes, err := m.metadataStore.ListJSONIndexes(ctx, target.JSONIndexLookup())
if err != nil {
return nil, err
}
// aggr keys are full sub-column paths: index.BaseColumn carries the
// column prefix and index.Name the bare path.
aggr := map[string][]promotetypes.WrappedIndex{}
for _, index := range indexes {
aggr[index.Name] = append(aggr[index.Name], promotetypes.WrappedIndex{
fullPath := index.BaseColumn + index.Name
aggr[fullPath] = append(aggr[fullPath], promotetypes.WrappedIndex{
FieldDataType: index.FieldDataType,
Type: index.IndexType,
Granularity: index.Granularity,
})
}
promotedPaths, err := m.listPromotedPaths(ctx)
if err != nil {
return nil, err
}
response := []promotetypes.PromotePath{}
for _, path := range promotedPaths {
fullPath := logstelemetryschema.BodyPromotedColumnPrefix + path
path = telemetrytypes.BodyJSONStringSearchPrefix + path
item := promotetypes.PromotePath{
Path: path,
Promote: true,
}
indexes, ok := aggr[fullPath]
if ok {
item.Indexes = indexes
for i := range response {
fullPath := target.PromotedColumnPrefix() + response[i].Path
if indexes, ok := aggr[fullPath]; ok {
response[i].Indexes = indexes
delete(aggr, fullPath)
}
response = append(response, item)
}
// add the paths that are not promoted but have indexes
for path, indexes := range aggr {
path := strings.TrimPrefix(path, logstelemetryschema.BodyV2ColumnPrefix)
path = telemetrytypes.BodyJSONStringSearchPrefix + path
for fullPath, indexes := range aggr {
path := strings.TrimPrefix(fullPath, target.BaseColumnPrefix())
path = strings.TrimPrefix(path, target.PromotedColumnPrefix())
response = append(response, promotetypes.PromotePath{
Signal: target.Entry.Signal.StringValue(),
Context: target.Entry.FieldContext.StringValue(),
Path: path,
Indexes: indexes,
})
@@ -78,68 +101,94 @@ func (m *module) ListPromotedAndIndexedPaths(ctx context.Context) ([]promotetype
return response, nil
}
func (m *module) listPromotedPaths(ctx context.Context) ([]string, error) {
paths, err := m.metadataStore.GetPromotedPaths(ctx)
if err != nil {
func (m *module) IndexMaterializedPaths(ctx context.Context, params promotetypes.IndexMaterializedPathsParams) ([]*promotetypes.PromotePath, error) {
filters := params.Filters()
keysBySignal := map[string][]*telemetrytypes.TelemetryFieldKey{}
response := []*promotetypes.PromotePath{}
for _, target := range promotetypes.Targets() {
if !filters.MatchesTarget(target) {
continue
}
keys, ok := keysBySignal[target.Entry.Signal.StringValue()]
if !ok {
var err error
keys, err = m.metadataStore.GetMaterializedKeys(ctx, target.Entry.Signal)
if err != nil {
return nil, err
}
keysBySignal[target.Entry.Signal.StringValue()] = keys
}
paths := []*promotetypes.PromotePath{}
names := []string{}
for _, key := range keys {
if key.FieldContext != target.Entry.FieldContext || key.FieldDataType != telemetrytypes.FieldDataTypeString || target.IsDefaultMaterialized(key.Name) {
continue
}
if path, ok := promotetypes.NewMaterializedPromotePath(target, key.Name); ok {
paths = append(paths, path)
names = append(names, path.Path)
}
}
if len(paths) == 0 {
continue
}
existing, err := m.metadataStore.ListJSONIndexes(ctx, target.JSONIndexLookup(), names...)
if err != nil {
return nil, err
}
existingByPath := map[string][]telemetrytypes.TelemetryFieldKeySkipIndex{}
for _, index := range existing {
existingByPath[index.Name] = append(existingByPath[index.Name], index)
}
for _, path := range paths {
path.SkipExistingIndexes(existingByPath[path.Path])
if len(path.Indexes) > 0 {
response = append(response, path)
}
}
}
if params.DryRun || len(response) == 0 {
return response, nil
}
if err := m.PromotePaths(ctx, response...); err != nil {
return nil, err
}
return slices.Collect(maps.Keys(paths)), nil
return response, nil
}
// PromotePaths inserts provided JSON paths into the promoted paths table for logs queries.
func (m *module) PromotePaths(ctx context.Context, paths []string) error {
if len(paths) == 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "paths cannot be empty")
}
return m.metadataStore.PromotePaths(ctx, paths...)
}
// createIndexes creates string ngram + token filter indexes on JSON path subcolumns for LIKE queries.
func (m *module) createIndexes(ctx context.Context, indexes []schemamigrator.Index) error {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalLogs.StringValue(),
instrumentationtypes.CodeNamespace: "promote",
instrumentationtypes.CodeFunctionName: "createIndexes",
})
if len(indexes) == 0 {
return nil
}
for _, index := range indexes {
alterStmt := schemamigrator.AlterTableAddIndex{
Database: logstelemetryschema.DBName,
Table: logstelemetryschema.LogsV2LocalTableName,
Index: index,
func (m *module) PromotePaths(ctx context.Context, paths ...*promotetypes.PromotePath) error {
byTarget := map[promotetypes.Target][]*promotetypes.PromotePath{}
targets := []promotetypes.Target{}
for _, path := range paths {
target, err := path.Target()
if err != nil {
return err
}
op := alterStmt.OnCluster(m.telemetryStore.Cluster())
if err := m.telemetryStore.ClickhouseDB().Exec(ctx, op.ToSQL()); err != nil {
return errors.WrapInternalf(err, CodeFailedToCreateIndex, "failed to create index")
if _, ok := byTarget[target]; !ok {
targets = append(targets, target)
}
byTarget[target] = append(byTarget[target], path)
}
for _, target := range targets {
if err := m.promotePaths(ctx, target, byTarget[target]...); err != nil {
return err
}
}
return nil
}
// PromoteAndIndexPaths handles promoting paths and creating indexes in one call.
func (m *module) PromoteAndIndexPaths(
ctx context.Context,
paths ...*promotetypes.PromotePath,
) error {
if len(paths) == 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "paths cannot be empty")
}
func (m *module) promotePaths(ctx context.Context, target promotetypes.Target, paths ...*promotetypes.PromotePath) error {
pathsStr := []string{}
// validate the paths
for _, path := range paths {
if err := path.ValidateAndSetDefaults(); err != nil {
return err
}
pathsStr = append(pathsStr, path.Path)
}
existingPromotedPaths, err := m.metadataStore.GetPromotedPaths(ctx, pathsStr...)
existingPromotedPaths, err := m.metadataStore.GetPromotedPaths(ctx, target.Entry, pathsStr...)
if err != nil {
return err
}
@@ -153,27 +202,20 @@ func (m *module) PromoteAndIndexPaths(
}
}
if len(it.Indexes) > 0 {
parentColumn := logstelemetryschema.LogsV2BodyV2Column
parentColumn := target.BaseColumn
// if the path is already promoted or is being promoted, add it to the promoted column
if _, promoted := existingPromotedPaths[it.Path]; promoted || it.Promote {
parentColumn = logstelemetryschema.LogsV2BodyPromotedColumn
parentColumn = target.PromotedColumn()
}
for _, index := range it.Indexes {
var typeIndex schemamigrator.IndexType
switch {
case strings.HasPrefix(index.Type, string(schemamigrator.IndexTypeNGramBF)):
typeIndex = schemamigrator.IndexTypeNGramBF
case strings.HasPrefix(index.Type, string(schemamigrator.IndexTypeTokenBF)):
typeIndex = schemamigrator.IndexTypeTokenBF
case strings.HasPrefix(index.Type, string(schemamigrator.IndexTypeMinMax)):
typeIndex = schemamigrator.IndexTypeMinMax
default:
typeIndex, ok := index.SkipIndexType()
if !ok {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid index type: %s", index.Type)
}
indexes = append(indexes, schemamigrator.Index{
Name: schemamigrator.JSONSubColumnIndexName(parentColumn, it.Path, index.JSONDataType.StringValue(), typeIndex),
Expression: schemamigrator.JSONSubColumnIndexExpr(parentColumn, it.Path, index.JSONDataType.StringValue()),
Expression: target.IndexExpression(parentColumn, it.Path, index.JSONDataType.StringValue()),
Type: index.Type,
Granularity: index.Granularity,
})
@@ -182,17 +224,42 @@ func (m *module) PromoteAndIndexPaths(
}
if len(toInsert) > 0 {
err := m.PromotePaths(ctx, toInsert)
err := m.metadataStore.PromotePaths(ctx, target.Entry, toInsert...)
if err != nil {
return err
}
}
if len(indexes) > 0 {
if err := m.createIndexes(ctx, indexes); err != nil {
if err := m.createIndexes(ctx, target, indexes); err != nil {
return err
}
}
return nil
}
func (m *module) createIndexes(ctx context.Context, target promotetypes.Target, indexes []schemamigrator.Index) error {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: target.Entry.Signal.StringValue(),
instrumentationtypes.CodeNamespace: "promote",
instrumentationtypes.CodeFunctionName: "createIndexes",
})
if len(indexes) == 0 {
return nil
}
for _, index := range indexes {
alterStmt := schemamigrator.AlterTableAddIndex{
Database: target.DBName,
Table: target.LocalTableName,
Index: index,
}
op := alterStmt.OnCluster(m.telemetryStore.Cluster())
if err := m.telemetryStore.ClickhouseDB().Exec(ctx, op.ToSQL()); err != nil {
return errors.WrapInternalf(err, CodeFailedToCreateIndex, "failed to create index")
}
}
return nil
}

View File

@@ -0,0 +1,448 @@
package implpromote
import (
"context"
"regexp"
"testing"
sqlmock "github.com/DATA-DOG/go-sqlmock"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
"github.com/SigNoz/signoz/pkg/types/promotetypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestPromotePaths(t *testing.T) {
ctx := context.Background()
testCases := []struct {
name string
paths []*promotetypes.PromotePath
promoteTwice bool
wantErr bool
wantPromoted []string
}{
{
name: "PromotesNewAttributes_Idempotent",
paths: []*promotetypes.PromotePath{
{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true},
{Signal: "traces", Context: "attribute", Path: "span.operation", Promote: true},
},
promoteTwice: true,
wantPromoted: []string{"http.method", "span.operation"},
},
{
name: "MixedDomains_RecordedPerTarget",
paths: []*promotetypes.PromotePath{
{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true},
{Signal: "logs", Context: "body", Path: "user.name", Promote: true},
},
wantPromoted: []string{"http.method", "user.name"},
},
{
name: "NonPromoteEntries_NotRecorded",
paths: []*promotetypes.PromotePath{{Signal: "traces", Context: "attribute", Path: "http.method"}},
},
{
name: "InvalidSignal_Rejected",
paths: []*promotetypes.PromotePath{{Signal: "events", Context: "attribute", Path: "http.method", Promote: true}},
wantErr: true,
},
{
name: "UnsupportedDomain_Rejected",
paths: []*promotetypes.PromotePath{{Signal: "metrics", Context: "attribute", Path: "http.method", Promote: true}},
wantErr: true,
},
{
name: "PromotesBareBodyPath_AsIs",
paths: []*promotetypes.PromotePath{{Signal: "logs", Context: "body", Path: "user.name", Promote: true}},
wantPromoted: []string{"user.name"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
store := telemetrytypestest.NewMockMetadataStore()
m := NewModule(store, nil)
err := m.PromotePaths(ctx, testCase.paths...)
if testCase.wantErr {
assert.Error(t, err)
assert.Empty(t, store.PromotedPathsMap)
return
}
require.NoError(t, err)
require.Len(t, store.PromotedPathsMap, len(testCase.wantPromoted))
for _, path := range testCase.wantPromoted {
assert.True(t, store.PromotedPathsMap[path], path)
}
if testCase.promoteTwice {
require.NoError(t, m.PromotePaths(ctx, testCase.paths...))
assert.Len(t, store.PromotedPathsMap, len(testCase.wantPromoted))
}
})
}
}
func TestPromotePathsCreatesIndexes(t *testing.T) {
ctx := context.Background()
testCases := []struct {
name string
promoted map[string]bool
path *promotetypes.PromotePath
wantDDLColumn string
}{
{
name: "LogsNewPromotion_IndexesPromotedColumn",
path: &promotetypes.PromotePath{
Signal: "logs",
Context: "body",
Path: "user.name",
Promote: true,
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1, JSONDataType: telemetrytypes.String},
},
},
wantDDLColumn: "dynamicElement(body_promoted.user.name",
},
{
name: "LogsAlreadyPromoted_IndexesPromotedColumn",
promoted: map[string]bool{"user.name": true},
path: &promotetypes.PromotePath{
Signal: "logs",
Context: "body",
Path: "user.name",
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1, JSONDataType: telemetrytypes.String},
},
},
wantDDLColumn: "dynamicElement(body_promoted.user.name",
},
{
name: "LogsUnpromotedPath_IndexesBaseColumn",
path: &promotetypes.PromotePath{
Signal: "logs",
Context: "body",
Path: "user.name",
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1, JSONDataType: telemetrytypes.String},
},
},
wantDDLColumn: "dynamicElement(body_v2.user.name",
},
{
name: "TracesNewPromotion_IndexesPromotedColumn",
path: &promotetypes.PromotePath{
Signal: "traces",
Context: "attribute",
Path: "http.method",
Promote: true,
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1, JSONDataType: telemetrytypes.String},
},
},
wantDDLColumn: "`attributes_promoted.http.method_String_ngrambf_v1` attributes_promoted.http.method::String",
},
{
name: "TracesUnpromotedPath_IndexesBaseColumn",
path: &promotetypes.PromotePath{
Signal: "traces",
Context: "attribute",
Path: "http.method",
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1, JSONDataType: telemetrytypes.String},
},
},
wantDDLColumn: "`attributes.http.method_String_ngrambf_v1` attributes.http.method::String",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
ts := telemetrystoretest.New(telemetrystore.Config{}, sqlmock.QueryMatcherRegexp)
store := telemetrytypestest.NewMockMetadataStore()
if testCase.promoted != nil {
store.PromotedPathsMap = testCase.promoted
}
m := NewModule(store, ts)
ts.Mock().ExpectExec("ADD INDEX (.+)" + regexp.QuoteMeta(testCase.wantDDLColumn)).WillReturnError(nil)
require.NoError(t, m.PromotePaths(ctx, testCase.path))
assert.NoError(t, ts.Mock().ExpectationsWereMet())
})
}
}
func TestListPromotedPaths(t *testing.T) {
ctx := context.Background()
trueValue := true
falseValue := false
testCases := []struct {
name string
filters promotetypes.ListPromotedPathsFilters
promoted map[string]bool
indexes []telemetrytypes.TelemetryFieldKeySkipIndex
wantPaths []promotetypes.PromotePath
}{
{
name: "PromotedPaths_EveryDomainAnnotated",
promoted: map[string]bool{"http.method": true},
wantPaths: []promotetypes.PromotePath{
{Signal: "logs", Context: "body", Path: "http.method", Promote: true},
{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true},
},
},
{
name: "SignalFilter_SkipsOtherDomains",
filters: promotetypes.ListPromotedPathsFilters{Signal: "traces"},
promoted: map[string]bool{"http.method": true},
wantPaths: []promotetypes.PromotePath{{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true}},
},
{
name: "ContextFilter_SkipsOtherDomains",
filters: promotetypes.ListPromotedPathsFilters{Context: "body"},
promoted: map[string]bool{"http.method": true},
wantPaths: []promotetypes.PromotePath{{Signal: "logs", Context: "body", Path: "http.method", Promote: true}},
},
{
name: "LogsIndexes_MergedWithLogsPaths",
promoted: map[string]bool{"user.name": true},
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
{
Name: "user.name",
FieldContext: telemetrytypes.FieldContextBody,
FieldDataType: telemetrytypes.FieldDataTypeString,
BaseColumn: "body_promoted.",
IndexType: "ngrambf_v1(4, 1024, 2, 0)",
Granularity: 1,
},
{
Name: "request.duration",
FieldContext: telemetrytypes.FieldContextBody,
FieldDataType: telemetrytypes.FieldDataTypeFloat64,
BaseColumn: "body_v2.",
IndexType: "minmax",
Granularity: 1,
},
},
wantPaths: []promotetypes.PromotePath{
{
Signal: "logs",
Context: "body",
Path: "user.name",
Promote: true,
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
},
},
{
Signal: "logs",
Context: "body",
Path: "request.duration",
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeFloat64, Type: "minmax", Granularity: 1},
},
},
{
Signal: "traces",
Context: "attribute",
Path: "user.name",
Promote: true,
},
},
},
{
name: "PromotedFalseFilter_IndexOnlyPaths",
filters: promotetypes.ListPromotedPathsFilters{Promoted: &falseValue},
promoted: map[string]bool{"user.name": true},
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
{
Name: "request.duration",
FieldContext: telemetrytypes.FieldContextBody,
FieldDataType: telemetrytypes.FieldDataTypeFloat64,
BaseColumn: "body_v2.",
IndexType: "minmax",
Granularity: 1,
},
},
wantPaths: []promotetypes.PromotePath{
{
Signal: "logs",
Context: "body",
Path: "request.duration",
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeFloat64, Type: "minmax", Granularity: 1},
},
},
},
},
{
name: "IndexesTrueFilter_PathsWithIndexes",
filters: promotetypes.ListPromotedPathsFilters{Indexes: &trueValue},
promoted: map[string]bool{"user.name": true},
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
{
Name: "user.name",
FieldContext: telemetrytypes.FieldContextBody,
FieldDataType: telemetrytypes.FieldDataTypeString,
BaseColumn: "body_promoted.",
IndexType: "ngrambf_v1(4, 1024, 2, 0)",
Granularity: 1,
},
},
wantPaths: []promotetypes.PromotePath{
{
Signal: "logs",
Context: "body",
Path: "user.name",
Promote: true,
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
},
},
},
},
{
name: "TracesIndexes_MergedWithTracesPaths",
promoted: map[string]bool{"http.method": true},
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
{
Name: "http.method",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
BaseColumn: "attributes_promoted.",
IndexType: "ngrambf_v1(4, 1024, 2, 0)",
Granularity: 1,
},
{
Name: "http.status_code",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeFloat64,
BaseColumn: "attributes.",
IndexType: "minmax",
Granularity: 1,
},
},
wantPaths: []promotetypes.PromotePath{
{
Signal: "logs",
Context: "body",
Path: "http.method",
Promote: true,
},
{
Signal: "traces",
Context: "attribute",
Path: "http.method",
Promote: true,
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
},
},
{
Signal: "traces",
Context: "attribute",
Path: "http.status_code",
Indexes: []promotetypes.WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeFloat64, Type: "minmax", Granularity: 1},
},
},
},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
store := telemetrytypestest.NewMockMetadataStore()
store.PromotedPathsMap = testCase.promoted
store.LogsJSONIndexes = testCase.indexes
m := NewModule(store, nil)
paths, err := m.ListPromotedPaths(ctx, testCase.filters)
require.NoError(t, err)
require.Len(t, paths, len(testCase.wantPaths))
byDomainPath := map[string]promotetypes.PromotePath{}
for _, path := range paths {
byDomainPath[path.Signal+"/"+path.Context+"/"+path.Path] = path
}
for _, want := range testCase.wantPaths {
key := want.Signal + "/" + want.Context + "/" + want.Path
require.Contains(t, byDomainPath, key)
assert.Equal(t, want, byDomainPath[key])
}
})
}
}
func TestIndexMaterializedPaths(t *testing.T) {
ctx := context.Background()
materializedKeys := []*telemetrytypes.TelemetryFieldKey{
{Name: "user.id", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "http.route", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "retry.count", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeNumber},
{Name: "k8s.pod.name", Signal: telemetrytypes.SignalLogs, FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "tenant", Signal: telemetrytypes.SignalLogs, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
}
testCases := []struct {
name string
params promotetypes.IndexMaterializedPathsParams
existing []telemetrytypes.TelemetryFieldKeySkipIndex
wantPaths []string
wantDDL []string
}{
{
name: "DryRun_NoDDL",
params: promotetypes.IndexMaterializedPathsParams{DryRun: true},
wantPaths: []string{"traces/attribute/user.id"},
},
{
name: "Run_IndexesBaseColumns",
wantPaths: []string{"traces/attribute/user.id"},
wantDDL: []string{"`attributes.user.id_String_bloom_filter` attributes.user.id::String TYPE bloom_filter(0.01) GRANULARITY 64"},
},
{
name: "SignalFilter_SkipsOtherSignals",
params: promotetypes.IndexMaterializedPathsParams{Signal: "logs", DryRun: true},
wantPaths: []string{},
},
{
name: "IndexAlreadyBuilt_Skipped",
params: promotetypes.IndexMaterializedPathsParams{Signal: "traces"},
existing: []telemetrytypes.TelemetryFieldKeySkipIndex{{Name: "user.id", FieldContext: telemetrytypes.FieldContextAttribute, IndexType: "bloom_filter(0.05)"}},
wantPaths: []string{},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
ts := telemetrystoretest.New(telemetrystore.Config{}, sqlmock.QueryMatcherRegexp)
store := telemetrytypestest.NewMockMetadataStore()
store.MaterializedKeys = materializedKeys
store.LogsJSONIndexes = testCase.existing
m := NewModule(store, ts)
for _, ddl := range testCase.wantDDL {
ts.Mock().ExpectExec("ADD INDEX (.+)" + regexp.QuoteMeta(ddl)).WillReturnError(nil)
}
paths, err := m.IndexMaterializedPaths(ctx, testCase.params)
require.NoError(t, err)
got := []string{}
for _, path := range paths {
assert.False(t, path.Promote)
got = append(got, path.Signal+"/"+path.Context+"/"+path.Path)
}
assert.Equal(t, testCase.wantPaths, got)
assert.NoError(t, ts.Mock().ExpectationsWereMet())
})
}
}

View File

@@ -8,11 +8,14 @@ import (
)
type Module interface {
ListPromotedAndIndexedPaths(ctx context.Context) ([]promotetypes.PromotePath, error)
PromoteAndIndexPaths(ctx context.Context, paths ...*promotetypes.PromotePath) error
ListPromotedPaths(ctx context.Context, filters promotetypes.ListPromotedPathsFilters) ([]promotetypes.PromotePath, error)
PromotePaths(ctx context.Context, paths ...*promotetypes.PromotePath) error
IndexMaterializedPaths(ctx context.Context, params promotetypes.IndexMaterializedPathsParams) ([]*promotetypes.PromotePath, error)
}
type Handler interface {
HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Request)
ListPromotedAndIndexedPaths(w http.ResponseWriter, r *http.Request)
PromotePaths(w http.ResponseWriter, r *http.Request)
ListPromotedPaths(w http.ResponseWriter, r *http.Request)
IndexMaterializedPaths(w http.ResponseWriter, r *http.Request)
}

View File

@@ -259,6 +259,7 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddAIObservabilityQuickFiltersFactory(sqlstore),
sqlmigration.NewAddChannelSpecFactory(sqlschema),
sqlmigration.NewAddUserTuplesFactory(sqlstore),
sqlmigration.NewAddFieldTuplesFactory(sqlstore),
)
}

View File

@@ -0,0 +1,144 @@
package sqlmigration
import (
"context"
"database/sql"
"time"
"github.com/oklog/ulid/v2"
"github.com/uptrace/bun"
"github.com/uptrace/bun/dialect"
"github.com/uptrace/bun/migrate"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlstore"
)
type addFieldTuples struct {
sqlstore sqlstore.SQLStore
}
func NewAddFieldTuplesFactory(sqlstore sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(factory.MustNewName("add_field_tuples"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addFieldTuples{sqlstore: sqlstore}, nil
})
}
func (migration *addFieldTuples) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func (migration *addFieldTuples) Up(ctx context.Context, db *bun.DB) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
var storeID string
err = tx.QueryRowContext(ctx, `SELECT id FROM store WHERE name = ? LIMIT 1`, "signoz").Scan(&storeID)
if err != nil {
return err
}
var orgIDs []string
err = tx.NewSelect().
Table("organizations").
Column("id").
Scan(ctx, &orgIDs)
if err != nil && err != sql.ErrNoRows {
return err
}
isPG := migration.sqlstore.BunDB().Dialect().Name() == dialect.PG
tuples := []migrationTuple{
{"signoz-admin", "metaresource", "logs-field", "read"},
{"signoz-admin", "metaresource", "logs-field", "update"},
{"signoz-admin", "metaresource", "logs-field", "list"},
{"signoz-editor", "metaresource", "logs-field", "read"},
{"signoz-editor", "metaresource", "logs-field", "update"},
{"signoz-editor", "metaresource", "logs-field", "list"},
{"signoz-viewer", "metaresource", "logs-field", "read"},
{"signoz-viewer", "metaresource", "logs-field", "list"},
{"signoz-admin", "metaresource", "traces-field", "read"},
{"signoz-admin", "metaresource", "traces-field", "update"},
{"signoz-admin", "metaresource", "traces-field", "list"},
{"signoz-editor", "metaresource", "traces-field", "read"},
{"signoz-editor", "metaresource", "traces-field", "update"},
{"signoz-editor", "metaresource", "traces-field", "list"},
{"signoz-viewer", "metaresource", "traces-field", "read"},
{"signoz-viewer", "metaresource", "traces-field", "list"},
}
for _, orgID := range orgIDs {
for _, tuple := range tuples {
entropy := ulid.DefaultEntropy()
now := time.Now().UTC()
tupleID := ulid.MustNew(ulid.Timestamp(now), entropy).String()
objectID := "organization/" + orgID + "/" + tuple.objectName + "/*"
roleSubject := "organization/" + orgID + "/role/" + tuple.roleName
if isPG {
user := "role:" + roleSubject + "#assignee"
result, err := tx.ExecContext(ctx, `
INSERT INTO tuple (store, object_type, object_id, relation, _user, user_type, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, object_type, object_id, relation, _user) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, user, "userset", tupleID, now,
)
if err != nil {
return err
}
rowsAffected, err := result.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
continue
}
_, err = tx.ExecContext(ctx, `
INSERT INTO changelog (store, object_type, object_id, relation, _user, operation, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, ulid, object_type) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, user, 0, tupleID, now,
)
if err != nil {
return err
}
} else {
result, err := tx.ExecContext(ctx, `
INSERT INTO tuple (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation, user_type, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, "role", roleSubject, "assignee", "userset", tupleID, now,
)
if err != nil {
return err
}
rowsAffected, err := result.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
continue
}
_, err = tx.ExecContext(ctx, `
INSERT INTO changelog (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation, operation, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, ulid, object_type) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, "role", roleSubject, "assignee", 0, tupleID, now,
)
if err != nil {
return err
}
}
}
}
return tx.Commit()
}
func (migration *addFieldTuples) Down(context.Context, *bun.DB) error {
return nil
}

View File

@@ -5,17 +5,18 @@ import (
"fmt"
"log/slog"
"reflect"
"regexp"
"strings"
"time"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz-otel-collector/constants"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
"github.com/SigNoz/signoz/pkg/types/promotetypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/huandu/go-sqlbuilder"
)
@@ -33,6 +34,26 @@ var (
CodeFailedToAppendPath = errors.MustNewCode("failed_to_append_path_promoted_paths")
)
var logsBodyPromotedEntry = promotetypes.NewLogsBodyTarget().Entry
var logsBodyIndexLookup = promotetypes.NewLogsBodyTarget().JSONIndexLookup()
// ClickHouse stores a `col.path::Type` skip index expression as CAST(col.path, 'Type').
var simpleJSONSubColumnIndexExprRe = regexp.MustCompile(`^CAST\((?P<expr>[^()]+), '(?P<type>.+)'\)$`)
// unfoldJSONSubColumnIndexExpr accepts both the folded (lower/assumeNotNull)
// and the bare type-cast expression forms.
func unfoldJSONSubColumnIndexExpr(expr string) (string, string, error) {
if columnExpr, columnType, err := schemamigrator.UnfoldJSONSubColumnIndexExpr(expr); err == nil {
return columnExpr, columnType, nil
}
matches := simpleJSONSubColumnIndexExprRe.FindStringSubmatch(expr)
if matches == nil {
return "", "", errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid expression: %s", expr)
}
return matches[1], matches[2], nil
}
// enrichJSONKeys enriches body-context keys with promoted path info, indexes,
// and JSON access plans. parentTypeCache contains parent array types (ArrayJSON/ArrayDynamic)
// pre-fetched in the main UNION query.
@@ -67,7 +88,7 @@ func (t *telemetryMetaStore) enrichJSONKeys(ctx context.Context, selectors []*te
}
// fetch promoted paths
promoted, err := t.GetPromotedPaths(ctx, paths...)
promoted, err := t.GetPromotedPaths(ctx, logsBodyPromotedEntry, paths...)
if err != nil {
return err
}
@@ -120,7 +141,7 @@ func (t *telemetryMetaStore) getJSONPathIndexes(ctx context.Context, paths ...st
}
// list indexes for the paths
indexes, err := t.ListLogsJSONIndexes(ctx, filteredPaths...)
indexes, err := t.ListJSONIndexes(ctx, logsBodyIndexLookup, filteredPaths...)
if err != nil {
return nil, errors.WrapInternalf(err, CodeFailLoadLogsJSONIndexes, "failed to list JSON path indexes")
}
@@ -134,16 +155,16 @@ func (t *telemetryMetaStore) getJSONPathIndexes(ctx context.Context, paths ...st
return fieldPathToIndexes, nil
}
func buildListLogsJSONIndexesQuery(cluster string, filters ...string) (string, []any) {
func buildListJSONIndexesQuery(cluster string, lookup telemetrytypes.JSONIndexLookup, filters ...string) (string, []any) {
sb := sqlbuilder.Select(
"name", "type_full", "expr", "granularity",
).From(fmt.Sprintf("clusterAllReplicas('%s', %s)", cluster, SkipIndexTableName))
sb.Where(sb.Equal("database", logstelemetryschema.DBName))
sb.Where(sb.Equal("table", logstelemetryschema.LogsV2LocalTableName))
sb.Where(sb.Equal("database", lookup.DBName))
sb.Where(sb.Equal("table", lookup.LocalTableName))
sb.Where(sb.Or(
sb.ILike("expr", fmt.Sprintf("%%%s%%", querybuilder.FormatValueForContains(constants.BodyV2ColumnPrefix))),
sb.ILike("expr", fmt.Sprintf("%%%s%%", querybuilder.FormatValueForContains(constants.BodyPromotedColumnPrefix))),
sb.ILike("expr", fmt.Sprintf("%%%s%%", querybuilder.FormatValueForContains(lookup.BaseColumnPrefix))),
sb.ILike("expr", fmt.Sprintf("%%%s%%", querybuilder.FormatValueForContains(lookup.PromotedColumnPrefix))),
))
filterExprs := []string{}
@@ -156,9 +177,9 @@ func buildListLogsJSONIndexesQuery(cluster string, filters ...string) (string, [
return sb.BuildWithFlavor(sqlbuilder.ClickHouse)
}
func (t *telemetryMetaStore) ListLogsJSONIndexes(ctx context.Context, filters ...string) ([]telemetrytypes.TelemetryFieldKeySkipIndex, error) {
ctx = withTelemetryContext(ctx, "ListLogsJSONIndexes")
query, args := buildListLogsJSONIndexesQuery(t.telemetrystore.Cluster(), filters...)
func (t *telemetryMetaStore) ListJSONIndexes(ctx context.Context, lookup telemetrytypes.JSONIndexLookup, filters ...string) ([]telemetrytypes.TelemetryFieldKeySkipIndex, error) {
ctx = withTelemetryContext(ctx, lookup.Signal, "ListJSONIndexes")
query, args := buildListJSONIndexesQuery(t.telemetrystore.Cluster(), lookup, filters...)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return nil, errors.WrapInternalf(err, CodeFailLoadLogsJSONIndexes, "failed to load string indexed columns")
@@ -175,7 +196,7 @@ func (t *telemetryMetaStore) ListLogsJSONIndexes(ctx context.Context, filters ..
return nil, errors.WrapInternalf(err, CodeFailLoadLogsJSONIndexes, "failed to scan string indexed column")
}
columnExpr, columnType, err := schemamigrator.UnfoldJSONSubColumnIndexExpr(expr)
columnExpr, columnType, err := unfoldJSONSubColumnIndexExpr(expr)
if err != nil {
return nil, errors.WrapInternalf(err, CodeFailLoadLogsJSONIndexes, "failed to unfold JSON sub column index expression: %s", expr)
}
@@ -189,18 +210,18 @@ func (t *telemetryMetaStore) ListLogsJSONIndexes(ctx context.Context, filters ..
baseColumn := ""
fieldName := ""
switch {
case strings.HasPrefix(columnExpr, logstelemetryschema.BodyV2ColumnPrefix):
baseColumn = logstelemetryschema.BodyV2ColumnPrefix
fieldName = strings.TrimPrefix(columnExpr, logstelemetryschema.BodyV2ColumnPrefix)
case strings.HasPrefix(columnExpr, logstelemetryschema.BodyPromotedColumnPrefix):
baseColumn = logstelemetryschema.BodyPromotedColumnPrefix
fieldName = strings.TrimPrefix(columnExpr, logstelemetryschema.BodyPromotedColumnPrefix)
case strings.HasPrefix(columnExpr, lookup.BaseColumnPrefix):
baseColumn = lookup.BaseColumnPrefix
fieldName = strings.TrimPrefix(columnExpr, lookup.BaseColumnPrefix)
case strings.HasPrefix(columnExpr, lookup.PromotedColumnPrefix):
baseColumn = lookup.PromotedColumnPrefix
fieldName = strings.TrimPrefix(columnExpr, lookup.PromotedColumnPrefix)
}
fieldName = strings.ReplaceAll(fieldName, "`", "")
indexes = append(indexes, telemetrytypes.TelemetryFieldKeySkipIndex{
Name: fieldName,
FieldContext: telemetrytypes.FieldContextBody,
FieldContext: lookup.FieldContext,
FieldDataType: fdt,
BaseColumn: baseColumn,
IndexName: name,
@@ -215,14 +236,14 @@ func (t *telemetryMetaStore) ListLogsJSONIndexes(ctx context.Context, filters ..
// TODO(Piyush): Remove this if not used in future.
func (t *telemetryMetaStore) ListJSONValues(ctx context.Context, path string, limit int) (*telemetrytypes.TelemetryFieldValues, bool, error) {
ctx = withTelemetryContext(ctx, "ListJSONValues")
ctx = withTelemetryContext(ctx, telemetrytypes.SignalLogs, "ListJSONValues")
path = CleanPathPrefixes(path)
if strings.Contains(path, telemetrytypes.ArraySep) || strings.Contains(path, telemetrytypes.ArrayAnyIndex) {
return nil, false, errors.NewInvalidInputf(errors.CodeInvalidInput, "array paths are not supported")
}
promoted, err := t.IsPathPromoted(ctx, path)
promoted, err := t.isPathPromoted(ctx, logsBodyPromotedEntry, path)
if err != nil {
return nil, false, err
}
@@ -376,13 +397,12 @@ func derefValue(v any) any {
return val.Interface()
}
// IsPathPromoted checks if a specific path is promoted (Column Evolution table: field_name for logs body).
func (t *telemetryMetaStore) IsPathPromoted(ctx context.Context, path string) (bool, error) {
ctx = withTelemetryContext(ctx, "IsPathPromoted")
func (t *telemetryMetaStore) isPathPromoted(ctx context.Context, entry telemetrytypes.EvolutionEntry, path string) (bool, error) {
ctx = withTelemetryContext(ctx, entry.Signal, "isPathPromoted")
split := strings.Split(path, telemetrytypes.ArraySep)
pathSegment := split[0]
query := fmt.Sprintf("SELECT 1 FROM %s.%s WHERE signal = ? AND column_name = ? AND field_context = ? AND field_name = ? LIMIT 1", DBName, PromotedPathsTableName)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, telemetrytypes.SignalLogs, logstelemetryschema.LogsV2BodyPromotedColumn, telemetrytypes.FieldContextBody, pathSegment)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, entry.Signal, entry.ColumnName, entry.FieldContext, pathSegment)
if err != nil {
return false, errors.WrapInternalf(err, CodeFailCheckPathPromoted, "failed to check if path %s is promoted", path)
}
@@ -391,14 +411,13 @@ func (t *telemetryMetaStore) IsPathPromoted(ctx context.Context, path string) (b
return rows.Next(), nil
}
// GetPromotedPaths returns promoted paths from the Column Evolution table (field_name for logs body).
func (t *telemetryMetaStore) GetPromotedPaths(ctx context.Context, paths ...string) (map[string]bool, error) {
ctx = withTelemetryContext(ctx, "GetPromotedPaths")
func (t *telemetryMetaStore) GetPromotedPaths(ctx context.Context, entry telemetrytypes.EvolutionEntry, paths ...string) (map[string]bool, error) {
ctx = withTelemetryContext(ctx, entry.Signal, "GetPromotedPaths")
sb := sqlbuilder.Select("field_name").From(fmt.Sprintf("%s.%s", DBName, PromotedPathsTableName))
conditions := []string{
sb.Equal("signal", telemetrytypes.SignalLogs),
sb.Equal("column_name", logstelemetryschema.LogsV2BodyPromotedColumn),
sb.Equal("field_context", telemetrytypes.FieldContextBody),
sb.Equal("signal", entry.Signal),
sb.Equal("column_name", entry.ColumnName),
sb.Equal("field_context", entry.FieldContext),
sb.NotEqual("field_name", "__all__"),
}
if len(paths) > 0 {
@@ -438,9 +457,8 @@ func CleanPathPrefixes(path string) string {
return path
}
// PromotePaths inserts promoted paths into the Column Evolution table (same schema as signoz-otel-collector metadata_migrations).
func (t *telemetryMetaStore) PromotePaths(ctx context.Context, paths ...string) error {
ctx = withTelemetryContext(ctx, "PromotePaths")
func (t *telemetryMetaStore) PromotePaths(ctx context.Context, entry telemetrytypes.EvolutionEntry, paths ...string) error {
ctx = withTelemetryContext(ctx, entry.Signal, "PromotePaths")
batch, err := t.telemetrystore.ClickhouseDB().PrepareBatch(ctx,
fmt.Sprintf("INSERT INTO %s.%s (signal, column_name, column_type, field_context, field_name, version, release_time) VALUES", DBName,
PromotedPathsTableName))
@@ -454,7 +472,7 @@ func (t *telemetryMetaStore) PromotePaths(ctx context.Context, paths ...string)
if trimmed == "" {
continue
}
if err := batch.Append(telemetrytypes.SignalLogs, logstelemetryschema.LogsV2BodyPromotedColumn, "JSON()", telemetrytypes.FieldContextBody, trimmed, 0, releaseTime); err != nil {
if err := batch.Append(entry.Signal, entry.ColumnName, entry.ColumnType, entry.FieldContext, trimmed, entry.Version, releaseTime); err != nil {
_ = batch.Abort()
return errors.WrapInternalf(err, CodeFailedToAppendPath, "failed to append path")
}
@@ -466,9 +484,9 @@ func (t *telemetryMetaStore) PromotePaths(ctx context.Context, paths ...string)
return nil
}
func withTelemetryContext(ctx context.Context, functionName string) context.Context {
func withTelemetryContext(ctx context.Context, signal telemetrytypes.Signal, functionName string) context.Context {
return ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalLogs.StringValue(),
instrumentationtypes.TelemetrySignal: signal.StringValue(),
instrumentationtypes.CodeNamespace: "metadata",
instrumentationtypes.CodeFunctionName: functionName,
})

View File

@@ -50,7 +50,7 @@ func TestBuildListLogsJSONIndexesQuery(t *testing.T) {
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
query, args := buildListLogsJSONIndexesQuery(tc.cluster, tc.filters...)
query, args := buildListJSONIndexesQuery(tc.cluster, logsBodyIndexLookup, tc.filters...)
require.Equal(t, tc.expectedSQL, query)
require.Equal(t, tc.expectedArgs, args)

View File

@@ -118,6 +118,17 @@ func NewTelemetryMetaStore(
return t
}
func (t *telemetryMetaStore) GetMaterializedKeys(ctx context.Context, signal telemetrytypes.Signal) ([]*telemetrytypes.TelemetryFieldKey, error) {
switch signal {
case telemetrytypes.SignalTraces:
return t.tracesTblStatementToFieldKeys(ctx)
case telemetrytypes.SignalLogs:
return t.logsTblStatementToFieldKeys(ctx)
default:
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "materialized keys are not supported for signal %s", signal.StringValue())
}
}
// tracesTblStatementToFieldKeys returns materialised attribute/resource/scope keys from the traces table.
func (t *telemetryMetaStore) tracesTblStatementToFieldKeys(ctx context.Context) ([]*telemetrytypes.TelemetryFieldKey, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{

View File

@@ -0,0 +1,86 @@
package promotetypes
import (
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
const (
DefaultMaterializedIndexType = "bloom_filter(0.01)"
DefaultMaterializedIndexGranularity = 64
)
// defaultMaterializedPaths are keyed by signal/context.
var defaultMaterializedPaths = map[string]map[string]struct{}{
"traces/attribute": {
"db.system": {},
"http.route": {},
"messaging.operation": {},
"messaging.system": {},
"peer.service": {},
"rpc.method": {},
"rpc.service": {},
"rpc.system": {},
"gen_ai.agent.name": {},
"gen_ai.provider.name": {},
"gen_ai.request.model": {},
"gen_ai.tool.name": {},
"gen_ai.usage.cache_creation.input_tokens": {},
"gen_ai.usage.cache_read.input_tokens": {},
"gen_ai.usage.input_tokens": {},
"gen_ai.usage.output_tokens": {},
"signoz.gen_ai.usage.cache_read.input_tokens.cost": {},
"signoz.gen_ai.usage.cache_write.input_tokens.cost": {},
"signoz.gen_ai.usage.input_tokens.cost": {},
"signoz.gen_ai.usage.output_tokens.cost": {},
"signoz.gen_ai.usage.tokens.cost": {},
},
}
type IndexMaterializedPathsParams struct {
Signal string `query:"signal" json:"signal"`
Context string `query:"context" json:"context"`
DryRun bool `query:"dryRun" json:"dryRun"`
}
func (p *IndexMaterializedPathsParams) Filters() ListPromotedPathsFilters {
return ListPromotedPathsFilters{Signal: p.Signal, Context: p.Context}
}
func NewMaterializedPromotePath(target Target, name string) (*PromotePath, bool) {
path := &PromotePath{
Signal: target.Entry.Signal.StringValue(),
Context: target.Entry.FieldContext.StringValue(),
Path: name,
Promote: false,
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeString, Type: DefaultMaterializedIndexType, Granularity: DefaultMaterializedIndexGranularity}},
}
if err := path.ValidateAndSetDefaults(target); err != nil {
return nil, false
}
return path, true
}
// SkipExistingIndexes matches on index type, ignoring its parameters.
func (i *PromotePath) SkipExistingIndexes(existing []telemetrytypes.TelemetryFieldKeySkipIndex) {
built := map[string]struct{}{}
for _, index := range existing {
if indexType, ok := (WrappedIndex{Type: index.IndexType}).SkipIndexType(); ok {
built[string(indexType)] = struct{}{}
}
}
indexes := make([]WrappedIndex, 0, len(i.Indexes))
for _, index := range i.Indexes {
indexType, _ := index.SkipIndexType()
if _, ok := built[string(indexType)]; !ok {
indexes = append(indexes, index)
}
}
i.Indexes = indexes
}
func (t Target) IsDefaultMaterialized(path string) bool {
_, ok := defaultMaterializedPaths[t.Entry.Signal.StringValue()+"/"+t.Entry.FieldContext.StringValue()][path]
return ok
}

View File

@@ -0,0 +1,92 @@
package promotetypes
import (
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
func TestNewMaterializedPromotePath(t *testing.T) {
testCases := []struct {
name string
path string
wantOK bool
}{
{
name: "BarePath_DefaultIndex",
path: "user.id",
wantOK: true,
},
{
name: "CardinalPath_Rejected",
path: "session.550e8400-e29b-41d4-a716-446655440000",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
path, ok := NewMaterializedPromotePath(NewTracesAttributesTarget(), testCase.path)
require.Equal(t, testCase.wantOK, ok)
if !ok {
return
}
assert.Equal(t, &PromotePath{
Signal: "traces",
Context: "attribute",
Path: testCase.path,
Promote: false,
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "bloom_filter(0.01)", Granularity: 64, JSONDataType: telemetrytypes.String}},
}, path)
})
}
}
func TestSkipExistingIndexes(t *testing.T) {
testCases := []struct {
name string
existing []telemetrytypes.TelemetryFieldKeySkipIndex
want []WrappedIndex
}{
{
name: "SameTypeOtherParams_Skipped",
existing: []telemetrytypes.TelemetryFieldKeySkipIndex{{IndexType: "bloom_filter(0.05)"}},
want: []WrappedIndex{{Type: "tokenbf_v1(1024, 2, 0)", Granularity: 1}},
},
{
name: "OtherType_Kept",
existing: []telemetrytypes.TelemetryFieldKeySkipIndex{{IndexType: "ngrambf_v1(4, 1024, 2, 0)"}},
want: []WrappedIndex{{Type: "bloom_filter(0.01)", Granularity: 64}, {Type: "tokenbf_v1(1024, 2, 0)", Granularity: 1}},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
path := &PromotePath{Indexes: []WrappedIndex{{Type: "bloom_filter(0.01)", Granularity: 64}, {Type: "tokenbf_v1(1024, 2, 0)", Granularity: 1}}}
path.SkipExistingIndexes(testCase.existing)
assert.Equal(t, testCase.want, path.Indexes)
})
}
}
func TestIsDefaultMaterialized(t *testing.T) {
testCases := []struct {
name string
target Target
path string
want bool
}{
{name: "TracesDefaultAttribute_Default", target: NewTracesAttributesTarget(), path: "http.route", want: true},
{name: "TracesGenAIAttribute_Default", target: NewTracesAttributesTarget(), path: "gen_ai.request.model", want: true},
{name: "TracesCustomAttribute_NotDefault", target: NewTracesAttributesTarget(), path: "user.id"},
{name: "LogsBodySamePath_NotDefault", target: NewLogsBodyTarget(), path: "http.route"},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.want, testCase.target.IsDefaultMaterialized(testCase.path))
})
}
}

View File

@@ -0,0 +1,72 @@
package promotetypes
import (
"encoding/json"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/coretypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// PromotePathsResources resolves the field resources of the signals in a promote request body.
func PromotePathsResources(ec coretypes.ExtractorContext) ([]coretypes.ResourceWithID, error) {
var paths []PromotePath
if err := json.Unmarshal(ec.RequestBody, &paths); err != nil {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid promote paths request body")
}
signals := make([]telemetrytypes.Signal, 0, len(paths))
for _, path := range paths {
signal, ok := telemetrytypes.SignalFromText(path.Signal)
if !ok {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid signal: %s", path.Signal)
}
signals = append(signals, signal)
}
return fieldResources(signals)
}
// FilteredSignalsResources falls back to every signal when the filters match none.
func FilteredSignalsResources(ec coretypes.ExtractorContext) ([]coretypes.ResourceWithID, error) {
filters := ListPromotedPathsFilters{}
if ec.Request != nil {
filters.Signal = ec.Request.URL.Query().Get("signal")
filters.Context = ec.Request.URL.Query().Get("context")
}
if err := filters.Validate(); err != nil {
return nil, err
}
signals := make([]telemetrytypes.Signal, 0)
for _, target := range Targets() {
if filters.MatchesTarget(target) {
signals = append(signals, target.Entry.Signal)
}
}
if len(signals) == 0 {
for _, target := range Targets() {
signals = append(signals, target.Entry.Signal)
}
}
return fieldResources(signals)
}
func fieldResources(signals []telemetrytypes.Signal) ([]coretypes.ResourceWithID, error) {
resources := make([]coretypes.ResourceWithID, 0, len(signals))
seen := make(map[string]struct{})
for _, signal := range signals {
resource, ok := signal.FieldResource()
if !ok {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "promotion is not supported for signal %s", signal.StringValue())
}
if _, ok := seen[resource.Kind().String()]; ok {
continue
}
seen[resource.Kind().String()] = struct{}{}
resources = append(resources, coretypes.ResourceWithID{Resource: resource})
}
return resources, nil
}

View File

@@ -0,0 +1,119 @@
package promotetypes
import (
"net/http/httptest"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/SigNoz/signoz/pkg/types/coretypes"
)
func resourceKinds(resources []coretypes.ResourceWithID) []string {
kinds := make([]string, 0, len(resources))
for _, resource := range resources {
kinds = append(kinds, resource.Resource.Kind().String())
}
return kinds
}
func TestPromotePathsResources(t *testing.T) {
testCases := []struct {
name string
body string
want []string
wantErr bool
}{
{
name: "Traces_TracesField",
body: `[{"signal":"traces","context":"attribute","path":"http.method"}]`,
want: []string{"traces-field"},
},
{
name: "LogsAndTraces_BothFieldsOnce",
body: `[{"signal":"logs","context":"body","path":"a"},{"signal":"traces","context":"attribute","path":"b"},{"signal":"logs","context":"body","path":"c"}]`,
want: []string{"logs-field", "traces-field"},
},
{
name: "Empty_NoResources",
body: `[]`,
want: []string{},
},
{
name: "InvalidSignal_Rejected",
body: `[{"signal":"events","context":"attribute","path":"a"}]`,
wantErr: true,
},
{
name: "UnsupportedSignal_Rejected",
body: `[{"signal":"metrics","context":"attribute","path":"a"}]`,
wantErr: true,
},
{
name: "MalformedBody_Rejected",
body: `{"signal":"traces"}`,
wantErr: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
resources, err := PromotePathsResources(coretypes.ExtractorContext{RequestBody: []byte(testCase.body)})
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, resourceKinds(resources))
})
}
}
func TestFilteredSignalsResources(t *testing.T) {
testCases := []struct {
name string
query string
want []string
wantErr bool
}{
{
name: "NoFilters_EveryField",
query: "",
want: []string{"logs-field", "traces-field"},
},
{
name: "TracesSignal_TracesField",
query: "?signal=traces",
want: []string{"traces-field"},
},
{
name: "BodyContext_LogsField",
query: "?context=body",
want: []string{"logs-field"},
},
{
name: "NoMatchingDomain_EveryField",
query: "?signal=metrics",
want: []string{"logs-field", "traces-field"},
},
{
name: "InvalidSignal_Rejected",
query: "?signal=events",
wantErr: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
request := httptest.NewRequest("GET", "/api/v1/promoted_path"+testCase.query, nil)
resources, err := FilteredSignalsResources(coretypes.ExtractorContext{Request: request})
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, resourceKinds(resources))
})
}
}

View File

@@ -0,0 +1,152 @@
package promotetypes
import (
"fmt"
"strings"
schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz-otel-collector/pkg/keycheck"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/tracestelemetryschema"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// Target identifies a promotion domain.
type Target struct {
Entry telemetrytypes.EvolutionEntry // evolution row template; FieldName and ReleaseTime are set per write
DBName string // index DDL database
LocalTableName string // index DDL local table
BaseColumn string // column holding every path; indexes for unpromoted paths are created on it
}
func (t Target) PromotedColumn() string { return t.Entry.ColumnName }
func (t Target) BaseColumnPrefix() string { return t.BaseColumn + "." }
func (t Target) PromotedColumnPrefix() string { return t.PromotedColumn() + "." }
// IndexExpression folds logs strings to lower case over assumeNotNull for
// case-insensitive LIKE searches; traces indexes are a bare type cast.
func (t Target) IndexExpression(column, path, jsonDataType string) string {
switch t.Entry.Signal {
case telemetrytypes.SignalLogs:
return schemamigrator.JSONSubColumnIndexExpr(column, path, jsonDataType)
default:
return simpleJSONSubColumnIndexExpr(column, path, jsonDataType)
}
}
func (t Target) JSONIndexLookup() telemetrytypes.JSONIndexLookup {
return telemetrytypes.JSONIndexLookup{
Signal: t.Entry.Signal,
FieldContext: t.Entry.FieldContext,
DBName: t.DBName,
LocalTableName: t.LocalTableName,
BaseColumnPrefix: t.BaseColumnPrefix(),
PromotedColumnPrefix: t.PromotedColumnPrefix(),
}
}
func NewTarget(entry telemetrytypes.EvolutionEntry, dbName, localTableName, baseColumn string) Target {
return Target{
Entry: entry,
DBName: dbName,
LocalTableName: localTableName,
BaseColumn: baseColumn,
}
}
// NewLogsBodyTarget returns the logs body domain (body_v2 -> body_promoted).
func NewLogsBodyTarget() Target {
return NewTarget(
telemetrytypes.EvolutionEntry{
Signal: telemetrytypes.SignalLogs,
ColumnName: logstelemetryschema.LogsV2BodyPromotedColumn,
ColumnType: "JSON()",
FieldContext: telemetrytypes.FieldContextBody,
},
logstelemetryschema.DBName,
logstelemetryschema.LogsV2LocalTableName,
logstelemetryschema.LogsV2BodyV2Column,
)
}
// NewTracesAttributesTarget returns the spans attributes domain (attributes
// -> attributes_promoted).
func NewTracesAttributesTarget() Target {
return NewTarget(
telemetrytypes.EvolutionEntry{
Signal: telemetrytypes.SignalTraces,
ColumnName: tracestelemetryschema.SpanAttributesPromotedColumn,
ColumnType: "JSON()",
FieldContext: telemetrytypes.FieldContextAttribute,
},
tracestelemetryschema.DBName,
tracestelemetryschema.SpanIndexV3LocalTableName,
tracestelemetryschema.SpanAttributesColumn,
)
}
func NewTargetFromText(signal, context string) (Target, error) {
parsedSignal, ok := telemetrytypes.SignalFromText(signal)
if !ok {
return Target{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid signal: %s", signal)
}
parsedContext, ok := telemetrytypes.FieldContextFromText(context)
if !ok {
return Target{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid context: %s", context)
}
target, ok := TargetFor(parsedSignal, parsedContext)
if !ok {
return Target{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "promotion is not supported for %s %s", parsedSignal.StringValue(), parsedContext.StringValue())
}
return target, nil
}
func Targets() []Target {
return []Target{
NewLogsBodyTarget(),
NewTracesAttributesTarget(),
}
}
func TargetFor(signal telemetrytypes.Signal, context telemetrytypes.FieldContext) (Target, bool) {
for _, target := range Targets() {
if target.Entry.Signal.StringValue() == signal.StringValue() &&
target.Entry.FieldContext.StringValue() == context.StringValue() {
return target, true
}
}
return Target{}, false
}
// simpleJSONSubColumnIndexExpr renders `column.path::Type`; the cast unwraps
// the Nullable the sub-column access returns, which bloom filter indexes
// reject.
func simpleJSONSubColumnIndexExpr(column, path, jsonDataType string) string {
parts := strings.Split(column+"."+path, ".")
for idx, part := range parts {
if keycheck.IsBacktickRequired(part) {
part := strings.Trim(part, "`")
parts[idx] = "`" + part + "`"
}
}
return fmt.Sprintf("%s::%s", strings.Join(parts, "."), jsonDataType)
}
// reservedPathPrefix returns the domain prefix path carries, if any: the base
// or promoted column prefix, or the logs body search alias, which is stripped
// from queries before metadata lookup and can never match.
func (t Target) reservedPathPrefix(path string) (string, bool) {
prefixes := []string{t.BaseColumnPrefix(), t.PromotedColumnPrefix()}
if t.Entry.Signal == telemetrytypes.SignalLogs {
prefixes = append(prefixes, telemetrytypes.BodyJSONStringSearchPrefix)
}
for _, prefix := range prefixes {
if strings.HasPrefix(path, prefix) {
return prefix, true
}
}
return "", false
}

View File

@@ -0,0 +1,57 @@
package promotetypes
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestIndexExpression(t *testing.T) {
testCases := []struct {
name string
target Target
column string
path string
jsonDataType string
want string
}{
{
name: "LogsString_LoweredOverAssumeNotNull",
target: NewLogsBodyTarget(),
column: "body_promoted",
path: "user.name",
jsonDataType: "String",
want: "lower(assumeNotNull(dynamicElement(body_promoted.user.name, 'String')))",
},
{
name: "LogsNumber_AssumeNotNullOnly",
target: NewLogsBodyTarget(),
column: "body_v2",
path: "request.duration",
jsonDataType: "Float64",
want: "assumeNotNull(dynamicElement(body_v2.request.duration, 'Float64'))",
},
{
name: "TracesString_TypeCastOnly",
target: NewTracesAttributesTarget(),
column: "attributes_promoted",
path: "http.method",
jsonDataType: "String",
want: "attributes_promoted.http.method::String",
},
{
name: "TracesPathNeedingBackticks_Backticked",
target: NewTracesAttributesTarget(),
column: "attributes",
path: "user-name",
jsonDataType: "String",
want: "attributes.`user-name`::String",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.want, testCase.target.IndexExpression(testCase.column, testCase.path, testCase.jsonDataType))
})
}
}

View File

@@ -3,12 +3,15 @@ package promotetypes
import (
"strings"
"github.com/SigNoz/signoz-otel-collector/constants"
schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz-otel-collector/pkg/keycheck"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// TODO: use schemamigrator.IndexTypeBloomFilter once SigNoz/signoz-otel-collector#942 is released.
const IndexTypeBloomFilter schemamigrator.IndexType = "bloom_filter"
type WrappedIndex struct {
JSONDataType telemetrytypes.JSONDataType `json:"-"`
FieldDataType telemetrytypes.FieldDataType `json:"fieldDataType"`
@@ -17,13 +20,71 @@ type WrappedIndex struct {
}
type PromotePath struct {
Path string `json:"path"`
Signal string `json:"signal" required:"true"`
Context string `json:"context" required:"true"`
Path string `json:"path" required:"true"`
Promote bool `json:"promote,omitempty"`
Indexes []WrappedIndex `json:"indexes,omitempty"`
}
func (i *PromotePath) ValidateAndSetDefaults() error {
func (w WrappedIndex) SkipIndexType() (schemamigrator.IndexType, bool) {
for _, indexType := range []schemamigrator.IndexType{schemamigrator.IndexTypeNGramBF, schemamigrator.IndexTypeTokenBF, schemamigrator.IndexTypeMinMax, IndexTypeBloomFilter} {
if strings.HasPrefix(w.Type, string(indexType)) {
return indexType, true
}
}
return "", false
}
func (i *PromotePath) Target() (Target, error) {
return NewTargetFromText(i.Signal, i.Context)
}
type ListPromotedPathsFilters struct {
Signal string `query:"signal" json:"signal"`
Context string `query:"context" json:"context"`
Promoted *bool `query:"promoted" json:"promoted"`
Indexes *bool `query:"indexes" json:"indexes"`
}
// Validate checks the signal and context words are known; the pair need not
// name a supported domain.
func (f *ListPromotedPathsFilters) Validate() error {
if f.Signal != "" {
if _, ok := telemetrytypes.SignalFromText(f.Signal); !ok {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid signal: %s", f.Signal)
}
}
if f.Context != "" {
if _, ok := telemetrytypes.FieldContextFromText(f.Context); !ok {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid context: %s", f.Context)
}
}
return nil
}
func (f *ListPromotedPathsFilters) MatchesTarget(target Target) bool {
if f.Signal != "" && f.Signal != target.Entry.Signal.StringValue() {
return false
}
if f.Context != "" && f.Context != target.Entry.FieldContext.StringValue() {
return false
}
return true
}
func (f *ListPromotedPathsFilters) MatchesPath(path PromotePath) bool {
if f.Promoted != nil && *f.Promoted != path.Promote {
return false
}
if f.Indexes != nil && *f.Indexes != (len(path.Indexes) > 0) {
return false
}
return true
}
func (i *PromotePath) ValidateAndSetDefaults(target Target) error {
if i.Path == "" {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "path is required")
}
@@ -36,17 +97,10 @@ func (i *PromotePath) ValidateAndSetDefaults() error {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "array paths can not be promoted or indexed")
}
if strings.HasPrefix(i.Path, constants.BodyV2ColumnPrefix) || strings.HasPrefix(i.Path, constants.BodyPromotedColumnPrefix) {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "`%s`, `%s` don't add these prefixes to the path", constants.BodyV2ColumnPrefix, constants.BodyPromotedColumnPrefix)
if prefix, ok := target.reservedPathPrefix(i.Path); ok {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "path must be a bare attribute name, without the `%s` prefix", prefix)
}
if !strings.HasPrefix(i.Path, telemetrytypes.BodyJSONStringSearchPrefix) {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "path must start with `body.`")
}
// remove the "body." prefix from the path
i.Path = strings.TrimPrefix(i.Path, telemetrytypes.BodyJSONStringSearchPrefix)
isCardinal := keycheck.IsCardinal(i.Path)
if isCardinal {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "cardinal paths can not be promoted or indexed")

View File

@@ -0,0 +1,349 @@
package promotetypes
import (
"testing"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestPromotePathTarget(t *testing.T) {
testCases := []struct {
name string
path *PromotePath
want Target
wantErr bool
}{
{
name: "LogsBody_Resolved",
path: &PromotePath{Signal: "logs", Context: "body", Path: "body.user.name"},
want: NewLogsBodyTarget(),
},
{
name: "TracesAttribute_Resolved",
path: &PromotePath{Signal: "traces", Context: "attribute", Path: "http.method"},
want: NewTracesAttributesTarget(),
},
{
name: "InvalidSignal_Rejected",
path: &PromotePath{Signal: "events", Context: "attribute", Path: "http.method"},
wantErr: true,
},
{
name: "InvalidContext_Rejected",
path: &PromotePath{Signal: "logs", Context: "span", Path: "user.name"},
wantErr: true,
},
{
name: "UnsupportedDomain_Rejected",
path: &PromotePath{Signal: "metrics", Context: "attribute", Path: "http.method"},
wantErr: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
target, err := testCase.path.Target()
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, target)
})
}
}
func TestValidateAndSetDefaultsLogsBody(t *testing.T) {
target := NewLogsBodyTarget()
testCases := []struct {
name string
path *PromotePath
wantErr bool
wantPath string
wantJSONDataType telemetrytypes.JSONDataType
}{
{
name: "BarePath_KeptAsIs",
path: &PromotePath{Path: "user.name", Promote: true},
wantPath: "user.name",
},
{
name: "BodyPrefixedPath_Rejected",
path: &PromotePath{Path: "body.user.name", Promote: true},
wantErr: true,
},
{
name: "BodyV2PrefixedPath_Rejected",
path: &PromotePath{Path: "body_v2.user.name", Promote: true},
wantErr: true,
},
{
name: "BodyPromotedPrefixedPath_Rejected",
path: &PromotePath{Path: "body_promoted.user.name", Promote: true},
wantErr: true,
},
{
name: "EmptyPath_Rejected",
path: &PromotePath{Path: "", Promote: true},
wantErr: true,
},
{
name: "SpacedPath_Rejected",
path: &PromotePath{Path: "body.my path", Promote: true},
wantErr: true,
},
{
name: "ArrayIndexPath_Rejected",
path: &PromotePath{Path: "body.users[].id", Promote: true},
wantErr: true,
},
{
name: "ArrayWildcardPath_Rejected",
path: &PromotePath{Path: "body.users[*].id", Promote: true},
wantErr: true,
},
{
name: "CardinalPath_Rejected",
path: &PromotePath{Path: "body.request.550e8400-e29b-41d4-a716-446655440000", Promote: true},
wantErr: true,
},
{
name: "ValidIndex_JSONDataTypeDefaulted",
path: &PromotePath{
Path: "user.name",
Indexes: []WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
},
},
wantPath: "user.name",
wantJSONDataType: telemetrytypes.String,
},
{
name: "UnsupportedColumnTypeIndex_Rejected",
path: &PromotePath{
Path: "user.active",
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeBool, Type: "minmax", Granularity: 1}},
},
wantErr: true,
},
{
name: "IndexWithoutType_Rejected",
path: &PromotePath{
Path: "user.name",
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeString, Granularity: 1}},
},
wantErr: true,
},
{
name: "IndexWithoutGranularity_Rejected",
path: &PromotePath{
Path: "user.name",
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "minmax"}},
},
wantErr: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
err := testCase.path.ValidateAndSetDefaults(target)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.wantPath, testCase.path.Path)
if testCase.wantJSONDataType != (telemetrytypes.JSONDataType{}) {
require.Len(t, testCase.path.Indexes, 1)
assert.Equal(t, testCase.wantJSONDataType, testCase.path.Indexes[0].JSONDataType)
}
})
}
}
func TestValidateAndSetDefaultsTracesAttributes(t *testing.T) {
target := NewTracesAttributesTarget()
testCases := []struct {
name string
path *PromotePath
wantErr bool
wantPath string
wantJSONDataType telemetrytypes.JSONDataType
}{
{
name: "BareAttributeName_KeptAsIs",
path: &PromotePath{Path: "http.method", Promote: true},
wantPath: "http.method",
},
{
name: "ValidIndex_JSONDataTypeDefaulted",
path: &PromotePath{
Path: "http.method",
Indexes: []WrappedIndex{
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
},
},
wantPath: "http.method",
wantJSONDataType: telemetrytypes.String,
},
{
name: "AttributesPrefixedPath_Rejected",
path: &PromotePath{Path: "attributes.http.method", Promote: true},
wantErr: true,
},
{
name: "AttributesPromotedPrefixedPath_Rejected",
path: &PromotePath{Path: "attributes_promoted.http.method", Promote: true},
wantErr: true,
},
{
name: "EmptyPath_Rejected",
path: &PromotePath{Path: "", Promote: true},
wantErr: true,
},
{
name: "SpacedPath_Rejected",
path: &PromotePath{Path: "my attr", Promote: true},
wantErr: true,
},
{
name: "ArrayIndexPath_Rejected",
path: &PromotePath{Path: "tags[].id", Promote: true},
wantErr: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
err := testCase.path.ValidateAndSetDefaults(target)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.wantPath, testCase.path.Path)
if testCase.wantJSONDataType != (telemetrytypes.JSONDataType{}) {
require.Len(t, testCase.path.Indexes, 1)
assert.Equal(t, testCase.wantJSONDataType, testCase.path.Indexes[0].JSONDataType)
}
})
}
}
func TestListPromotedPathsFiltersValidate(t *testing.T) {
testCases := []struct {
name string
filters ListPromotedPathsFilters
wantErr bool
}{
{
name: "EmptyFilters_Valid",
filters: ListPromotedPathsFilters{},
},
{
name: "KnownSignalAndContext_Valid",
filters: ListPromotedPathsFilters{Signal: "traces", Context: "attribute"},
},
{
name: "InvalidSignal_Rejected",
filters: ListPromotedPathsFilters{Signal: "events"},
wantErr: true,
},
{
name: "InvalidContext_Rejected",
filters: ListPromotedPathsFilters{Context: "json"},
wantErr: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
err := testCase.filters.Validate()
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
})
}
}
func TestListPromotedPathsFiltersMatch(t *testing.T) {
trueValue := true
falseValue := false
testCases := []struct {
name string
filters ListPromotedPathsFilters
target Target
path PromotePath
wantTarget bool
wantPath bool
}{
{
name: "EmptyFilters_MatchEverything",
filters: ListPromotedPathsFilters{},
target: NewTracesAttributesTarget(),
path: PromotePath{Path: "http.method", Promote: true},
wantTarget: true,
wantPath: true,
},
{
name: "SignalFilter_MatchesSameSignal",
filters: ListPromotedPathsFilters{Signal: "traces"},
target: NewTracesAttributesTarget(),
wantTarget: true,
},
{
name: "SignalFilter_SkipsOtherSignals",
filters: ListPromotedPathsFilters{Signal: "traces"},
target: NewLogsBodyTarget(),
wantTarget: false,
},
{
name: "ContextFilter_SkipsOtherContexts",
filters: ListPromotedPathsFilters{Context: "body"},
target: NewTracesAttributesTarget(),
wantTarget: false,
},
{
name: "PromotedFalseFilter_MatchesUnpromotedPath",
filters: ListPromotedPathsFilters{Promoted: &falseValue},
path: PromotePath{Path: "request.duration"},
wantPath: true,
},
{
name: "PromotedFalseFilter_SkipsPromotedPath",
filters: ListPromotedPathsFilters{Promoted: &falseValue},
path: PromotePath{Path: "http.method", Promote: true},
wantPath: false,
},
{
name: "IndexesTrueFilter_MatchesIndexedPath",
filters: ListPromotedPathsFilters{Indexes: &trueValue},
path: PromotePath{Path: "user.name", Indexes: []WrappedIndex{{Type: "minmax"}}},
wantPath: true,
},
{
name: "IndexesTrueFilter_SkipsUnindexedPath",
filters: ListPromotedPathsFilters{Indexes: &trueValue},
path: PromotePath{Path: "http.method", Promote: true},
wantPath: false,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
if testCase.target.Entry.Signal.StringValue() != "" {
assert.Equal(t, testCase.wantTarget, testCase.filters.MatchesTarget(testCase.target))
}
if testCase.path.Path != "" {
assert.Equal(t, testCase.wantPath, testCase.filters.MatchesPath(testCase.path))
}
})
}
}

View File

@@ -398,6 +398,19 @@ func NewTelemetryFieldKey(name string, fieldContext FieldContext, fieldDataType
}
}
// JSONIndexLookup locates one domain's JSON sub-column indexes in
// system.data_skipping_indices: the table to read, the base and promoted
// column prefixes to match index expressions against, and the signal and
// context to stamp on the results.
type JSONIndexLookup struct {
Signal Signal
FieldContext FieldContext
DBName string
LocalTableName string
BaseColumnPrefix string
PromotedColumnPrefix string
}
type TelemetryFieldKeySkipIndex struct {
Name string `json:"name"` // Name is TelemetryFieldKey.Name not IndexName from ClickHouse
FieldContext FieldContext `json:"fieldContext,omitzero"`

View File

@@ -1,6 +1,9 @@
package telemetrytypes
import "github.com/SigNoz/signoz/pkg/valuer"
import (
"github.com/SigNoz/signoz/pkg/types/coretypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
type Signal struct {
valuer.String
@@ -22,3 +25,30 @@ func (Signal) Enum() []any {
SignalUnspecified,
}
}
// FieldResource is the authz resource guarding the signal's field configuration.
func (signal Signal) FieldResource() (coretypes.Resource, bool) {
switch signal {
case SignalLogs:
return coretypes.ResourceMetaResourceLogsField, true
case SignalTraces:
return coretypes.ResourceMetaResourceTracesField, true
default:
return nil, false
}
}
// SignalFromText resolves a signal word to its Signal; ok is false for an
// unknown word.
func SignalFromText(text string) (Signal, bool) {
s := Signal{valuer.NewString(text)}
switch s {
case SignalTraces:
return SignalTraces, true
case SignalLogs:
return SignalLogs, true
case SignalMetrics:
return SignalMetrics, true
}
return Signal{}, false
}

View File

@@ -34,14 +34,19 @@ type MetadataStore interface {
FetchTemporalityAndTypeMulti(ctx context.Context, orgID valuer.UUID, queryTimeRangeStartTs, queryTimeRangeEndTs uint64, metricNames ...string) (map[string]metrictypes.Temporality, map[string]metrictypes.Type, map[string]bool, error)
// ListLogsJSONIndexes lists the JSON indexes for the logs table.
ListLogsJSONIndexes(ctx context.Context, filters ...string) ([]TelemetryFieldKeySkipIndex, error)
// ListJSONIndexes lists the per-path JSON skip indexes of the given source.
ListJSONIndexes(ctx context.Context, lookup JSONIndexLookup, filters ...string) ([]TelemetryFieldKeySkipIndex, error)
// ListPromotedPaths lists the promoted paths.
GetPromotedPaths(ctx context.Context, paths ...string) (map[string]bool, error)
// GetPromotedPaths lists the promoted paths recorded in the column
// evolution table for the entry's signal, column and field context.
GetPromotedPaths(ctx context.Context, entry EvolutionEntry, paths ...string) (map[string]bool, error)
// PromotePaths promotes the paths.
PromotePaths(ctx context.Context, paths ...string) error
// PromotePaths records promoted paths in the column evolution table as
// rows templated by entry; FieldName and ReleaseTime are set per path.
PromotePaths(ctx context.Context, entry EvolutionEntry, paths ...string) error
// GetMaterializedKeys lists the keys materialized as dedicated columns.
GetMaterializedKeys(ctx context.Context, signal Signal) ([]*TelemetryFieldKey, error)
// GetFirstSeenFromMetricMetadata gets the first seen timestamp for a metric metadata lookup key.
GetFirstSeenFromMetricMetadata(ctx context.Context, lookupKeys []MetricMetadataLookupKey) (map[MetricMetadataLookupKey]int64, error)

View File

@@ -20,6 +20,7 @@ type MockMetadataStore struct {
ReducedMap map[string]bool
PromotedPathsMap map[string]bool
LogsJSONIndexes []telemetrytypes.TelemetryFieldKeySkipIndex
MaterializedKeys []*telemetrytypes.TelemetryFieldKey
ColumnEvolutionMetadataMap map[string][]*telemetrytypes.EvolutionEntry
LookupKeysMap map[telemetrytypes.MetricMetadataLookupKey]int64
// StaticFields holds signal-specific intrinsic field definitions (e.g. logstelemetryschema.IntrinsicFields).
@@ -361,7 +362,7 @@ func (m *MockMetadataStore) SetTemporality(metricName string, temporality metric
}
// PromotePaths promotes the paths.
func (m *MockMetadataStore) PromotePaths(ctx context.Context, paths ...string) error {
func (m *MockMetadataStore) PromotePaths(_ context.Context, _ telemetrytypes.EvolutionEntry, paths ...string) error {
for _, path := range paths {
m.PromotedPathsMap[path] = true
}
@@ -369,13 +370,29 @@ func (m *MockMetadataStore) PromotePaths(ctx context.Context, paths ...string) e
}
// GetPromotedPaths returns the promoted paths.
func (m *MockMetadataStore) GetPromotedPaths(ctx context.Context, paths ...string) (map[string]bool, error) {
func (m *MockMetadataStore) GetPromotedPaths(_ context.Context, _ telemetrytypes.EvolutionEntry, _ ...string) (map[string]bool, error) {
return m.PromotedPathsMap, nil
}
// ListLogsJSONIndexes lists the JSON indexes for the logs table.
func (m *MockMetadataStore) ListLogsJSONIndexes(ctx context.Context, filters ...string) ([]telemetrytypes.TelemetryFieldKeySkipIndex, error) {
return m.LogsJSONIndexes, nil
func (m *MockMetadataStore) GetMaterializedKeys(_ context.Context, signal telemetrytypes.Signal) ([]*telemetrytypes.TelemetryFieldKey, error) {
keys := []*telemetrytypes.TelemetryFieldKey{}
for _, key := range m.MaterializedKeys {
if key.Signal == signal {
keys = append(keys, key)
}
}
return keys, nil
}
// ListJSONIndexes narrows the stored indexes to the lookup's field context like the real query.
func (m *MockMetadataStore) ListJSONIndexes(ctx context.Context, lookup telemetrytypes.JSONIndexLookup, filters ...string) ([]telemetrytypes.TelemetryFieldKeySkipIndex, error) {
indexes := []telemetrytypes.TelemetryFieldKeySkipIndex{}
for _, index := range m.LogsJSONIndexes {
if index.FieldContext.StringValue() == lookup.FieldContext.StringValue() {
indexes = append(indexes, index)
}
}
return indexes, nil
}
func (m *MockMetadataStore) updateColumnEvolutionMetadataForKeys(_ context.Context, keysToUpdate []*telemetrytypes.TelemetryFieldKey) map[string][]*telemetrytypes.EvolutionEntry {

View File

@@ -32,6 +32,7 @@ pytest_plugins = [
"fixtures.alerts",
"fixtures.cloudintegrations",
"fixtures.jsontypes",
"fixtures.materialized",
"fixtures.seeder",
"fixtures.serviceaccount",
"fixtures.role",

35
tests/fixtures/materialized.py vendored Normal file
View File

@@ -0,0 +1,35 @@
from collections.abc import Callable, Generator
import pytest
from fixtures import types
# signal -> (database, local table, distributed table)
MATERIALIZED_TABLES = {
"traces": ("signoz_traces", "signoz_index_v3", "distributed_signoz_index_v3"),
}
@pytest.fixture(name="materialize_attribute", scope="function")
def materialize_attribute(signoz: types.SigNoz) -> Generator[Callable[[str, str], str]]:
"""Teardown drops the columns and every index whose expression mentions a materialized key."""
cluster = signoz.telemetrystore.env["SIGNOZ_TELEMETRYSTORE_CLICKHOUSE_CLUSTER"]
created: list[tuple[str, str, str]] = []
def materialize(signal: str, key: str) -> str:
database, local_table, distributed_table = MATERIALIZED_TABLES[signal]
column = f"attribute_string_{key.replace('.', '$$')}"
for table in (local_table, distributed_table):
signoz.telemetrystore.conn.query(f"ALTER TABLE {database}.{table} ON CLUSTER '{cluster}' ADD COLUMN IF NOT EXISTS `{column}` LowCardinality(String) DEFAULT attributes_string['{key}']")
created.append((signal, key, column))
return column
yield materialize
for signal, key, column in created:
database, local_table, distributed_table = MATERIALIZED_TABLES[signal]
indexes = signoz.telemetrystore.conn.query(f"SELECT name FROM system.data_skipping_indices WHERE database = '{database}' AND table = '{local_table}' AND expr LIKE '%{key}%'").result_rows
for (index_name,) in indexes:
signoz.telemetrystore.conn.query(f"ALTER TABLE {database}.{local_table} ON CLUSTER '{cluster}' DROP INDEX IF EXISTS `{index_name}`")
for table in (local_table, distributed_table):
signoz.telemetrystore.conn.query(f"ALTER TABLE {database}.{table} ON CLUSTER '{cluster}' DROP COLUMN IF EXISTS `{column}`")

View File

@@ -0,0 +1,164 @@
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,
add_license,
change_user_role,
create_active_user,
)
from fixtures.role import find_role_by_name, transaction_group
PROMOTED_PATH_BASE = "/api/v1/promoted_path"
_EDITOR_EMAIL = "editor+promote@integration.test"
_VIEWER_EMAIL = "viewer+promote@integration.test"
_TRACES_ONLY_ROLE_NAME = "promote-traces-only"
_TRACES_ONLY_EMAIL = "customrole+promote@integration.test"
_PASSWORD = "password123Z$"
_TRACES_PATH = {"signal": "traces", "context": "attribute", "path": "authz.probe", "promote": True}
_LOGS_PATH = {"signal": "logs", "context": "body", "path": "authz.probe", "promote": True}
def test_apply_license(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list[Mapping]], None],
get_token: Callable[[str, str], str],
) -> None:
add_license(signoz, make_http_mocks, get_token)
def test_setup_users(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_role: Callable[..., str],
) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
create_active_user(signoz, admin_token, email=_EDITOR_EMAIL, role="signoz-editor", password=_PASSWORD, name="promote editor")
create_active_user(signoz, admin_token, email=_VIEWER_EMAIL, role="signoz-viewer", password=_PASSWORD, name="promote viewer")
create_role(
admin_token,
_TRACES_ONLY_ROLE_NAME,
[
transaction_group("update", "metaresource", "traces-field", ["*"]),
transaction_group("list", "metaresource", "traces-field", ["*"]),
],
)
user_id = create_active_user(signoz, admin_token, email=_TRACES_ONLY_EMAIL, role="signoz-viewer", password=_PASSWORD, name="promote traces only")
change_user_role(signoz, admin_token, user_id, "signoz-viewer", _TRACES_ONLY_ROLE_NAME)
def test_managed_roles(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
cases = [
(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD, HTTPStatus.CREATED),
(_EDITOR_EMAIL, _PASSWORD, HTTPStatus.CREATED),
(_VIEWER_EMAIL, _PASSWORD, HTTPStatus.FORBIDDEN),
]
for email, password, promote_status in cases:
token = get_token(email, password)
resp = requests.post(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
json=[_TRACES_PATH, _LOGS_PATH],
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == promote_status, f"{email} promote: expected {promote_status}, got {resp.status_code}: {resp.text}"
resp = requests.get(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.OK, f"{email} list: {resp.text}"
def test_promote_scoped_to_granted_signal(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
token = get_token(_TRACES_ONLY_EMAIL, _PASSWORD)
resp = requests.post(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
json=[_TRACES_PATH],
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.CREATED, f"promote traces: {resp.text}"
for body in ([_LOGS_PATH], [_TRACES_PATH, _LOGS_PATH]):
resp = requests.post(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
json=body,
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"promote {body}: expected 403, got {resp.status_code}: {resp.text}"
def test_list_scoped_to_granted_signal(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
token = get_token(_TRACES_ONLY_EMAIL, _PASSWORD)
resp = requests.get(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
params={"signal": "traces"},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.OK, f"list traces: {resp.text}"
assert {(path["signal"], path["path"]) for path in resp.json()["data"]} >= {("traces", "authz.probe")}
for params in ({"signal": "logs"}, {}):
resp = requests.get(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
params=params,
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"list {params}: expected 403, got {resp.status_code}: {resp.text}"
def test_revoke_update(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
role_id = find_role_by_name(signoz, admin_token, _TRACES_ONLY_ROLE_NAME)
resp = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/roles/{role_id}"),
json={"description": "", "transactionGroups": [transaction_group("list", "metaresource", "traces-field", ["*"])]},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.NO_CONTENT, resp.text
token = get_token(_TRACES_ONLY_EMAIL, _PASSWORD)
resp = requests.post(
signoz.self.host_configs["8080"].get(PROMOTED_PATH_BASE),
json=[_TRACES_PATH],
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"promote after revoke: expected 403, got {resp.status_code}: {resp.text}"

View File

@@ -0,0 +1,59 @@
from collections.abc import Callable
from http import HTTPStatus
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
MATERIALIZED_PATH = "/api/v1/promoted_path/materialized"
PROBE_INDEX_EXPRS_QUERY = "SELECT expr FROM system.data_skipping_indices WHERE database = 'signoz_traces' AND table = 'signoz_index_v3' AND expr LIKE 'CAST(%materialized.probe%'"
def test_index_materialized_paths(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
materialize_attribute: Callable[[str, str], str],
) -> None:
materialize_attribute("traces", "materialized.probe")
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
resp = requests.post(
signoz.self.host_configs["8080"].get(MATERIALIZED_PATH),
params={"dryRun": "true"},
headers={"Authorization": f"Bearer {token}"},
timeout=30,
)
assert resp.status_code == HTTPStatus.OK, resp.text
paths = {(path["signal"], path["context"], path["path"]): path["indexes"] for path in resp.json()["data"]}
assert paths[("traces", "attribute", "materialized.probe")] == [{"fieldDataType": "string", "type": "bloom_filter(0.01)", "granularity": 64}]
assert ("traces", "attribute", "http.route") not in paths
assert not signoz.telemetrystore.conn.query(PROBE_INDEX_EXPRS_QUERY).result_rows
resp = requests.post(
signoz.self.host_configs["8080"].get(MATERIALIZED_PATH),
params={"signal": "logs"},
headers={"Authorization": f"Bearer {token}"},
timeout=30,
)
assert resp.status_code == HTTPStatus.OK, resp.text
assert not [path for path in resp.json()["data"] if path["path"] == "materialized.probe"]
assert not signoz.telemetrystore.conn.query(PROBE_INDEX_EXPRS_QUERY).result_rows
resp = requests.post(
signoz.self.host_configs["8080"].get(MATERIALIZED_PATH),
headers={"Authorization": f"Bearer {token}"},
timeout=30,
)
assert resp.status_code == HTTPStatus.OK, resp.text
assert ("traces", "attribute", "materialized.probe") in {(path["signal"], path["context"], path["path"]) for path in resp.json()["data"]}
assert [row[0] for row in signoz.telemetrystore.conn.query(PROBE_INDEX_EXPRS_QUERY).result_rows] == ["CAST(attributes.materialized.probe, 'String')"]
resp = requests.post(
signoz.self.host_configs["8080"].get(MATERIALIZED_PATH),
headers={"Authorization": f"Bearer {token}"},
timeout=30,
)
assert resp.status_code == HTTPStatus.OK, resp.text
assert not [path for path in resp.json()["data"] if path["path"] == "materialized.probe"]