mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-30 23:30:40 +01:00
Compare commits
5 Commits
feat/saved
...
issue_6090
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
86878658ca | ||
|
|
e374d03e54 | ||
|
|
d445b6c296 | ||
|
|
fd8aaac300 | ||
|
|
5dae0b975a |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -68,6 +68,7 @@ jobs:
|
||||
- semconvfamilies
|
||||
- serviceaccount
|
||||
- spanmapper
|
||||
- tracedetail
|
||||
- querier_json_body
|
||||
- querier_skip_resource_fingerprint
|
||||
- ttl
|
||||
|
||||
@@ -1771,12 +1771,15 @@ components:
|
||||
additionalProperties: {}
|
||||
nullable: true
|
||||
type: object
|
||||
syncState:
|
||||
$ref: '#/components/schemas/CloudintegrationtypesSyncState'
|
||||
timestampMillis:
|
||||
format: int64
|
||||
type: integer
|
||||
required:
|
||||
- timestampMillis
|
||||
- data
|
||||
- syncState
|
||||
type: object
|
||||
CloudintegrationtypesAzureAccountConfig:
|
||||
properties:
|
||||
@@ -2021,6 +2024,8 @@ components:
|
||||
format: date-time
|
||||
nullable: true
|
||||
type: string
|
||||
syncState:
|
||||
$ref: '#/components/schemas/CloudintegrationtypesSyncState'
|
||||
required:
|
||||
- account_id
|
||||
- cloud_account_id
|
||||
@@ -2030,6 +2035,7 @@ components:
|
||||
- providerAccountId
|
||||
- integrationConfig
|
||||
- removedAt
|
||||
- syncState
|
||||
type: object
|
||||
CloudintegrationtypesGettableServicesMetadata:
|
||||
properties:
|
||||
@@ -2129,6 +2135,9 @@ components:
|
||||
type: object
|
||||
providerAccountId:
|
||||
type: string
|
||||
syncedVersion:
|
||||
nullable: true
|
||||
type: integer
|
||||
required:
|
||||
- data
|
||||
type: object
|
||||
@@ -2141,6 +2150,18 @@ components:
|
||||
gcp:
|
||||
$ref: '#/components/schemas/CloudintegrationtypesGCPIntegrationConfig'
|
||||
type: object
|
||||
CloudintegrationtypesRegionState:
|
||||
enum:
|
||||
- enabled
|
||||
- disabled
|
||||
type: string
|
||||
CloudintegrationtypesRegionSyncState:
|
||||
properties:
|
||||
state:
|
||||
$ref: '#/components/schemas/CloudintegrationtypesRegionState'
|
||||
required:
|
||||
- state
|
||||
type: object
|
||||
CloudintegrationtypesService:
|
||||
properties:
|
||||
assets:
|
||||
@@ -2282,6 +2303,23 @@ components:
|
||||
metrics:
|
||||
type: boolean
|
||||
type: object
|
||||
CloudintegrationtypesSyncState:
|
||||
nullable: true
|
||||
properties:
|
||||
inSync:
|
||||
type: boolean
|
||||
regions:
|
||||
additionalProperties:
|
||||
$ref: '#/components/schemas/CloudintegrationtypesRegionSyncState'
|
||||
type: object
|
||||
version:
|
||||
format: int64
|
||||
type: integer
|
||||
required:
|
||||
- version
|
||||
- inSync
|
||||
- regions
|
||||
type: object
|
||||
CloudintegrationtypesUpdatableAccount:
|
||||
properties:
|
||||
config:
|
||||
@@ -9655,6 +9693,29 @@ components:
|
||||
required:
|
||||
- aggregations
|
||||
type: object
|
||||
SpantypesGettableTraceSummary:
|
||||
properties:
|
||||
ai:
|
||||
$ref: '#/components/schemas/SpantypesTraceAISummary'
|
||||
endTimestampMillis:
|
||||
minimum: 0
|
||||
type: integer
|
||||
hasMissingSpans:
|
||||
type: boolean
|
||||
rootServiceEntryPoint:
|
||||
type: string
|
||||
rootServiceName:
|
||||
type: string
|
||||
startTimestampMillis:
|
||||
minimum: 0
|
||||
type: integer
|
||||
totalErrorSpansCount:
|
||||
minimum: 0
|
||||
type: integer
|
||||
totalSpansCount:
|
||||
minimum: 0
|
||||
type: integer
|
||||
type: object
|
||||
SpantypesGettableWaterfallTrace:
|
||||
properties:
|
||||
endTimestampMillis:
|
||||
@@ -9969,6 +10030,32 @@ components:
|
||||
nullable: true
|
||||
type: object
|
||||
type: object
|
||||
SpantypesTraceAISummary:
|
||||
properties:
|
||||
tokens:
|
||||
$ref: '#/components/schemas/SpantypesTraceAITokens'
|
||||
totalCost:
|
||||
nullable: true
|
||||
type: number
|
||||
type: object
|
||||
SpantypesTraceAITokens:
|
||||
properties:
|
||||
cacheRead:
|
||||
minimum: 0
|
||||
type: integer
|
||||
cacheWrite:
|
||||
minimum: 0
|
||||
type: integer
|
||||
input:
|
||||
minimum: 0
|
||||
type: integer
|
||||
output:
|
||||
minimum: 0
|
||||
type: integer
|
||||
reasoning:
|
||||
minimum: 0
|
||||
type: integer
|
||||
type: object
|
||||
SpantypesUpdatableSpanMapper:
|
||||
properties:
|
||||
config:
|
||||
@@ -15693,6 +15780,67 @@ paths:
|
||||
tags:
|
||||
- tracedetail
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/traces/{traceID}/summary:
|
||||
get:
|
||||
deprecated: false
|
||||
description: Returns the trace-level fields of the waterfall (time range, root,
|
||||
span counts, missing spans) and, when the trace has gen_ai spans, its token
|
||||
and cost totals. Computed in one aggregate query.
|
||||
operationId: GetTraceSummary
|
||||
parameters:
|
||||
- in: path
|
||||
name: traceID
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/SpantypesGettableTraceSummary'
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"404":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Not Found
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: Get summary for a trace
|
||||
tags:
|
||||
- tracedetail
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/user/me:
|
||||
get:
|
||||
deprecated: true
|
||||
|
||||
@@ -118,7 +118,7 @@ router.Handle("/api/v1/service_accounts", handler.New(
|
||||
The pieces:
|
||||
|
||||
- **`CheckResources(handlerFn, roles...)`** — the resource-aware authorization wrapper from [pkg/http/middleware/authz.go](/pkg/http/middleware/authz.go). The role list is the community-edition fallback: which managed roles may call this route when per-resource checks are unavailable.
|
||||
- **`ResourceDef`** — declares the resource, verb, audit category, how to extract the instance ID, and how to turn that ID into selectors. ID extractors live in [pkg/types/coretypes/extractor.go](/pkg/types/coretypes/extractor.go): `PathParam("id")`, `BodyField(func(req *T) string)` / `BodyFields(func(req *T) []string)` reading the request body the resource middleware decoded into the route's `OpenAPIDef.Request` type `T`, and `ResponseJSONPath("data.id")` for IDs only known after the handler runs (e.g. `create`). A handler on such a route reads the same decoded value with `coretypes.BodyFromContext[T](r.Context())`.
|
||||
- **`ResourceDef`** — declares the resource, verb, audit category, how to extract the instance ID, and how to turn that ID into selectors. ID extractors live in [pkg/types/coretypes/extractor.go](/pkg/types/coretypes/extractor.go): `PathParam("id")`, `BodyJSONPath("data.id")`, `BodyJSONArray("ids")`, and `ResponseJSONPath("data.id")` for IDs only known after the handler runs (e.g. `create`).
|
||||
- **`SecuritySchemes`** — advertises the required scope (`resource.Scope(verb)`, e.g. `serviceaccount:create`) in the OpenAPI spec.
|
||||
|
||||
For routes that link two resources, use `AttachDetachSiblingResourceDef` (both sides are authz-checked, e.g. attaching a role to a service account requires `attach` on **both** the service account and the role). For parent-child routes (e.g. creating an API key under a service account), both sides are checked too, but with different verbs: declare a `BasicResourceDef` checking the child with `create`/`delete`, alongside an `AttachDetachParentChildResourceDef` checking the parent with `attach`/`detach` (within that def the child is only recorded for audit) — see the `/api/v1/service_accounts/{id}/keys` route in [pkg/apiserver/signozapiserver/serviceaccount.go](/pkg/apiserver/signozapiserver/serviceaccount.go).
|
||||
|
||||
@@ -183,32 +183,52 @@ func (module *module) AgentCheckIn(ctx context.Context, orgID valuer.UUID, provi
|
||||
return nil, errors.New(errors.TypeAlreadyExists, cloudintegrationtypes.ErrCodeCloudIntegrationAlreadyConnected, errMessage)
|
||||
}
|
||||
|
||||
account, err := module.store.GetAccountByID(ctx, orgID, req.CloudIntegrationID, provider)
|
||||
storableAccount, err := module.store.GetAccountByID(ctx, orgID, req.CloudIntegrationID, provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
account, err := cloudintegrationtypes.NewAccountFromStorable(storableAccount)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
syncState := account.NextSyncState(req.SyncedVersion)
|
||||
|
||||
// If account has been removed (disconnected), return a minimal response with empty integration config.
|
||||
// The agent uses this response to clean up resources
|
||||
if account.RemovedAt != nil {
|
||||
// Heartbeat stays frozen after removal, only the sync state is updated.
|
||||
if account.AgentReport != nil && syncState != nil {
|
||||
account.UpdateSyncState(syncState)
|
||||
|
||||
storableAccount, err = cloudintegrationtypes.NewStorableCloudIntegration(account)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = module.store.UpdateAgentReport(ctx, storableAccount)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
return cloudintegrationtypes.NewAgentCheckInResponse(
|
||||
req.ProviderAccountID,
|
||||
account.ID.StringValue(),
|
||||
new(cloudintegrationtypes.ProviderIntegrationConfig),
|
||||
account.RemovedAt,
|
||||
syncState,
|
||||
), nil
|
||||
}
|
||||
|
||||
// update account with cloud provider account id and agent report (heartbeat)
|
||||
account.Update(&req.ProviderAccountID, cloudintegrationtypes.NewAgentReport(req.Data))
|
||||
account.UpdateAgentReport(&req.ProviderAccountID, cloudintegrationtypes.NewAgentReport(req.Data, syncState))
|
||||
|
||||
err = module.store.UpdateAccount(ctx, account)
|
||||
storableAccount, err = cloudintegrationtypes.NewStorableCloudIntegration(account)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Get account as domain object for config access (enabled regions, etc.)
|
||||
domainAccount, err := cloudintegrationtypes.NewAccountFromStorable(account)
|
||||
err = module.store.UpdateAgentReport(ctx, storableAccount)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -223,8 +243,7 @@ func (module *module) AgentCheckIn(ctx context.Context, orgID valuer.UUID, provi
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Delegate integration config building entirely to the provider module
|
||||
integrationConfig, err := cloudProvider.BuildIntegrationConfig(ctx, domainAccount, storedServices)
|
||||
integrationConfig, err := cloudProvider.BuildIntegrationConfig(ctx, account, storedServices)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -234,6 +253,7 @@ func (module *module) AgentCheckIn(ctx context.Context, orgID valuer.UUID, provi
|
||||
account.ID.StringValue(),
|
||||
integrationConfig,
|
||||
account.RemovedAt,
|
||||
syncState,
|
||||
), nil
|
||||
}
|
||||
|
||||
|
||||
@@ -3396,6 +3396,37 @@ export interface CloudintegrationtypesAWSServiceConfigDTO {
|
||||
metrics?: CloudintegrationtypesAWSServiceMetricsConfigDTO;
|
||||
}
|
||||
|
||||
export enum CloudintegrationtypesRegionStateDTO {
|
||||
enabled = 'enabled',
|
||||
disabled = 'disabled',
|
||||
}
|
||||
export interface CloudintegrationtypesRegionSyncStateDTO {
|
||||
state: CloudintegrationtypesRegionStateDTO;
|
||||
}
|
||||
|
||||
export type CloudintegrationtypesSyncStateDTORegions = {
|
||||
[key: string]: CloudintegrationtypesRegionSyncStateDTO;
|
||||
};
|
||||
|
||||
/**
|
||||
* @nullable
|
||||
*/
|
||||
export type CloudintegrationtypesSyncStateDTO = {
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
inSync: boolean;
|
||||
/**
|
||||
* @type object
|
||||
*/
|
||||
regions: CloudintegrationtypesSyncStateDTORegions;
|
||||
/**
|
||||
* @type integer
|
||||
* @format int64
|
||||
*/
|
||||
version: number;
|
||||
} | null;
|
||||
|
||||
export type CloudintegrationtypesAgentReportDTODataAnyOf = {
|
||||
[key: string]: unknown;
|
||||
};
|
||||
@@ -3414,6 +3445,7 @@ export type CloudintegrationtypesAgentReportDTO = {
|
||||
* @type object,null
|
||||
*/
|
||||
data: CloudintegrationtypesAgentReportDTOData;
|
||||
syncState: CloudintegrationtypesSyncStateDTO | null;
|
||||
/**
|
||||
* @type integer
|
||||
* @format int64
|
||||
@@ -3842,6 +3874,7 @@ export interface CloudintegrationtypesGettableAgentCheckInDTO {
|
||||
* @format date-time
|
||||
*/
|
||||
removedAt: string | null;
|
||||
syncState: CloudintegrationtypesSyncStateDTO | null;
|
||||
}
|
||||
|
||||
export interface CloudintegrationtypesServiceMetadataDTO {
|
||||
@@ -3912,6 +3945,10 @@ export interface CloudintegrationtypesPostableAgentCheckInDTO {
|
||||
* @type string
|
||||
*/
|
||||
providerAccountId?: string;
|
||||
/**
|
||||
* @type integer,null
|
||||
*/
|
||||
syncedVersion?: number | null;
|
||||
}
|
||||
|
||||
export interface CloudintegrationtypesStorableIntegrationDashboardDTO {
|
||||
@@ -11157,6 +11194,78 @@ export interface SpantypesGettableTraceAggregationsDTO {
|
||||
aggregations: SpantypesSpanAggregationResultDTO[];
|
||||
}
|
||||
|
||||
export interface SpantypesTraceAITokensDTO {
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
cacheRead?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
cacheWrite?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
input?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
output?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
reasoning?: number;
|
||||
}
|
||||
|
||||
export interface SpantypesTraceAISummaryDTO {
|
||||
tokens?: SpantypesTraceAITokensDTO;
|
||||
/**
|
||||
* @type number,null
|
||||
*/
|
||||
totalCost?: number | null;
|
||||
}
|
||||
|
||||
export interface SpantypesGettableTraceSummaryDTO {
|
||||
ai?: SpantypesTraceAISummaryDTO;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
endTimestampMillis?: number;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
hasMissingSpans?: boolean;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
rootServiceEntryPoint?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
rootServiceName?: string;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
startTimestampMillis?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
totalErrorSpansCount?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
totalSpansCount?: number;
|
||||
}
|
||||
|
||||
export interface SpantypesOtelSpanRefDTO {
|
||||
/**
|
||||
* @type string
|
||||
@@ -12845,6 +12954,17 @@ export type GetTraceAggregations200 = {
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type GetTraceSummaryPathParameters = {
|
||||
traceID: string;
|
||||
};
|
||||
export type GetTraceSummary200 = {
|
||||
data: SpantypesGettableTraceSummaryDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListUserPreferences200 = {
|
||||
/**
|
||||
* @type array
|
||||
|
||||
@@ -4,11 +4,17 @@
|
||||
* * regenerate with 'pnpm generate:api'
|
||||
* SigNoz
|
||||
*/
|
||||
import { useMutation } from 'react-query';
|
||||
import { useMutation, useQuery } from 'react-query';
|
||||
import type {
|
||||
InvalidateOptions,
|
||||
MutationFunction,
|
||||
QueryClient,
|
||||
QueryFunction,
|
||||
QueryKey,
|
||||
UseMutationOptions,
|
||||
UseMutationResult,
|
||||
UseQueryOptions,
|
||||
UseQueryResult,
|
||||
} from 'react-query';
|
||||
|
||||
import type {
|
||||
@@ -16,6 +22,8 @@ import type {
|
||||
GetFlamegraphPathParameters,
|
||||
GetTraceAggregations200,
|
||||
GetTraceAggregationsPathParameters,
|
||||
GetTraceSummary200,
|
||||
GetTraceSummaryPathParameters,
|
||||
GetWaterfallV4200,
|
||||
GetWaterfallV4PathParameters,
|
||||
RenderErrorResponseDTO,
|
||||
@@ -27,6 +35,26 @@ 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;
|
||||
};
|
||||
|
||||
/**
|
||||
* Computes span aggregations grouped by requested field.
|
||||
* @summary Get aggregations for a trace
|
||||
@@ -127,6 +155,108 @@ export const useGetTraceAggregations = <
|
||||
> => {
|
||||
return useMutation(getGetTraceAggregationsMutationOptions(options));
|
||||
};
|
||||
/**
|
||||
* Returns the trace-level fields of the waterfall (time range, root, span counts, missing spans) and, when the trace has gen_ai spans, its token and cost totals. Computed in one aggregate query.
|
||||
* @summary Get summary for a trace
|
||||
*/
|
||||
export const getTraceSummary = (
|
||||
{ traceID }: GetTraceSummaryPathParameters,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<GetTraceSummary200>({
|
||||
url: `/api/v1/traces/${traceID}/summary`,
|
||||
method: 'GET',
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getGetTraceSummaryQueryKey = ({
|
||||
traceID,
|
||||
}: GetTraceSummaryPathParameters) => {
|
||||
return [`/api/v1/traces/${traceID}/summary`] as const;
|
||||
};
|
||||
|
||||
export const getGetTraceSummaryQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof getTraceSummary>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
{ traceID }: GetTraceSummaryPathParameters,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getTraceSummary>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey =
|
||||
queryOptions?.queryKey ?? getGetTraceSummaryQueryKey({ traceID });
|
||||
|
||||
const queryFn: QueryFunction<Awaited<ReturnType<typeof getTraceSummary>>> = ({
|
||||
signal,
|
||||
}) => getTraceSummary({ traceID }, signal);
|
||||
|
||||
return {
|
||||
queryKey,
|
||||
queryFn,
|
||||
enabled: traceID !== null && traceID !== undefined,
|
||||
...queryOptions,
|
||||
} as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getTraceSummary>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type GetTraceSummaryQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof getTraceSummary>>
|
||||
>;
|
||||
export type GetTraceSummaryQueryError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Get summary for a trace
|
||||
*/
|
||||
|
||||
export function useGetTraceSummary<
|
||||
TData = Awaited<ReturnType<typeof getTraceSummary>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
{ traceID }: GetTraceSummaryPathParameters,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getTraceSummary>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getGetTraceSummaryQueryOptions({ traceID }, options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return withQueryKey(query, queryOptions.queryKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary Get summary for a trace
|
||||
*/
|
||||
export const invalidateGetTraceSummary = async (
|
||||
queryClient: QueryClient,
|
||||
{ traceID }: GetTraceSummaryPathParameters,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getGetTraceSummaryQueryKey({ traceID }) },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* Returns the flamegraph view of spans for a given trace ID.
|
||||
* @summary Get flamegraph view for a trace
|
||||
|
||||
@@ -24,6 +24,7 @@ const accountsResponse: ListAccounts200 = {
|
||||
agentReport: {
|
||||
timestampMillis: 1747114366214,
|
||||
data: null,
|
||||
syncState: null,
|
||||
},
|
||||
providerAccountId: PROVIDER_ACCOUNT_ID,
|
||||
removedAt: null,
|
||||
|
||||
@@ -295,7 +295,11 @@ const account = (
|
||||
provider,
|
||||
providerAccountId: ACCOUNTS[provider][index],
|
||||
config: accountConfig(provider),
|
||||
agentReport: { timestampMillis: Date.now() - 45 * 1000, data: null },
|
||||
agentReport: {
|
||||
timestampMillis: Date.now() - 45 * 1000,
|
||||
data: null,
|
||||
syncState: null,
|
||||
},
|
||||
createdAt: new Date(Date.now() - 21 * 24 * 60 * 60 * 1000).toISOString(),
|
||||
updatedAt: new Date(Date.now() - 60 * 60 * 1000).toISOString(),
|
||||
removedAt: null,
|
||||
|
||||
@@ -1,15 +1,18 @@
|
||||
package signozapiserver
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"slices"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/http/handler"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/tidwall/gjson"
|
||||
)
|
||||
|
||||
func (provider *provider) addAuthDomainRoutes(router *mux.Router) error {
|
||||
@@ -74,7 +77,7 @@ func (provider *provider) addAuthDomainRoutes(router *mux.Router) error {
|
||||
SourceIDs: coretypes.OneID(coretypes.ResponseJSONPath("data.id")),
|
||||
SourceSelector: coretypes.WildcardSelector,
|
||||
TargetResource: coretypes.ResourceRole,
|
||||
TargetIDs: authDomainPostableRoleNamesExtractor(),
|
||||
TargetIDs: authDomainRoleNamesExtractor(),
|
||||
TargetSelector: coretypes.IDSelector,
|
||||
},
|
||||
),
|
||||
@@ -146,7 +149,7 @@ func (provider *provider) addAuthDomainRoutes(router *mux.Router) error {
|
||||
SourceIDs: coretypes.OneID(coretypes.PathParam("id")),
|
||||
SourceSelector: coretypes.IDSelector,
|
||||
TargetResource: coretypes.ResourceRole,
|
||||
TargetIDs: authDomainUpdatableRoleNamesExtractor(),
|
||||
TargetIDs: authDomainRoleNamesExtractor(),
|
||||
TargetSelector: coretypes.IDSelector,
|
||||
},
|
||||
handler.AttachDetachSiblingResourceDef{
|
||||
@@ -196,16 +199,20 @@ func (provider *provider) addAuthDomainRoutes(router *mux.Router) error {
|
||||
|
||||
// The extracted names are the roles the request body's mapping grants at SSO
|
||||
// login — see authDomainEffectiveRoleNames.
|
||||
func authDomainPostableRoleNamesExtractor() coretypes.ResourceIDsExtractor {
|
||||
return coretypes.BodyFields(func(req *authtypes.PostableAuthDomain) []string {
|
||||
return authDomainEffectiveRoleNames(req.RoleMapping)
|
||||
})
|
||||
}
|
||||
func authDomainRoleNamesExtractor() coretypes.ResourceIDsExtractor {
|
||||
return coretypes.ResourceIDsExtractor{Phase: coretypes.PhaseRequest, Fn: func(ec coretypes.ExtractorContext) ([]string, error) {
|
||||
roleMappingJSON := gjson.GetBytes(ec.RequestBody, "roleMapping")
|
||||
if !roleMappingJSON.Exists() || roleMappingJSON.Type == gjson.Null {
|
||||
return authDomainEffectiveRoleNames(nil), nil
|
||||
}
|
||||
|
||||
func authDomainUpdatableRoleNamesExtractor() coretypes.ResourceIDsExtractor {
|
||||
return coretypes.BodyFields(func(req *authtypes.UpdatableAuthDomain) []string {
|
||||
return authDomainEffectiveRoleNames(req.RoleMapping)
|
||||
})
|
||||
roleMapping := new(authtypes.RoleMapping)
|
||||
if err := json.Unmarshal([]byte(roleMappingJSON.Raw), roleMapping); err != nil {
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid role mapping: %v", err)
|
||||
}
|
||||
|
||||
return authDomainEffectiveRoleNames(roleMapping), nil
|
||||
}}
|
||||
}
|
||||
|
||||
// The extracted names are the roles the stored domain's mapping grants at SSO
|
||||
|
||||
@@ -350,7 +350,7 @@ func (provider *provider) addCloudIntegrationRoutes(router *mux.Router) error {
|
||||
Resource: coretypes.ResourceMetaResourceCloudIntegration,
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
ID: coretypes.BodyField(func(req *citypes.PostableAgentCheckIn) string { return req.ID }),
|
||||
ID: coretypes.BodyJSONPath("account_id"),
|
||||
Selector: coretypes.IDSelector,
|
||||
}),
|
||||
)).Methods(http.MethodPost).GetError(); err != nil {
|
||||
@@ -377,12 +377,7 @@ func (provider *provider) addCloudIntegrationRoutes(router *mux.Router) error {
|
||||
Resource: coretypes.ResourceMetaResourceCloudIntegration,
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
ID: coretypes.BodyField(func(req *citypes.PostableAgentCheckIn) string {
|
||||
if req.CloudIntegrationID.IsZero() {
|
||||
return ""
|
||||
}
|
||||
return req.CloudIntegrationID.StringValue()
|
||||
}),
|
||||
ID: coretypes.BodyJSONPath("cloudIntegrationId"),
|
||||
Selector: coretypes.IDSelector,
|
||||
}),
|
||||
)).Methods(http.MethodPost).GetError(); err != nil {
|
||||
|
||||
@@ -332,7 +332,7 @@ func (provider *provider) addGatewayRoutes(router *mux.Router) error {
|
||||
Verb: coretypes.VerbAttach,
|
||||
Category: coretypes.ActionCategoryConfigurationChange,
|
||||
ParentResource: coretypes.ResourceMetaResourceIngestionKey,
|
||||
ParentID: coretypes.BodyField(func(req *gatewaytypes.PostableIngestionKeyLimit) string { return req.KeyID }),
|
||||
ParentID: coretypes.BodyJSONPath("keyId"),
|
||||
ParentSelector: coretypes.IDSelector,
|
||||
ChildResource: coretypes.ResourceMetaResourceIngestionLimit,
|
||||
ChildIDs: coretypes.OneID(coretypes.ResponseJSONPath("data.id")),
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"strings"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/handler"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/prometheus"
|
||||
@@ -80,14 +79,6 @@ func (h *prometheusOpenAPIHandler) ResourceDefs() []handler.ResourceDef {
|
||||
}}
|
||||
}
|
||||
|
||||
func (h *prometheusOpenAPIHandler) Request() any {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *prometheusOpenAPIHandler) BindBodyOptions() []binding.BindBodyOption {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (provider *provider) addPrometheusRoutes(router *mux.Router) error {
|
||||
if err := router.Handle("/prometheus/api/v1/query", &prometheusOpenAPIHandler{
|
||||
handlerFunc: provider.authzMiddleware.CheckResources(provider.prometheusHandler.Query, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName),
|
||||
|
||||
@@ -461,11 +461,10 @@ func (provider *provider) addQuerierRoutes(router *mux.Router) error {
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest},
|
||||
SecuritySchemes: newScopedSecuritySchemes(telemetryReadScopes()),
|
||||
}, handler.WithResourceDefs(handler.TelemetryResourceDef{
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
Selector: querybuilder.TelemetrySelector,
|
||||
Resources: querybuilder.QueryRangeResources,
|
||||
RequiresBody: true,
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
Selector: querybuilder.TelemetrySelector,
|
||||
Resources: querybuilder.QueryRangeResources,
|
||||
}))).Methods(http.MethodPost).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -484,11 +483,10 @@ func (provider *provider) addQuerierRoutes(router *mux.Router) error {
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest},
|
||||
SecuritySchemes: newScopedSecuritySchemes(telemetryReadScopes()),
|
||||
}, handler.WithResourceDefs(handler.TelemetryResourceDef{
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
Selector: querybuilder.TelemetrySelector,
|
||||
Resources: querybuilder.QueryRangeResources,
|
||||
RequiresBody: true,
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
Selector: querybuilder.TelemetrySelector,
|
||||
Resources: querybuilder.QueryRangeResources,
|
||||
}))).Methods(http.MethodPost).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"net/http"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
pkghandler "github.com/SigNoz/signoz/pkg/http/handler"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/gorilla/mux"
|
||||
@@ -56,14 +55,6 @@ func (handler *healthOpenAPIHandler) ResourceDefs() []pkghandler.ResourceDef {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (handler *healthOpenAPIHandler) Request() any {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (handler *healthOpenAPIHandler) BindBodyOptions() []binding.BindBodyOption {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (provider *provider) addRegistryRoutes(router *mux.Router) error {
|
||||
if err := router.Handle("/api/v2/healthz", newHealthOpenAPIHandler(
|
||||
provider.authzMiddleware.OpenAccess(provider.factoryHandler.Healthz),
|
||||
|
||||
@@ -358,20 +358,10 @@ func (provider *provider) addServiceAccountRoutes(router *mux.Router) error {
|
||||
Verb: coretypes.VerbAttach,
|
||||
Category: coretypes.ActionCategoryAccessControl,
|
||||
SourceResource: coretypes.ResourceServiceAccount,
|
||||
SourceIDs: coretypes.OneID(coretypes.BodyField(func(req *serviceaccounttypes.PostableServiceAccountRole) string {
|
||||
if req.ServiceAccountID.IsZero() {
|
||||
return ""
|
||||
}
|
||||
return req.ServiceAccountID.StringValue()
|
||||
})),
|
||||
SourceIDs: coretypes.OneID(coretypes.BodyJSONPath("serviceAccountId")),
|
||||
SourceSelector: coretypes.IDSelector,
|
||||
TargetResource: coretypes.ResourceRole,
|
||||
TargetIDs: coretypes.OneID(coretypes.BodyField(func(req *serviceaccounttypes.PostableServiceAccountRole) string {
|
||||
if req.RoleID.IsZero() {
|
||||
return ""
|
||||
}
|
||||
return req.RoleID.StringValue()
|
||||
})),
|
||||
TargetIDs: coretypes.OneID(coretypes.BodyJSONPath("roleId")),
|
||||
TargetSelector: provider.roleSelector,
|
||||
}),
|
||||
)).Methods(http.MethodPost).GetError(); err != nil {
|
||||
|
||||
@@ -10,6 +10,23 @@ import (
|
||||
)
|
||||
|
||||
func (provider *provider) addTraceDetailRoutes(router *mux.Router) error {
|
||||
if err := router.Handle("/api/v1/traces/{traceID}/summary", handler.New(
|
||||
provider.authzMiddleware.ViewAccess(provider.traceDetailHandler.GetTraceSummary),
|
||||
handler.OpenAPIDef{
|
||||
ID: "GetTraceSummary",
|
||||
Tags: []string{"tracedetail"},
|
||||
Summary: "Get summary for a trace",
|
||||
Description: "Returns the trace-level fields of the waterfall (time range, root, span counts, missing spans) and, when the trace has gen_ai spans, its token and cost totals. Computed in one aggregate query.",
|
||||
Response: new(spantypes.GettableTraceSummary),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusOK,
|
||||
ErrorStatusCodes: []int{http.StatusNotFound},
|
||||
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
|
||||
},
|
||||
)).Methods(http.MethodGet).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v4/traces/{traceID}/waterfall", handler.New(
|
||||
provider.authzMiddleware.ViewAccess(provider.traceDetailHandler.GetWaterfallV4),
|
||||
handler.OpenAPIDef{
|
||||
|
||||
@@ -68,7 +68,7 @@ func (provider *provider) addZeusRoutes(router *mux.Router) error {
|
||||
Resource: coretypes.ResourceMetaResourceDeploymentHost,
|
||||
Verb: coretypes.VerbUpdate,
|
||||
Category: coretypes.ActionCategoryConfigurationChange,
|
||||
ID: coretypes.BodyField(func(req *zeustypes.PostableHost) string { return req.Name }),
|
||||
ID: coretypes.BodyJSONPath("name"),
|
||||
Selector: coretypes.WildcardSelector,
|
||||
}))).Methods(http.MethodPut).GetError(); err != nil {
|
||||
return err
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/gatewaytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
@@ -285,8 +284,8 @@ func (handler *handler) CreateIngestionKeyLimit(rw http.ResponseWriter, r *http.
|
||||
|
||||
orgID := valuer.MustNewUUID(claims.OrgID)
|
||||
|
||||
req, err := coretypes.BodyFromContext[gatewaytypes.PostableIngestionKeyLimit](r.Context())
|
||||
if err != nil {
|
||||
var req gatewaytypes.PostableIngestionKeyLimit
|
||||
if err := binding.JSON.BindBody(r.Body, &req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1,13 +1,10 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"reflect"
|
||||
"slices"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/swaggest/openapi-go"
|
||||
"github.com/swaggest/openapi-go/openapi3"
|
||||
@@ -19,15 +16,12 @@ type Handler interface {
|
||||
http.Handler
|
||||
ServeOpenAPI(openapi.OperationContext)
|
||||
ResourceDefs() []ResourceDef
|
||||
Request() any
|
||||
BindBodyOptions() []binding.BindBodyOption
|
||||
}
|
||||
|
||||
type handler struct {
|
||||
handlerFunc http.HandlerFunc
|
||||
openAPIDef OpenAPIDef
|
||||
resourceDefs []ResourceDef
|
||||
bindBodyOptions []binding.BindBodyOption
|
||||
handlerFunc http.HandlerFunc
|
||||
openAPIDef OpenAPIDef
|
||||
resourceDefs []ResourceDef
|
||||
}
|
||||
|
||||
func New(handlerFunc http.HandlerFunc, openAPIDef OpenAPIDef, opts ...Option) Handler {
|
||||
@@ -53,10 +47,6 @@ func New(handlerFunc http.HandlerFunc, openAPIDef OpenAPIDef, opts ...Option) Ha
|
||||
opt(handler)
|
||||
}
|
||||
|
||||
if RequiresBody(handler.resourceDefs) && (openAPIDef.Request == nil || reflect.TypeOf(openAPIDef.Request).Kind() != reflect.Pointer) {
|
||||
panic(fmt.Sprintf("handler %s: a body extractor needs OpenAPIDef.Request to be a pointer, got %T", openAPIDef.ID, openAPIDef.Request))
|
||||
}
|
||||
|
||||
return handler
|
||||
}
|
||||
|
||||
@@ -145,11 +135,3 @@ func (handler *handler) ServeOpenAPI(opCtx openapi.OperationContext) {
|
||||
func (handler *handler) ResourceDefs() []ResourceDef {
|
||||
return handler.resourceDefs
|
||||
}
|
||||
|
||||
func (handler *handler) Request() any {
|
||||
return handler.openAPIDef.Request
|
||||
}
|
||||
|
||||
func (handler *handler) BindBodyOptions() []binding.BindBodyOption {
|
||||
return handler.bindBodyOptions
|
||||
}
|
||||
|
||||
@@ -4,8 +4,6 @@ import (
|
||||
"net/http"
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -24,41 +22,6 @@ func (bespokeOpenAPIHandler) ServeOpenAPI(opCtx openapi.OperationContext) {
|
||||
|
||||
func (bespokeOpenAPIHandler) ResourceDefs() []ResourceDef { return nil }
|
||||
|
||||
func (bespokeOpenAPIHandler) Request() any { return nil }
|
||||
|
||||
func (bespokeOpenAPIHandler) BindBodyOptions() []binding.BindBodyOption { return nil }
|
||||
|
||||
func TestNewPanicsWhenBodyExtractorHasNoPointerRequest(t *testing.T) {
|
||||
type body struct{ ID string }
|
||||
bodyDef := BasicResourceDef{Resource: coretypes.ResourceRole, Verb: coretypes.VerbRead, ID: coretypes.BodyField(func(req *body) string { return req.ID }), Selector: coretypes.IDSelector}
|
||||
pathDef := BasicResourceDef{Resource: coretypes.ResourceRole, Verb: coretypes.VerbRead, ID: coretypes.PathParam("id"), Selector: coretypes.IDSelector}
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
request any
|
||||
def ResourceDef
|
||||
panics bool
|
||||
}{
|
||||
{name: "BodyExtractor_ValueRequest_Panics", request: body{}, def: bodyDef, panics: true},
|
||||
{name: "BodyExtractor_NilRequest_Panics", request: nil, def: bodyDef, panics: true},
|
||||
{name: "BodyExtractor_PointerRequest_Registers", request: new(body), def: bodyDef, panics: false},
|
||||
{name: "PathExtractor_ValueRequest_Registers", request: body{}, def: pathDef, panics: false},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
register := func() {
|
||||
New(func(http.ResponseWriter, *http.Request) {}, OpenAPIDef{ID: testCase.name, Request: testCase.request}, WithResourceDefs(testCase.def))
|
||||
}
|
||||
if testCase.panics {
|
||||
assert.Panics(t, register)
|
||||
} else {
|
||||
assert.NotPanics(t, register)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAttachStabilities(t *testing.T) {
|
||||
router := mux.NewRouter()
|
||||
router.Handle("/development", New(func(http.ResponseWriter, *http.Request) {}, OpenAPIDef{ID: "Development", SuccessStatusCode: http.StatusOK, Stability: StabilityDevelopment})).Methods(http.MethodGet)
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
package handler
|
||||
|
||||
import "github.com/SigNoz/signoz/pkg/http/binding"
|
||||
|
||||
type Option func(*handler)
|
||||
|
||||
func WithResourceDefs(defs ...ResourceDef) Option {
|
||||
@@ -9,9 +7,3 @@ func WithResourceDefs(defs ...ResourceDef) Option {
|
||||
h.resourceDefs = append(h.resourceDefs, defs...)
|
||||
}
|
||||
}
|
||||
|
||||
func WithBindBodyOptions(opts ...binding.BindBodyOption) Option {
|
||||
return func(h *handler) {
|
||||
h.bindBodyOptions = append(h.bindBodyOptions, opts...)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,7 +9,6 @@ type ResourceDef interface {
|
||||
// resolveRequest is unexported to seal the interface. It returns a slice so a
|
||||
// single def can fan out (e.g. a telemetry query touching multiple signals).
|
||||
resolveRequest(ec coretypes.ExtractorContext) []coretypes.ResolvedResource
|
||||
requiresBody() bool
|
||||
}
|
||||
|
||||
func ResolveRequest(defs []ResourceDef, ec coretypes.ExtractorContext) []coretypes.ResolvedResource {
|
||||
@@ -21,17 +20,6 @@ func ResolveRequest(defs []ResourceDef, ec coretypes.ExtractorContext) []coretyp
|
||||
return resolved
|
||||
}
|
||||
|
||||
// RequiresBody reports whether any def needs the decoded request body.
|
||||
func RequiresBody(defs []ResourceDef) bool {
|
||||
for _, def := range defs {
|
||||
if def.requiresBody() {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
// BasicResourceDef checks a single resource for one verb.
|
||||
type BasicResourceDef struct {
|
||||
Resource coretypes.Resource
|
||||
@@ -54,10 +42,6 @@ func (def BasicResourceDef) resolveRequest(ec coretypes.ExtractorContext) []core
|
||||
}
|
||||
}
|
||||
|
||||
func (def BasicResourceDef) requiresBody() bool {
|
||||
return def.ID.RequiresBody
|
||||
}
|
||||
|
||||
// AttachDetachSiblingResourceDef checks an attach/detach between peer resources;
|
||||
// both source and target are authz-checked.
|
||||
type AttachDetachSiblingResourceDef struct {
|
||||
@@ -88,10 +72,6 @@ func (def AttachDetachSiblingResourceDef) resolveRequest(ec coretypes.ExtractorC
|
||||
}
|
||||
}
|
||||
|
||||
func (def AttachDetachSiblingResourceDef) requiresBody() bool {
|
||||
return def.SourceIDs.RequiresBody || def.TargetIDs.RequiresBody
|
||||
}
|
||||
|
||||
// AttachDetachParentChildResourceDef authz-checks only the parent; the child
|
||||
// rides along for audit context.
|
||||
type AttachDetachParentChildResourceDef struct {
|
||||
@@ -121,20 +101,11 @@ func (def AttachDetachParentChildResourceDef) resolveRequest(ec coretypes.Extrac
|
||||
}
|
||||
}
|
||||
|
||||
func (def AttachDetachParentChildResourceDef) requiresBody() bool {
|
||||
return def.ParentID.RequiresBody || def.ChildIDs.RequiresBody
|
||||
}
|
||||
|
||||
type TelemetryResourceDef struct {
|
||||
Verb coretypes.Verb
|
||||
Category coretypes.ActionCategory
|
||||
Selector coretypes.SelectorFunc
|
||||
Resources coretypes.ResourceExtractor
|
||||
RequiresBody bool
|
||||
}
|
||||
|
||||
func (def TelemetryResourceDef) requiresBody() bool {
|
||||
return def.RequiresBody
|
||||
Verb coretypes.Verb
|
||||
Category coretypes.ActionCategory
|
||||
Selector coretypes.SelectorFunc
|
||||
Resources coretypes.ResourceExtractor
|
||||
}
|
||||
|
||||
func (def TelemetryResourceDef) resolveRequest(ec coretypes.ExtractorContext) []coretypes.ResolvedResource {
|
||||
|
||||
@@ -5,9 +5,7 @@ import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"reflect"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/handler"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/gorilla/mux"
|
||||
@@ -25,8 +23,8 @@ func NewResource(logger *slog.Logger) *Resource {
|
||||
|
||||
func (middleware *Resource) Wrap(next http.Handler) http.Handler {
|
||||
return http.HandlerFunc(func(rw http.ResponseWriter, req *http.Request) {
|
||||
provider := handlerFromRequest(req)
|
||||
if provider == nil || len(provider.ResourceDefs()) == 0 {
|
||||
defs := resourceDefsFromRequest(req)
|
||||
if len(defs) == 0 {
|
||||
next.ServeHTTP(rw, req)
|
||||
return
|
||||
}
|
||||
@@ -38,40 +36,18 @@ func (middleware *Resource) Wrap(next http.Handler) http.Handler {
|
||||
req.Body = io.NopCloser(bytes.NewReader(body))
|
||||
}
|
||||
|
||||
defs := provider.ResourceDefs()
|
||||
|
||||
var decoded any
|
||||
var decodeErr error
|
||||
if handler.RequiresBody(defs) {
|
||||
decoded, decodeErr = decodeBody(provider.Request(), body, provider.BindBodyOptions()...)
|
||||
extractorCtx := coretypes.ExtractorContext{
|
||||
Request: req,
|
||||
RequestBody: body,
|
||||
}
|
||||
resolved := handler.ResolveRequest(defs, extractorCtx)
|
||||
|
||||
extractorCtx := coretypes.ExtractorContext{Request: req, RequestBody: decoded}
|
||||
|
||||
var resolved []coretypes.ResolvedResource
|
||||
if decodeErr != nil {
|
||||
// authz renders the error inside the audit middleware, so the request is still logged
|
||||
resolved = []coretypes.ResolvedResource{coretypes.NewResolvedResourceWithError(coretypes.Verb{}, coretypes.ActionCategory{}, decodeErr)}
|
||||
} else {
|
||||
resolved = handler.ResolveRequest(defs, extractorCtx)
|
||||
}
|
||||
|
||||
ctx := coretypes.NewContextWithExtractorContext(req.Context(), extractorCtx)
|
||||
ctx = coretypes.NewContextWithResolvedResources(ctx, resolved)
|
||||
ctx := coretypes.NewContextWithResolvedResources(req.Context(), resolved)
|
||||
next.ServeHTTP(rw, req.WithContext(ctx))
|
||||
})
|
||||
}
|
||||
|
||||
func decodeBody(prototype any, body []byte, opts ...binding.BindBodyOption) (any, error) {
|
||||
decoded := reflect.New(reflect.TypeOf(prototype).Elem()).Interface()
|
||||
if err := binding.JSON.BindBody(bytes.NewReader(body), decoded, opts...); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return decoded, nil
|
||||
}
|
||||
|
||||
func handlerFromRequest(req *http.Request) handler.Handler {
|
||||
func resourceDefsFromRequest(req *http.Request) []handler.ResourceDef {
|
||||
route := mux.CurrentRoute(req)
|
||||
if route == nil {
|
||||
return nil
|
||||
@@ -87,5 +63,5 @@ func handlerFromRequest(req *http.Request) handler.Handler {
|
||||
return nil
|
||||
}
|
||||
|
||||
return provider
|
||||
return provider.ResourceDefs()
|
||||
}
|
||||
|
||||
@@ -5,11 +5,11 @@ import (
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/modules/authdomain"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
@@ -32,8 +32,8 @@ func (handler *handler) Create(rw http.ResponseWriter, req *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
body, err := coretypes.BodyFromContext[authtypes.PostableAuthDomain](req.Context())
|
||||
if err != nil {
|
||||
body := new(authtypes.PostableAuthDomain)
|
||||
if err := binding.JSON.BindBody(req.Body, body); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
@@ -142,8 +142,8 @@ func (handler *handler) Update(rw http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
body, err := coretypes.BodyFromContext[authtypes.UpdatableAuthDomain](r.Context())
|
||||
if err != nil {
|
||||
body := new(authtypes.UpdatableAuthDomain)
|
||||
if err := binding.JSON.BindBody(r.Body, body); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ func newConfig() factory.Config {
|
||||
Agent: AgentConfig{
|
||||
// we will maintain the latest version of cloud integration agent from here,
|
||||
// till we automate it externally or figure out a way to validate it.
|
||||
Version: "v0.0.14",
|
||||
Version: "v0.0.15",
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,7 +10,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/modules/cloudintegration"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/cloudintegrationtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
@@ -468,8 +467,8 @@ func (handler *handler) AgentCheckIn(rw http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
req, err := coretypes.BodyFromContext[cloudintegrationtypes.PostableAgentCheckIn](r.Context())
|
||||
if err != nil {
|
||||
req := new(cloudintegrationtypes.PostableAgentCheckIn)
|
||||
if err := binding.JSON.BindBody(r.Body, req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -134,6 +134,24 @@ func (store *store) UpdateAccount(ctx context.Context, account *cloudintegration
|
||||
BunDBCtx(ctx).
|
||||
NewUpdate().
|
||||
Model(account).
|
||||
Column("config").
|
||||
Column("updated_at").
|
||||
WherePK().
|
||||
Where("org_id = ?", account.OrgID).
|
||||
Where("provider = ?", account.Provider).
|
||||
Exec(ctx)
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
func (store *store) UpdateAgentReport(ctx context.Context, account *cloudintegrationtypes.StorableCloudIntegration) error {
|
||||
_, err := store.
|
||||
store.
|
||||
BunDBCtx(ctx).
|
||||
NewUpdate().
|
||||
Model(account).
|
||||
Column("account_id").
|
||||
Column("last_agent_report").
|
||||
WherePK().
|
||||
Where("org_id = ?", account.OrgID).
|
||||
Where("provider = ?", account.Provider).
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/modules/serviceaccount"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/serviceaccounttypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
@@ -223,8 +222,8 @@ func (handler *handler) CreateServiceAccountRole(rw http.ResponseWriter, r *http
|
||||
return
|
||||
}
|
||||
|
||||
req, err := coretypes.BodyFromContext[serviceaccounttypes.PostableServiceAccountRole](r.Context())
|
||||
if err != nil {
|
||||
req := new(serviceaccounttypes.PostableServiceAccountRole)
|
||||
if err := binding.JSON.BindBody(r.Body, req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -6,7 +6,9 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/modules/tracedetail"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
|
||||
@@ -18,6 +20,27 @@ func NewHandler(module tracedetail.Module) tracedetail.Handler {
|
||||
return &handler{module: module}
|
||||
}
|
||||
|
||||
func (h *handler) GetTraceSummary(rw http.ResponseWriter, r *http.Request) {
|
||||
claims, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
orgID, err := valuer.NewUUID(claims.OrgID)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
stats, err := h.module.GetTraceStats(r.Context(), orgID, mux.Vars(r)["traceID"])
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, spantypes.NewGettableTraceSummary(stats))
|
||||
}
|
||||
|
||||
func (h *handler) GetWaterfallV4(rw http.ResponseWriter, r *http.Request) {
|
||||
req := new(spantypes.PostableWaterfall)
|
||||
if err := binding.JSON.BindBody(r.Body, req); err != nil {
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/modules/tracedetail"
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"go.opentelemetry.io/otel/metric"
|
||||
)
|
||||
|
||||
@@ -39,6 +40,21 @@ func NewModule(traceStore spantypes.TraceStore, providerSettings factory.Provide
|
||||
return m
|
||||
}
|
||||
|
||||
func (m *module) GetTraceStats(ctx context.Context, orgID valuer.UUID, traceID string) (*spantypes.TraceStats, error) {
|
||||
summary, err := m.store.GetTraceSummary(ctx, traceID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
stats, err := m.store.GetTraceStats(ctx, orgID, traceID, summary)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if stats.TotalSpans == 0 {
|
||||
return nil, spantypes.ErrTraceNotFound
|
||||
}
|
||||
return stats, nil
|
||||
}
|
||||
|
||||
// GetWaterfallV4 is the OOM-safe V4 waterfall.
|
||||
// For large traces (NumSpans > effectiveLimit) it uses a two-step fetch:
|
||||
// minimal fields for all spans to build the tree, then full fields for the
|
||||
|
||||
@@ -10,9 +10,15 @@ import (
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/clickhousesql"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/flagger"
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/tracestelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
|
||||
const colServiceName = `resource_string_service$$$$name` // $ gets escaped so $$$$ converts to $$.
|
||||
@@ -38,10 +44,18 @@ type spanDurationRow struct {
|
||||
|
||||
type traceStore struct {
|
||||
telemetryStore telemetrystore.TelemetryStore
|
||||
metadataStore telemetrytypes.MetadataStore
|
||||
storage qbtypes.Storage
|
||||
flagger flagger.Flagger
|
||||
}
|
||||
|
||||
func NewTraceStore(ts telemetrystore.TelemetryStore) *traceStore {
|
||||
return &traceStore{telemetryStore: ts}
|
||||
func NewTraceStore(ts telemetrystore.TelemetryStore, metadataStore telemetrytypes.MetadataStore, fl flagger.Flagger) *traceStore {
|
||||
return &traceStore{
|
||||
telemetryStore: ts,
|
||||
metadataStore: metadataStore,
|
||||
storage: tracestelemetryschema.NewStorage(),
|
||||
flagger: fl,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *traceStore) GetTraceSummary(ctx context.Context, traceID string) (*spantypes.TraceSummary, error) {
|
||||
@@ -65,6 +79,131 @@ func (s *traceStore) GetTraceSummary(ctx context.Context, traceID string) (*span
|
||||
return &summary, nil
|
||||
}
|
||||
|
||||
func (s *traceStore) GetTraceStats(ctx context.Context, orgID valuer.UUID, traceID string, summary *spantypes.TraceSummary) (*spantypes.TraceStats, error) {
|
||||
table := fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable)
|
||||
spans := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
genAIColumns, err := s.genAISpanColumns(ctx, orgID, summary, spans)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// A span whose parent was never recorded hangs off a synthetic "Missing Span" root in the waterfall.
|
||||
ids := sqlbuilder.NewSelectBuilder()
|
||||
ids.Select("span_id")
|
||||
ids.From(table)
|
||||
ids.Where(
|
||||
ids.E("trace_id", traceID),
|
||||
ids.GE("ts_bucket_start", summary.Start.Unix()-1800),
|
||||
ids.LE("ts_bucket_start", summary.End.Unix()),
|
||||
)
|
||||
missingParent := fmt.Sprintf("parent_span_id <> '' AND parent_span_id GLOBAL NOT IN (%s)", spans.Var(ids))
|
||||
|
||||
spans.Select(
|
||||
"toUnixTimestamp64Nano(timestamp) AS span_start_ns",
|
||||
"span_start_ns + duration_nano AS span_end_ns",
|
||||
"span_id",
|
||||
"has_error",
|
||||
"("+missingParent+") AS has_missing_parent",
|
||||
"(parent_span_id = '' OR has_missing_parent) AS is_root",
|
||||
"if(parent_span_id = '', name, 'Missing Span') AS root_name",
|
||||
"if(parent_span_id = '', "+colServiceName+", '') AS root_service",
|
||||
)
|
||||
spans.SelectMore(genAIColumns...)
|
||||
spans.From(table)
|
||||
spans.Where(
|
||||
spans.E("trace_id", traceID),
|
||||
spans.GE("ts_bucket_start", summary.Start.Unix()-1800),
|
||||
spans.LE("ts_bucket_start", summary.End.Unix()),
|
||||
)
|
||||
spans.SQL("LIMIT 1 BY span_id")
|
||||
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
sb.Select(
|
||||
"toUInt64(min(span_start_ns)) AS start_ns",
|
||||
"toUInt64(max(span_end_ns)) AS end_ns",
|
||||
"count() AS total_spans",
|
||||
"countIf(has_error) AS total_error_spans",
|
||||
"countIf(has_missing_parent) > 0 AS has_missing_spans",
|
||||
"argMinIf(root_service, (span_start_ns, root_name), is_root) AS root_service_name",
|
||||
"argMinIf(root_name, (span_start_ns, root_name), is_root) AS root_entry_point",
|
||||
"countIf(is_gen_ai) AS gen_ai_span_count",
|
||||
"toUInt64(coalesce(sum(input_tokens_value), 0)) AS input_tokens",
|
||||
"toUInt64(coalesce(sum(output_tokens_value), 0)) AS output_tokens",
|
||||
"toUInt64(coalesce(sum(cache_read_tokens_value), 0)) AS cache_read_tokens",
|
||||
"toUInt64(coalesce(sum(cache_write_tokens_value), 0)) AS cache_write_tokens",
|
||||
"toUInt64(coalesce(sum(reasoning_tokens_value), 0)) AS reasoning_tokens",
|
||||
"sum(total_cost_value) AS total_cost",
|
||||
)
|
||||
sb.From(sb.BuilderAs(spans, "spans"))
|
||||
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
|
||||
var stats spantypes.TraceStats
|
||||
err = s.telemetryStore.ClickhouseDB().QueryRow(ctx, query, args...).Scan(
|
||||
&stats.StartNs, &stats.EndNs, &stats.TotalSpans, &stats.TotalErrorSpans, &stats.HasMissingSpans,
|
||||
&stats.RootServiceName, &stats.RootEntryPoint, &stats.GenAISpanCount,
|
||||
&stats.Tokens.Input, &stats.Tokens.Output, &stats.Tokens.CacheRead, &stats.Tokens.CacheWrite, &stats.Tokens.Reasoning,
|
||||
&stats.TotalCost,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error querying trace stats")
|
||||
}
|
||||
return &stats, nil
|
||||
}
|
||||
|
||||
// genAISpanColumns renders the per-span gen_ai gate and value reads through the shared
|
||||
// traces storage, so each attribute is read from the column its evolutions place it in
|
||||
// over the trace's own time window. Exists predicates bind their args into sb.
|
||||
func (s *traceStore) genAISpanColumns(ctx context.Context, orgID valuer.UUID, summary *spantypes.TraceSummary, sb *sqlbuilder.SelectBuilder) ([]string, error) {
|
||||
// no data type: metadata reports token counts as number, so a float64 request would
|
||||
// miss them and fall back to a map read without evolutions
|
||||
attributeKey := func(name string) *telemetrytypes.TelemetryFieldKey {
|
||||
return &telemetrytypes.TelemetryFieldKey{Name: name, Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute}
|
||||
}
|
||||
|
||||
selectors := make([]*telemetrytypes.FieldKeySelector, 0, len(aiobservabilitytypes.GenAISpanGateKeys)+len(spantypes.TraceStatsGenAIColumns))
|
||||
addSelector := func(name string) {
|
||||
selectors = append(selectors, &telemetrytypes.FieldKeySelector{
|
||||
Name: name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact,
|
||||
})
|
||||
}
|
||||
for _, name := range aiobservabilitytypes.GenAISpanGateKeys {
|
||||
addSelector(name)
|
||||
}
|
||||
for _, col := range spantypes.TraceStatsGenAIColumns {
|
||||
addSelector(col.Key)
|
||||
}
|
||||
keys, _, err := s.metadataStore.GetKeysMulti(ctx, orgID, querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, s.flagger, selectors))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
q := querybuilder.NewQueryInfo(ctx, orgID, s.flagger, telemetrytypes.SignalTraces, nil, uint64(summary.Start.UnixNano()), uint64(summary.End.UnixNano()))
|
||||
|
||||
gate := make([]string, 0, len(aiobservabilitytypes.GenAISpanGateKeys))
|
||||
for _, name := range aiobservabilitytypes.GenAISpanGateKeys {
|
||||
conds, _, err := querybuilder.Conditions(ctx, q, s.storage, attributeKey(name), qbtypes.FilterOperatorExists, nil, keys, false, sb)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
gate = append(gate, conds...)
|
||||
}
|
||||
columns := []string{sb.Or(gate...) + " AS is_gen_ai"}
|
||||
|
||||
for _, col := range spantypes.TraceStatsGenAIColumns {
|
||||
expr, err := querybuilder.ResolveColumn(ctx, q, s.storage, attributeKey(col.Key), telemetrytypes.FieldDataTypeFloat64, keys)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// a materialized column name carries `$$`, which Build would otherwise unescape
|
||||
columns = append(columns, sqlbuilder.Escape(expr)+" AS "+col.Column+"_value")
|
||||
}
|
||||
return columns, nil
|
||||
}
|
||||
|
||||
func (s *traceStore) GetTraceSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary) ([]spantypes.StorableSpan, error) {
|
||||
// DISTINCT ON (span_id) is ClickHouse-specific syntax not supported by sqlbuilder
|
||||
query := fmt.Sprintf(`
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -6,10 +6,12 @@ import (
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
|
||||
// Handler exposes HTTP handlers for trace detail APIs.
|
||||
type Handler interface {
|
||||
GetTraceSummary(http.ResponseWriter, *http.Request)
|
||||
GetWaterfallV4(http.ResponseWriter, *http.Request)
|
||||
GetTraceAggregations(http.ResponseWriter, *http.Request)
|
||||
GetFlamegraph(http.ResponseWriter, *http.Request)
|
||||
@@ -17,6 +19,7 @@ type Handler interface {
|
||||
|
||||
// Module defines the business logic for trace detail operations.
|
||||
type Module interface {
|
||||
GetTraceStats(ctx context.Context, orgID valuer.UUID, traceID string) (*spantypes.TraceStats, error)
|
||||
GetWaterfallV4(ctx context.Context, traceID string, selectedSpanID string, uncollapsedSpans []string) (*spantypes.GettableWaterfallTrace, error)
|
||||
GetTraceAggregations(ctx context.Context, traceID string, req *spantypes.PostableTraceAggregations) (*spantypes.GettableTraceAggregations, error)
|
||||
GetFlamegraph(ctx context.Context, traceID string, selectedSpanID string, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error)
|
||||
|
||||
@@ -14,7 +14,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
@@ -53,8 +52,8 @@ func (handler *handler) QueryRange(rw http.ResponseWriter, req *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
queryRangeRequest, err := coretypes.BodyFromContext[qbtypes.QueryRangeRequest](req.Context())
|
||||
if err != nil {
|
||||
var queryRangeRequest qbtypes.QueryRangeRequest
|
||||
if err := binding.JSON.BindBody(req.Body, &queryRangeRequest); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
@@ -71,7 +70,7 @@ func (handler *handler) QueryRange(rw http.ResponseWriter, req *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
queryRangeResponse, err := handler.querier.QueryRange(ctx, orgID, queryRangeRequest)
|
||||
queryRangeResponse, err := handler.querier.QueryRange(ctx, orgID, &queryRangeRequest)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
@@ -97,8 +96,8 @@ func (handler *handler) QueryRangePreview(rw http.ResponseWriter, req *http.Requ
|
||||
return
|
||||
}
|
||||
|
||||
queryRangeRequest, err := coretypes.BodyFromContext[qbtypes.QueryRangeRequest](req.Context())
|
||||
if err != nil {
|
||||
var queryRangeRequest qbtypes.QueryRangeRequest
|
||||
if err := json.NewDecoder(req.Body).Decode(&queryRangeRequest); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
@@ -119,7 +118,7 @@ func (handler *handler) QueryRangePreview(rw http.ResponseWriter, req *http.Requ
|
||||
return
|
||||
}
|
||||
|
||||
preview, err := handler.querier.QueryRangePreview(ctx, orgID, queryRangeRequest, previewOpts)
|
||||
preview, err := handler.querier.QueryRangePreview(ctx, orgID, &queryRangeRequest, previewOpts)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
|
||||
@@ -2,6 +2,7 @@ package querybuilder
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"strings"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
@@ -9,6 +10,7 @@ import (
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/tidwall/gjson"
|
||||
)
|
||||
|
||||
func TelemetrySelector(_ context.Context, resource coretypes.Resource, id string, _ valuer.UUID) ([]coretypes.Selector, error) {
|
||||
@@ -27,19 +29,20 @@ func TelemetrySelector(_ context.Context, resource coretypes.Resource, id string
|
||||
}
|
||||
|
||||
func QueryRangeResources(ec coretypes.ExtractorContext) ([]coretypes.ResourceWithID, error) {
|
||||
req, err := coretypes.BodyAs[qbtypes.QueryRangeRequest](ec)
|
||||
queries := gjson.GetBytes(ec.RequestBody, "compositeQuery.queries")
|
||||
if !queries.IsArray() || len(queries.Array()) == 0 {
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "atleast one query is required")
|
||||
}
|
||||
|
||||
variables, err := queryRangeVariables(ec.RequestBody)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if len(req.CompositeQuery.Queries) == 0 {
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "atleast one query is required")
|
||||
}
|
||||
|
||||
refs := make([]coretypes.ResourceWithID, 0, len(req.CompositeQuery.Queries))
|
||||
refs := make([]coretypes.ResourceWithID, 0, len(queries.Array()))
|
||||
seen := make(map[string]struct{})
|
||||
for _, query := range req.CompositeQuery.Queries {
|
||||
queryRefs, err := resourcesForQuery(query, req.Variables)
|
||||
for _, query := range queries.Array() {
|
||||
queryRefs, err := resourcesForQuery(query, variables)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -57,6 +60,21 @@ func QueryRangeResources(ec coretypes.ExtractorContext) ([]coretypes.ResourceWit
|
||||
return refs, nil
|
||||
}
|
||||
|
||||
func queryRangeVariables(body []byte) (map[string]qbtypes.VariableItem, error) {
|
||||
variables := make(map[string]qbtypes.VariableItem)
|
||||
|
||||
raw := gjson.GetBytes(body, "variables")
|
||||
if !raw.Exists() {
|
||||
return variables, nil
|
||||
}
|
||||
|
||||
if err := json.Unmarshal([]byte(raw.Raw), &variables); err != nil {
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid variables in query range request")
|
||||
}
|
||||
|
||||
return variables, nil
|
||||
}
|
||||
|
||||
// PromQLResources is the resource set of a bare PromQL query: metrics on
|
||||
// the promql wildcard, the same ID resourcesForQuery assigns to a PromQL
|
||||
// query inside a composite — one grant covers both entry points.
|
||||
@@ -67,53 +85,42 @@ func PromQLResources(coretypes.ExtractorContext) ([]coretypes.ResourceWithID, er
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func resourcesForQuery(query qbtypes.QueryEnvelope, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) {
|
||||
queryType := query.Type.StringValue()
|
||||
func resourcesForQuery(query gjson.Result, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) {
|
||||
queryType := query.Get("type").String()
|
||||
typeWildcard := queryType + "/" + coretypes.WildCardSelectorString
|
||||
|
||||
switch query.Type {
|
||||
case qbtypes.QueryTypeBuilder, qbtypes.QueryTypeSubQuery:
|
||||
return resourcesForBuilderQuery(queryType, query.Spec, variables)
|
||||
case qbtypes.QueryTypeBuilderAI:
|
||||
switch queryType {
|
||||
case qbtypes.QueryTypeBuilder.StringValue(), qbtypes.QueryTypeSubQuery.StringValue():
|
||||
return resourcesForBuilderQuery(queryType, query.Get("spec"), variables)
|
||||
case qbtypes.QueryTypeBuilderAI.StringValue():
|
||||
// always a traces query; the signal may be absent from the payload
|
||||
_, _, expression, err := builderQuerySpec(query.Spec)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return builderQueryResourceRefs(queryType, coretypes.ResourceTelemetryResourceTraces, expression, variables)
|
||||
case qbtypes.QueryTypePromQL:
|
||||
return builderQueryResourceRefs(queryType, coretypes.ResourceTelemetryResourceTraces, query.Get("spec"), variables)
|
||||
case qbtypes.QueryTypePromQL.StringValue():
|
||||
return []coretypes.ResourceWithID{{Resource: coretypes.ResourceTelemetryResourceMetrics, ID: typeWildcard}}, nil
|
||||
case qbtypes.QueryTypeClickHouseSQL:
|
||||
case qbtypes.QueryTypeClickHouseSQL.StringValue():
|
||||
return []coretypes.ResourceWithID{
|
||||
{Resource: coretypes.ResourceTelemetryResourceLogs, ID: typeWildcard},
|
||||
{Resource: coretypes.ResourceTelemetryResourceTraces, ID: typeWildcard},
|
||||
{Resource: coretypes.ResourceTelemetryResourceMetrics, ID: typeWildcard},
|
||||
{Resource: coretypes.ResourceTelemetryResourceMeterMetrics, ID: typeWildcard},
|
||||
}, nil
|
||||
case qbtypes.QueryTypeFormula, qbtypes.QueryTypeJoin, qbtypes.QueryTypeTraceOperator:
|
||||
case qbtypes.QueryTypeFormula.StringValue(), qbtypes.QueryTypeJoin.StringValue(), qbtypes.QueryTypeTraceOperator.StringValue():
|
||||
return nil, nil
|
||||
default:
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported query type %q", queryType)
|
||||
}
|
||||
}
|
||||
|
||||
func resourcesForBuilderQuery(queryType string, spec any, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) {
|
||||
signal, source, expression, err := builderQuerySpec(spec)
|
||||
func resourcesForBuilderQuery(queryType string, spec gjson.Result, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) {
|
||||
resource, err := builderQueryResource(spec)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
resource, err := builderQueryResource(signal, source)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return builderQueryResourceRefs(queryType, resource, expression, variables)
|
||||
return builderQueryResourceRefs(queryType, resource, spec, variables)
|
||||
}
|
||||
|
||||
func builderQueryResourceRefs(queryType string, resource coretypes.Resource, expression string, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) {
|
||||
ids, err := builderQuerySelectors(queryType, expression, variables)
|
||||
func builderQueryResourceRefs(queryType string, resource coretypes.Resource, spec gjson.Result, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) {
|
||||
ids, err := builderQuerySelectors(queryType, spec.Get("filter.expression").String(), variables)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -126,46 +133,27 @@ func builderQueryResourceRefs(queryType string, resource coretypes.Resource, exp
|
||||
return refs, nil
|
||||
}
|
||||
|
||||
func builderQueryResource(signal telemetrytypes.Signal, source telemetrytypes.Source) (coretypes.Resource, error) {
|
||||
switch signal {
|
||||
case telemetrytypes.SignalTraces:
|
||||
func builderQueryResource(spec gjson.Result) (coretypes.Resource, error) {
|
||||
source := spec.Get("source").String()
|
||||
|
||||
switch spec.Get("signal").String() {
|
||||
case telemetrytypes.SignalTraces.StringValue():
|
||||
return coretypes.ResourceTelemetryResourceTraces, nil
|
||||
case telemetrytypes.SignalLogs:
|
||||
if source == telemetrytypes.SourceAudit {
|
||||
case telemetrytypes.SignalLogs.StringValue():
|
||||
if source == telemetrytypes.SourceAudit.StringValue() {
|
||||
return coretypes.ResourceTelemetryResourceAuditLogs, nil
|
||||
}
|
||||
return coretypes.ResourceTelemetryResourceLogs, nil
|
||||
case telemetrytypes.SignalMetrics:
|
||||
if source == telemetrytypes.SourceMeter {
|
||||
case telemetrytypes.SignalMetrics.StringValue():
|
||||
if source == telemetrytypes.SourceMeter.StringValue() {
|
||||
return coretypes.ResourceTelemetryResourceMeterMetrics, nil
|
||||
}
|
||||
return coretypes.ResourceTelemetryResourceMetrics, nil
|
||||
default:
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported signal %q", signal.StringValue())
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported signal %q", spec.Get("signal").String())
|
||||
}
|
||||
}
|
||||
|
||||
func builderQuerySpec(spec any) (telemetrytypes.Signal, telemetrytypes.Source, string, error) {
|
||||
switch typed := spec.(type) {
|
||||
case qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]:
|
||||
return typed.Signal, typed.Source, filterExpression(typed.Filter), nil
|
||||
case qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]:
|
||||
return typed.Signal, typed.Source, filterExpression(typed.Filter), nil
|
||||
case qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]:
|
||||
return typed.Signal, typed.Source, filterExpression(typed.Filter), nil
|
||||
default:
|
||||
return telemetrytypes.Signal{}, telemetrytypes.Source{}, "", errors.Newf(errors.TypeInternal, errors.CodeInternal, "unexpected builder query spec %T", spec)
|
||||
}
|
||||
}
|
||||
|
||||
func filterExpression(filter *qbtypes.Filter) string {
|
||||
if filter == nil {
|
||||
return ""
|
||||
}
|
||||
|
||||
return filter.Expression
|
||||
}
|
||||
|
||||
func builderQuerySelectors(queryType, expression string, variables map[string]qbtypes.VariableItem) ([]string, error) {
|
||||
typeWildcard := queryType + "/" + coretypes.WildCardSelectorString
|
||||
|
||||
|
||||
@@ -2,24 +2,14 @@ package querybuilder
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/http/binding"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func queryRangeExtractorContext(t *testing.T, body string) coretypes.ExtractorContext {
|
||||
t.Helper()
|
||||
req := new(qbtypes.QueryRangeRequest)
|
||||
require.NoError(t, binding.JSON.BindBody(strings.NewReader(body), req))
|
||||
return coretypes.ExtractorContext{RequestBody: req}
|
||||
}
|
||||
|
||||
func builderQueryBody(signal, filterExpression string) string {
|
||||
return `{"compositeQuery":{"queries":[{"type":"builder_query","spec":{"signal":"` + signal + `","filter":{"expression":"` + filterExpression + `"}}}]}}`
|
||||
}
|
||||
@@ -152,13 +142,6 @@ func TestQueryRangeResources(t *testing.T) {
|
||||
{Resource: coretypes.ResourceTelemetryResourceLogs, ID: "builder_query/signoz.workspace.key.id/checkout"},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "DuplicateSignalKey_LastValueWins",
|
||||
body: `{"compositeQuery":{"queries":[{"type":"builder_query","spec":{"signal":"logs","signal":"traces","filter":{"expression":"signoz.workspace.key.id = 'a'"}}}]}}`,
|
||||
expected: []coretypes.ResourceWithID{
|
||||
{Resource: coretypes.ResourceTelemetryResourceTraces, ID: "builder_query/signoz.workspace.key.id/a"},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "duplicate queries dedupe",
|
||||
body: `{"compositeQuery":{"queries":[{"type":"builder_query","spec":{"signal":"logs","filter":{"expression":"signoz.workspace.key.id = 'a'"}}},{"type":"builder_query","spec":{"signal":"logs","filter":{"expression":"signoz.workspace.key.id='a'"}}}]}}`,
|
||||
@@ -170,7 +153,7 @@ func TestQueryRangeResources(t *testing.T) {
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
refs, err := QueryRangeResources(queryRangeExtractorContext(t, testCase.body))
|
||||
refs, err := QueryRangeResources(coretypes.ExtractorContext{RequestBody: []byte(testCase.body)})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.expected, refs)
|
||||
})
|
||||
@@ -182,20 +165,14 @@ func TestQueryRangeResourcesErrors(t *testing.T) {
|
||||
`{"compositeQuery":{"queries":[]}}`,
|
||||
`{}`,
|
||||
builderQueryBody("logs", "signoz.workspace.key.id = "),
|
||||
`{"compositeQuery":{"queries":[{"type":"builder_query","spec":{"signal":"unknown"}}]}}`,
|
||||
`{"compositeQuery":{"queries":[{"type":"unknown_type"}]}}`,
|
||||
}
|
||||
|
||||
for _, body := range bodies {
|
||||
_, err := QueryRangeResources(queryRangeExtractorContext(t, body))
|
||||
_, err := QueryRangeResources(coretypes.ExtractorContext{RequestBody: []byte(body)})
|
||||
assert.Error(t, err, "body %s", body)
|
||||
}
|
||||
|
||||
// rejected by the decode the middleware runs, before any extractor
|
||||
for _, body := range []string{
|
||||
`{"compositeQuery":{"queries":[{"type":"builder_query","spec":{"signal":"unknown"}}]}}`,
|
||||
`{"compositeQuery":{"queries":[{"type":"unknown_type"}]}}`,
|
||||
} {
|
||||
assert.Error(t, binding.JSON.BindBody(strings.NewReader(body), new(qbtypes.QueryRangeRequest)), "body %s", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelemetrySelector(t *testing.T) {
|
||||
|
||||
@@ -161,7 +161,7 @@ func NewModules(
|
||||
LogsPipeline: impllogspipeline.NewModule(sqlstore),
|
||||
RuleStateHistory: implrulestatehistory.NewModule(implrulestatehistory.NewStore(telemetryStore, telemetryMetadataStore, providerSettings.Logger), ruleStore),
|
||||
CloudIntegration: cloudIntegrationModule,
|
||||
TraceDetail: impltracedetail.NewModule(impltracedetail.NewTraceStore(telemetryStore), providerSettings, config.TraceDetail),
|
||||
TraceDetail: impltracedetail.NewModule(impltracedetail.NewTraceStore(telemetryStore, telemetryMetadataStore, fl), providerSettings, config.TraceDetail),
|
||||
SpanMapper: spanMapper,
|
||||
LLMPricingRule: impllmpricingrule.NewModule(impllmpricingrule.NewStore(sqlstore), querier),
|
||||
Tag: tagModule,
|
||||
|
||||
@@ -19,6 +19,7 @@ var (
|
||||
aiobservabilitytypes.GenAIUsageOutputTokens: genAIAttribute(aiobservabilitytypes.GenAIUsageOutputTokens, telemetrytypes.FieldDataTypeFloat64),
|
||||
aiobservabilitytypes.GenAIUsageCacheReadInputTokens: genAIAttribute(aiobservabilitytypes.GenAIUsageCacheReadInputTokens, telemetrytypes.FieldDataTypeFloat64),
|
||||
aiobservabilitytypes.GenAIUsageCacheCreationInputTokens: genAIAttribute(aiobservabilitytypes.GenAIUsageCacheCreationInputTokens, telemetrytypes.FieldDataTypeFloat64),
|
||||
aiobservabilitytypes.GenAIUsageReasoningOutputTokens: genAIAttribute(aiobservabilitytypes.GenAIUsageReasoningOutputTokens, telemetrytypes.FieldDataTypeFloat64),
|
||||
aiobservabilitytypes.SignozGenAITotalCost: genAIAttribute(aiobservabilitytypes.SignozGenAITotalCost, telemetrytypes.FieldDataTypeFloat64),
|
||||
|
||||
aiobservabilitytypes.GenAIInputMessages: genAIAttribute(aiobservabilitytypes.GenAIInputMessages, telemetrytypes.FieldDataTypeString),
|
||||
|
||||
@@ -15,6 +15,7 @@ const (
|
||||
GenAIUsageOutputTokens = "gen_ai.usage.output_tokens"
|
||||
GenAIUsageCacheReadInputTokens = "gen_ai.usage.cache_read.input_tokens"
|
||||
GenAIUsageCacheCreationInputTokens = "gen_ai.usage.cache_creation.input_tokens"
|
||||
GenAIUsageReasoningOutputTokens = "gen_ai.usage.reasoning.output_tokens"
|
||||
|
||||
GenAIInputMessages = "gen_ai.input.messages"
|
||||
GenAIOutputMessages = "gen_ai.output.messages"
|
||||
|
||||
@@ -3,6 +3,7 @@ package cloudintegrationtypes
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"maps"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
@@ -26,6 +27,17 @@ type Account struct {
|
||||
type AgentReport struct {
|
||||
TimestampMillis int64 `json:"timestampMillis" required:"true"`
|
||||
Data map[string]any `json:"data" required:"true" nullable:"true"`
|
||||
SyncState *SyncState `json:"syncState" required:"true" nullable:"true"`
|
||||
}
|
||||
|
||||
type SyncState struct {
|
||||
Version int64 `json:"version" required:"true"`
|
||||
InSync bool `json:"inSync" required:"true"`
|
||||
Regions map[string]*RegionSyncState `json:"regions" required:"true" nullable:"false"`
|
||||
}
|
||||
|
||||
type RegionSyncState struct {
|
||||
State RegionState `json:"state" required:"true"`
|
||||
}
|
||||
|
||||
type AccountConfig struct {
|
||||
@@ -150,6 +162,7 @@ func NewAccountFromStorable(storableAccount *StorableCloudIntegration) (*Account
|
||||
account.AgentReport = &AgentReport{
|
||||
TimestampMillis: storableAccount.LastAgentReport.TimestampMillis,
|
||||
Data: storableAccount.LastAgentReport.Data,
|
||||
SyncState: NewSyncStateFromStorable(storableAccount.LastAgentReport.SyncState),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -308,10 +321,28 @@ func NewAccountConfigFromUpdatable(provider CloudProviderType, config *Updatable
|
||||
}
|
||||
}
|
||||
|
||||
func NewAgentReport(data map[string]any) *AgentReport {
|
||||
func NewAgentReport(data map[string]any, syncState *SyncState) *AgentReport {
|
||||
return &AgentReport{
|
||||
TimestampMillis: time.Now().UnixMilli(),
|
||||
Data: data,
|
||||
SyncState: syncState,
|
||||
}
|
||||
}
|
||||
|
||||
func NewSyncStateFromStorable(storableSyncState *StorableSyncState) *SyncState {
|
||||
if storableSyncState == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
regions := make(map[string]*RegionSyncState, len(storableSyncState.Regions))
|
||||
for region, regionSyncState := range storableSyncState.Regions {
|
||||
regions[region] = &RegionSyncState{State: regionSyncState.State}
|
||||
}
|
||||
|
||||
return &SyncState{
|
||||
Version: storableSyncState.Version,
|
||||
InSync: storableSyncState.InSync,
|
||||
Regions: regions,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -335,6 +366,40 @@ func (account *Account) Update(provider CloudProviderType, config *AccountConfig
|
||||
return nil
|
||||
}
|
||||
|
||||
func (account *Account) UpdateAgentReport(providerAccountID *string, agentReport *AgentReport) {
|
||||
account.ProviderAccountID = providerAccountID
|
||||
account.AgentReport = agentReport
|
||||
}
|
||||
|
||||
// UpdateSyncState keeps the rest of the agent report, and is a no-op when the agent has never checked in.
|
||||
func (account *Account) UpdateSyncState(syncState *SyncState) {
|
||||
if account.AgentReport == nil {
|
||||
return
|
||||
}
|
||||
|
||||
account.AgentReport.SyncState = syncState
|
||||
}
|
||||
|
||||
// NextSyncState returns the sync state for this check-in, or nil for providers without one.
|
||||
func (account *Account) NextSyncState(syncedVersion *int64) *SyncState {
|
||||
if account.Provider != CloudProviderTypeAWS {
|
||||
return nil
|
||||
}
|
||||
|
||||
var previous *SyncState
|
||||
if account.AgentReport != nil {
|
||||
previous = account.AgentReport.SyncState
|
||||
}
|
||||
|
||||
regions := account.Config.AWS.Regions
|
||||
// Removed before the agent ever checked in: no region was sent to it, so there is nothing to clean up.
|
||||
if account.AgentReport == nil && account.RemovedAt != nil {
|
||||
regions = nil
|
||||
}
|
||||
|
||||
return newSyncState(previous, regions, account.RemovedAt != nil, syncedVersion)
|
||||
}
|
||||
|
||||
func (postableAccount *PostableAccount) UnmarshalJSON(data []byte) error {
|
||||
type Alias PostableAccount
|
||||
|
||||
@@ -406,3 +471,79 @@ func (config *AccountConfig) ToJSON() ([]byte, error) {
|
||||
func NewIngestionKeyName(provider CloudProviderType) string {
|
||||
return fmt.Sprintf("%s-integration", provider.StringValue())
|
||||
}
|
||||
|
||||
// newSyncState returns the sync state after a check-in without mutating previous.
|
||||
func newSyncState(previous *SyncState, regions []string, removed bool, syncedVersion *int64) *SyncState {
|
||||
if previous == nil {
|
||||
previous = newSyncStateFromRegions(regions)
|
||||
}
|
||||
|
||||
next := previous.copy()
|
||||
|
||||
// The agent synced this version, so its disabled regions are cleaned up and can be dropped.
|
||||
if syncedVersion != nil && *syncedVersion == next.Version {
|
||||
next.InSync = true
|
||||
maps.DeleteFunc(next.Regions, func(_ string, regionSyncState *RegionSyncState) bool {
|
||||
return regionSyncState.State == RegionStateDisabled
|
||||
})
|
||||
}
|
||||
|
||||
// Once the integration is removed, every region is disabled.
|
||||
if removed {
|
||||
regions = nil
|
||||
}
|
||||
|
||||
changed := false
|
||||
desiredRegionsMap := make(map[string]struct{}, len(regions))
|
||||
|
||||
for _, region := range regions {
|
||||
desiredRegionsMap[region] = struct{}{}
|
||||
|
||||
if regionSyncState, ok := next.Regions[region]; ok && regionSyncState.State == RegionStateEnabled {
|
||||
continue
|
||||
}
|
||||
|
||||
next.Regions[region] = &RegionSyncState{State: RegionStateEnabled}
|
||||
changed = true
|
||||
}
|
||||
|
||||
for region, regionSyncState := range next.Regions {
|
||||
_, ok := desiredRegionsMap[region]
|
||||
if ok && regionSyncState.State == RegionStateEnabled {
|
||||
continue
|
||||
}
|
||||
|
||||
if !ok && regionSyncState.State == RegionStateDisabled {
|
||||
continue
|
||||
}
|
||||
|
||||
regionSyncState.State = RegionStateDisabled
|
||||
changed = true
|
||||
}
|
||||
|
||||
if changed {
|
||||
next.Version++
|
||||
next.InSync = false
|
||||
}
|
||||
|
||||
return next
|
||||
}
|
||||
|
||||
// newSyncStateFromRegions is used on the first check-in, when the agent has already deployed regions, so it starts in sync.
|
||||
func newSyncStateFromRegions(regions []string) *SyncState {
|
||||
syncState := &SyncState{Version: 1, InSync: true, Regions: make(map[string]*RegionSyncState, len(regions))}
|
||||
for _, region := range regions {
|
||||
syncState.Regions[region] = &RegionSyncState{State: RegionStateEnabled}
|
||||
}
|
||||
|
||||
return syncState
|
||||
}
|
||||
|
||||
func (syncState *SyncState) copy() *SyncState {
|
||||
regions := make(map[string]*RegionSyncState, len(syncState.Regions))
|
||||
for region, regionSyncState := range syncState.Regions {
|
||||
regions[region] = &RegionSyncState{State: regionSyncState.State}
|
||||
}
|
||||
|
||||
return &SyncState{Version: syncState.Version, InSync: syncState.InSync, Regions: regions}
|
||||
}
|
||||
|
||||
@@ -12,7 +12,8 @@ type AgentCheckInRequest struct {
|
||||
ProviderAccountID string `json:"providerAccountId" required:"false"`
|
||||
CloudIntegrationID valuer.UUID `json:"cloudIntegrationId" required:"false"`
|
||||
|
||||
Data map[string]any `json:"data" required:"true" nullable:"true"`
|
||||
Data map[string]any `json:"data" required:"true" nullable:"true"`
|
||||
SyncedVersion *int64 `json:"syncedVersion" required:"false" nullable:"true"`
|
||||
}
|
||||
|
||||
type PostableAgentCheckIn struct {
|
||||
@@ -28,6 +29,7 @@ type AgentCheckInResponse struct {
|
||||
ProviderAccountID string `json:"providerAccountId" required:"true"`
|
||||
IntegrationConfig *ProviderIntegrationConfig `json:"integrationConfig" required:"true"`
|
||||
RemovedAt *time.Time `json:"removedAt" required:"true" nullable:"true"`
|
||||
SyncState *SyncState `json:"syncState" required:"true" nullable:"true"`
|
||||
}
|
||||
|
||||
type GettableAgentCheckIn struct {
|
||||
@@ -73,12 +75,13 @@ func NewGettableAgentCheckIn(provider CloudProviderType, resp *AgentCheckInRespo
|
||||
return gettable
|
||||
}
|
||||
|
||||
func NewAgentCheckInResponse(providerAccountID, cloudIntegrationID string, integrationConfig *ProviderIntegrationConfig, removedAt *time.Time) *AgentCheckInResponse {
|
||||
func NewAgentCheckInResponse(providerAccountID, cloudIntegrationID string, integrationConfig *ProviderIntegrationConfig, removedAt *time.Time, syncState *SyncState) *AgentCheckInResponse {
|
||||
return &AgentCheckInResponse{
|
||||
CloudIntegrationID: cloudIntegrationID,
|
||||
ProviderAccountID: providerAccountID,
|
||||
IntegrationConfig: integrationConfig,
|
||||
RemovedAt: removedAt,
|
||||
SyncState: syncState,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -25,6 +25,17 @@ var (
|
||||
ErrCodeServiceDefinitionNotFound = errors.MustNewCode("service_definition_not_found")
|
||||
)
|
||||
|
||||
var (
|
||||
RegionStateEnabled = RegionState{valuer.NewString("enabled")}
|
||||
RegionStateDisabled = RegionState{valuer.NewString("disabled")}
|
||||
)
|
||||
|
||||
type RegionState struct{ valuer.String }
|
||||
|
||||
func (RegionState) Enum() []any {
|
||||
return []any{RegionStateEnabled, RegionStateDisabled}
|
||||
}
|
||||
|
||||
// StorableCloudIntegration represents a cloud integration stored in the database.
|
||||
// This is also referred as "Account" in the context of cloud integrations.
|
||||
type StorableCloudIntegration struct {
|
||||
@@ -43,8 +54,16 @@ type StorableCloudIntegration struct {
|
||||
// StorableAgentReport represents the last heartbeat and arbitrary data sent by the agent
|
||||
// as of now there is no use case for Data field, but keeping it for backwards compatibility with older structure.
|
||||
type StorableAgentReport struct {
|
||||
TimestampMillis int64 `json:"timestamp_millis"` // backward compatibility
|
||||
Data map[string]any `json:"data"`
|
||||
TimestampMillis int64 `json:"timestamp_millis"` // backward compatibility
|
||||
Data map[string]any `json:"data"`
|
||||
SyncState *StorableSyncState `json:"sync_state,omitempty"`
|
||||
}
|
||||
|
||||
// StorableSyncState holds every region sent to the agent. A disabled region is dropped only after the agent acks Version.
|
||||
type StorableSyncState struct {
|
||||
Version int64 `json:"version"`
|
||||
InSync bool `json:"in_sync"`
|
||||
Regions map[string]*RegionSyncState `json:"regions"`
|
||||
}
|
||||
|
||||
// StorableCloudIntegrationService is to store service config for a cloud integration, which is a cloud provider specific configuration.
|
||||
@@ -148,12 +167,30 @@ func NewStorableCloudIntegration(account *Account) (*StorableCloudIntegration, e
|
||||
storableAccount.LastAgentReport = &StorableAgentReport{
|
||||
TimestampMillis: account.AgentReport.TimestampMillis,
|
||||
Data: account.AgentReport.Data,
|
||||
SyncState: NewStorableSyncState(account.AgentReport.SyncState),
|
||||
}
|
||||
}
|
||||
|
||||
return storableAccount, nil
|
||||
}
|
||||
|
||||
func NewStorableSyncState(syncState *SyncState) *StorableSyncState {
|
||||
if syncState == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
regions := make(map[string]*RegionSyncState, len(syncState.Regions))
|
||||
for region, regionSyncState := range syncState.Regions {
|
||||
regions[region] = &RegionSyncState{State: regionSyncState.State}
|
||||
}
|
||||
|
||||
return &StorableSyncState{
|
||||
Version: syncState.Version,
|
||||
InSync: syncState.InSync,
|
||||
Regions: regions,
|
||||
}
|
||||
}
|
||||
|
||||
// NewStorableCloudIntegrationService creates a new StorableCloudIntegrationService with
|
||||
// generated ID and timestamps from a CloudIntegrationService and its serialized config JSON.
|
||||
func NewStorableCloudIntegrationService(svc *CloudIntegrationService, configJSON string) *StorableCloudIntegrationService {
|
||||
@@ -172,6 +209,7 @@ func (account *StorableCloudIntegration) Update(providerAccountID *string, agent
|
||||
account.LastAgentReport = &StorableAgentReport{
|
||||
TimestampMillis: agentReport.TimestampMillis,
|
||||
Data: agentReport.Data,
|
||||
SyncState: NewStorableSyncState(agentReport.SyncState),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -25,9 +25,12 @@ type Store interface {
|
||||
// CreateAccount creates a new cloud integration account
|
||||
CreateAccount(ctx context.Context, account *StorableCloudIntegration) error
|
||||
|
||||
// UpdateAccount updates an existing cloud integration account
|
||||
// UpdateAccount updates the user updatable fields (config) of an existing cloud integration account
|
||||
UpdateAccount(ctx context.Context, account *StorableCloudIntegration) error
|
||||
|
||||
// UpdateAgentReport updates the provider account id and last agent report of an existing cloud integration account
|
||||
UpdateAgentReport(ctx context.Context, account *StorableCloudIntegration) error
|
||||
|
||||
// RemoveAccount marks a cloud integration account as removed by setting the RemovedAt field
|
||||
RemoveAccount(ctx context.Context, orgID, id valuer.UUID, provider CloudProviderType) error
|
||||
|
||||
|
||||
@@ -1,10 +1,8 @@
|
||||
package coretypes
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/tidwall/gjson"
|
||||
)
|
||||
@@ -14,70 +12,24 @@ const (
|
||||
PhaseResponse
|
||||
)
|
||||
|
||||
var (
|
||||
errCodeExtractorContextNotFound = errors.MustNewCode("extractor_context_not_found")
|
||||
errCodeRequestTypeUndeclared = errors.MustNewCode("request_type_undeclared")
|
||||
errCodeRequestTypeMismatch = errors.MustNewCode("request_type_mismatch")
|
||||
)
|
||||
|
||||
type ExtractPhase int
|
||||
|
||||
type extractorContextKey struct{}
|
||||
|
||||
// ExtractorContext carries everything an extractor may read: Request + RequestBody
|
||||
// are filled pre-handler, ResponseBody post-handler. RequestBody is the body
|
||||
// decoded by the resource middleware into the route's declared request type.
|
||||
// are filled pre-handler, ResponseBody post-handler.
|
||||
type ExtractorContext struct {
|
||||
Request *http.Request
|
||||
RequestBody any
|
||||
RequestBody []byte
|
||||
ResponseBody []byte
|
||||
}
|
||||
|
||||
func NewContextWithExtractorContext(ctx context.Context, ec ExtractorContext) context.Context {
|
||||
return context.WithValue(ctx, extractorContextKey{}, ec)
|
||||
}
|
||||
|
||||
func ExtractorContextFromContext(ctx context.Context) (ExtractorContext, error) {
|
||||
ec, ok := ctx.Value(extractorContextKey{}).(ExtractorContext)
|
||||
if !ok {
|
||||
return ExtractorContext{}, errors.New(errors.TypeInternal, errCodeExtractorContextNotFound, "extractor context not found in context")
|
||||
}
|
||||
|
||||
return ec, nil
|
||||
}
|
||||
|
||||
func BodyAs[T any](ec ExtractorContext) (*T, error) {
|
||||
if ec.RequestBody == nil {
|
||||
return nil, errors.New(errors.TypeInternal, errCodeRequestTypeUndeclared, "route does not declare a request type")
|
||||
}
|
||||
|
||||
typed, ok := ec.RequestBody.(*T)
|
||||
if !ok {
|
||||
return nil, errors.Newf(errors.TypeInternal, errCodeRequestTypeMismatch, "route declares request type %T, expected %T", ec.RequestBody, (*T)(nil))
|
||||
}
|
||||
|
||||
return typed, nil
|
||||
}
|
||||
|
||||
func BodyFromContext[T any](ctx context.Context) (*T, error) {
|
||||
ec, err := ExtractorContextFromContext(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return BodyAs[T](ec)
|
||||
}
|
||||
|
||||
type ResourceIDExtractor struct {
|
||||
Phase ExtractPhase
|
||||
RequiresBody bool
|
||||
Fn func(ExtractorContext) (string, error)
|
||||
Phase ExtractPhase
|
||||
Fn func(ExtractorContext) (string, error)
|
||||
}
|
||||
|
||||
type ResourceIDsExtractor struct {
|
||||
Phase ExtractPhase
|
||||
RequiresBody bool
|
||||
Fn func(ExtractorContext) ([]string, error)
|
||||
Phase ExtractPhase
|
||||
Fn func(ExtractorContext) ([]string, error)
|
||||
}
|
||||
|
||||
func NewResourceIDExtractor(phase ExtractPhase, fn func(ExtractorContext) (string, error)) ResourceIDExtractor {
|
||||
@@ -98,7 +50,7 @@ func OneID(extractor ResourceIDExtractor) ResourceIDsExtractor {
|
||||
return ResourceIDsExtractor{}
|
||||
}
|
||||
|
||||
return ResourceIDsExtractor{Phase: extractor.Phase, RequiresBody: extractor.RequiresBody, Fn: func(ec ExtractorContext) ([]string, error) {
|
||||
return ResourceIDsExtractor{Phase: extractor.Phase, Fn: func(ec ExtractorContext) ([]string, error) {
|
||||
id, err := extractor.Fn(ec)
|
||||
if err != nil || id == "" {
|
||||
return nil, err
|
||||
@@ -123,25 +75,26 @@ func PathParam(name string) ResourceIDExtractor {
|
||||
}}
|
||||
}
|
||||
|
||||
func BodyField[T any](pick func(*T) string) ResourceIDExtractor {
|
||||
return ResourceIDExtractor{Phase: PhaseRequest, RequiresBody: true, Fn: func(ec ExtractorContext) (string, error) {
|
||||
req, err := BodyAs[T](ec)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
return pick(req), nil
|
||||
func BodyJSONPath(path string) ResourceIDExtractor {
|
||||
return ResourceIDExtractor{Phase: PhaseRequest, Fn: func(ec ExtractorContext) (string, error) {
|
||||
return gjson.GetBytes(ec.RequestBody, path).String(), nil
|
||||
}}
|
||||
}
|
||||
|
||||
func BodyFields[T any](pick func(*T) []string) ResourceIDsExtractor {
|
||||
return ResourceIDsExtractor{Phase: PhaseRequest, RequiresBody: true, Fn: func(ec ExtractorContext) ([]string, error) {
|
||||
req, err := BodyAs[T](ec)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
func BodyJSONArray(path string) ResourceIDsExtractor {
|
||||
return ResourceIDsExtractor{Phase: PhaseRequest, Fn: func(ec ExtractorContext) ([]string, error) {
|
||||
result := gjson.GetBytes(ec.RequestBody, path)
|
||||
if !result.Exists() {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
return pick(req), nil
|
||||
array := result.Array()
|
||||
ids := make([]string, 0, len(array))
|
||||
for _, r := range array {
|
||||
ids = append(ids, r.String())
|
||||
}
|
||||
|
||||
return ids, nil
|
||||
}}
|
||||
}
|
||||
|
||||
|
||||
@@ -32,6 +32,7 @@ type SpanMapperStore interface {
|
||||
// TraceStore defines the data access interface for trace detail queries.
|
||||
type TraceStore interface {
|
||||
GetTraceSummary(ctx context.Context, traceID string) (*TraceSummary, error)
|
||||
GetTraceStats(ctx context.Context, orgID valuer.UUID, traceID string, summary *TraceSummary) (*TraceStats, error)
|
||||
GetTraceSpans(ctx context.Context, traceID string, summary *TraceSummary) ([]StorableSpan, error)
|
||||
GetMinimalSpans(ctx context.Context, traceID string, start, end time.Time) ([]MinimalSpan, error)
|
||||
GetTraceSpansByIDs(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
|
||||
|
||||
76
pkg/types/spantypes/trace_summary.go
Normal file
76
pkg/types/spantypes/trace_summary.go
Normal file
@@ -0,0 +1,76 @@
|
||||
package spantypes
|
||||
|
||||
import "github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
|
||||
|
||||
// TraceStatsGenAIColumns pairs each summed TraceStats column with the gen_ai attribute it sums.
|
||||
var TraceStatsGenAIColumns = []TraceStatsGenAIColumn{
|
||||
{Column: "input_tokens", Key: aiobservabilitytypes.GenAIUsageInputTokens},
|
||||
{Column: "output_tokens", Key: aiobservabilitytypes.GenAIUsageOutputTokens},
|
||||
{Column: "cache_read_tokens", Key: aiobservabilitytypes.GenAIUsageCacheReadInputTokens},
|
||||
{Column: "cache_write_tokens", Key: aiobservabilitytypes.GenAIUsageCacheCreationInputTokens},
|
||||
{Column: "reasoning_tokens", Key: aiobservabilitytypes.GenAIUsageReasoningOutputTokens},
|
||||
{Column: "total_cost", Key: aiobservabilitytypes.SignozGenAITotalCost},
|
||||
}
|
||||
|
||||
type TraceStatsGenAIColumn struct {
|
||||
Column string
|
||||
Key string
|
||||
}
|
||||
|
||||
// TraceStats is the single-row result of the trace summary aggregate query.
|
||||
type TraceStats struct {
|
||||
StartNs uint64
|
||||
EndNs uint64
|
||||
RootServiceName string
|
||||
RootEntryPoint string
|
||||
TotalSpans uint64
|
||||
TotalErrorSpans uint64
|
||||
HasMissingSpans bool
|
||||
GenAISpanCount uint64
|
||||
Tokens TraceAITokens
|
||||
TotalCost *float64
|
||||
}
|
||||
|
||||
// GettableTraceSummary is the response for the trace summary API; the trace-level
|
||||
// fields match the waterfall response.
|
||||
type GettableTraceSummary struct {
|
||||
StartTimestampMillis uint64 `json:"startTimestampMillis"`
|
||||
EndTimestampMillis uint64 `json:"endTimestampMillis"`
|
||||
RootServiceName string `json:"rootServiceName"`
|
||||
RootServiceEntryPoint string `json:"rootServiceEntryPoint"`
|
||||
TotalSpansCount uint64 `json:"totalSpansCount"`
|
||||
TotalErrorSpansCount uint64 `json:"totalErrorSpansCount"`
|
||||
HasMissingSpans bool `json:"hasMissingSpans"`
|
||||
AI *TraceAISummary `json:"ai,omitempty"`
|
||||
}
|
||||
|
||||
// TraceAISummary is present when any span carries a gen_ai gate key.
|
||||
type TraceAISummary struct {
|
||||
Tokens TraceAITokens `json:"tokens"`
|
||||
// TotalCost is null when no span carries a cost attribute.
|
||||
TotalCost *float64 `json:"totalCost" nullable:"true"`
|
||||
}
|
||||
|
||||
type TraceAITokens struct {
|
||||
Input uint64 `json:"input"`
|
||||
Output uint64 `json:"output"`
|
||||
CacheRead uint64 `json:"cacheRead"`
|
||||
CacheWrite uint64 `json:"cacheWrite"`
|
||||
Reasoning uint64 `json:"reasoning"`
|
||||
}
|
||||
|
||||
func NewGettableTraceSummary(stats *TraceStats) *GettableTraceSummary {
|
||||
summary := &GettableTraceSummary{
|
||||
StartTimestampMillis: stats.StartNs / 1_000_000,
|
||||
EndTimestampMillis: stats.EndNs / 1_000_000,
|
||||
RootServiceName: stats.RootServiceName,
|
||||
RootServiceEntryPoint: stats.RootEntryPoint,
|
||||
TotalSpansCount: stats.TotalSpans,
|
||||
TotalErrorSpansCount: stats.TotalErrorSpans,
|
||||
HasMissingSpans: stats.HasMissingSpans,
|
||||
}
|
||||
if stats.GenAISpanCount > 0 {
|
||||
summary.AI = &TraceAISummary{Tokens: stats.Tokens, TotalCost: stats.TotalCost}
|
||||
}
|
||||
return summary
|
||||
}
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/licensing"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/zeustypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
@@ -95,8 +94,8 @@ func (h *handler) PutHost(rw http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
req, err := coretypes.BodyFromContext[zeustypes.PostableHost](r.Context())
|
||||
if err != nil {
|
||||
req := new(zeustypes.PostableHost)
|
||||
if err := binding.JSON.BindBody(r.Body, req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
5
tests/fixtures/cloudintegrations.py
vendored
5
tests/fixtures/cloudintegrations.py
vendored
@@ -34,6 +34,8 @@ class ProviderAccountSpec:
|
||||
expected_config: Callable[[dict], dict]
|
||||
# only the suites that exercise updates need to supply it.
|
||||
updated_params: dict = field(default_factory=dict)
|
||||
# params -> the agentReport.syncState the API is expected to return after the first check-in.
|
||||
expected_sync_state: Callable[[dict], dict | None] = lambda p: None
|
||||
# id shown in parametrized test names; defaults to the provider slug.
|
||||
id: str = field(default="")
|
||||
|
||||
@@ -315,6 +317,7 @@ def simulate_agent_checkin(
|
||||
account_id: str,
|
||||
cloud_account_id: str,
|
||||
data: dict | None = None,
|
||||
synced_version: int | None = None,
|
||||
) -> requests.Response:
|
||||
endpoint = f"/api/v1/cloud_integrations/{cloud_provider}/accounts/check_in"
|
||||
|
||||
@@ -323,6 +326,8 @@ def simulate_agent_checkin(
|
||||
"providerAccountId": cloud_account_id,
|
||||
"data": data or {},
|
||||
}
|
||||
if synced_version is not None:
|
||||
checkin_payload["syncedVersion"] = synced_version
|
||||
|
||||
response = requests.post(
|
||||
signoz.self.host_configs["8080"].get(endpoint),
|
||||
|
||||
1
tests/fixtures/traces.py
vendored
1
tests/fixtures/traces.py
vendored
@@ -895,6 +895,7 @@ _TRACES_TABLES_TO_TRUNCATE = [
|
||||
"span_attributes_keys",
|
||||
"signoz_error_index_v2",
|
||||
"top_level_operations",
|
||||
"trace_summary",
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ from collections.abc import Callable
|
||||
from http import HTTPStatus
|
||||
|
||||
import pytest
|
||||
import requests
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD, add_license
|
||||
@@ -152,3 +153,230 @@ def test_duplicate_cloud_account_checkins(
|
||||
# Second check-in: account2 tries to claim the same provider account ID → 409
|
||||
response = simulate_agent_checkin(signoz, admin_token, spec.provider, account2["id"], same_provider_account_id)
|
||||
assert response.status_code == HTTPStatus.CONFLICT, f"Expected 409 for duplicate providerAccountId, got {response.status_code}: {response.text}"
|
||||
|
||||
|
||||
def test_sync_state_drops_removed_region_after_ack(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_cloud_integration_account: Callable,
|
||||
) -> None:
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
|
||||
provider_account_id = str(uuid.uuid4())
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
json={"config": {"aws": {"regions": ["us-east-1"]}}},
|
||||
timeout=10,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 2,
|
||||
"inSync": False,
|
||||
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}},
|
||||
}, "removed region should be marked disabled and the version bumped"
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=2)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 2,
|
||||
"inSync": True,
|
||||
"regions": {"us-east-1": {"state": "enabled"}},
|
||||
}, "acked removed region should be dropped"
|
||||
|
||||
|
||||
def test_sync_state_keeps_removed_region_without_ack(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_cloud_integration_account: Callable,
|
||||
) -> None:
|
||||
"""The agent failed to clean up or crashed, so it never acks: the removed region stays and the version stays put."""
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
|
||||
provider_account_id = str(uuid.uuid4())
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
json={"config": {"aws": {"regions": ["us-east-1"]}}},
|
||||
timeout=10,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
for _ in range(3):
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 2,
|
||||
"inSync": False,
|
||||
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}},
|
||||
}, "unacked removed region should stay without bumping the version"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("synced_version", [2, 9], ids=["stale", "ahead"])
|
||||
def test_sync_state_ignores_mismatched_ack(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_cloud_integration_account: Callable,
|
||||
synced_version: int,
|
||||
) -> None:
|
||||
"""An ack for any version other than the current one (v3) is ignored,
|
||||
so us-west-2, removed at v2 and still unacked, is not dropped.
|
||||
"""
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
|
||||
provider_account_id = str(uuid.uuid4())
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
for regions in (["us-east-1"], ["us-east-1", "eu-west-1"]):
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
json={"config": {"aws": {"regions": regions}}},
|
||||
timeout=10,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
expected_sync_state = {
|
||||
"version": 3,
|
||||
"inSync": False,
|
||||
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}, "eu-west-1": {"state": "enabled"}},
|
||||
}
|
||||
assert response.json()["data"]["syncState"] == expected_sync_state
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=synced_version)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == expected_sync_state, "an ack for another version should be ignored"
|
||||
|
||||
|
||||
def test_sync_state_applies_ack_before_config_change(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_cloud_integration_account: Callable,
|
||||
) -> None:
|
||||
"""The user changes regions while the agent syncs: the ack for the version it synced still lands."""
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
|
||||
provider_account_id = str(uuid.uuid4())
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
for regions, synced_version in ((["us-east-1"], None), (["us-east-1", "eu-west-1"], 2)):
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
json={"config": {"aws": {"regions": regions}}},
|
||||
timeout=10,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=synced_version)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 3,
|
||||
"inSync": False,
|
||||
"regions": {"us-east-1": {"state": "enabled"}, "eu-west-1": {"state": "enabled"}},
|
||||
}, "ack should drop the removed region before the new region bumps the version"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("synced_version", [1, None], ids=["agent_acks_synced_version", "agent_crashed"])
|
||||
def test_sync_state_region_removed_during_sync(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_cloud_integration_account: Callable,
|
||||
synced_version: int | None,
|
||||
) -> None:
|
||||
"""The user removes a region while the agent syncs v1; whether the agent acks v1 or crashed, the region must not be lost."""
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
|
||||
provider_account_id = str(uuid.uuid4())
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 1,
|
||||
"inSync": True,
|
||||
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "enabled"}},
|
||||
}
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
json={"config": {"aws": {"regions": ["us-east-1"]}}},
|
||||
timeout=10,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=synced_version)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 2,
|
||||
"inSync": False,
|
||||
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}},
|
||||
}, "region removed mid-sync should be marked disabled"
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=2)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 2,
|
||||
"inSync": True,
|
||||
"regions": {"us-east-1": {"state": "enabled"}},
|
||||
}
|
||||
|
||||
|
||||
def test_sync_state_after_disconnect(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_cloud_integration_account: Callable,
|
||||
) -> None:
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
|
||||
provider_account_id = str(uuid.uuid4())
|
||||
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
response = requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=10,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
for _ in range(2):
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["removedAt"] is not None, "removedAt should be set after disconnect"
|
||||
assert response.json()["data"]["syncState"] == {
|
||||
"version": 2,
|
||||
"inSync": False,
|
||||
"regions": {"us-east-1": {"state": "disabled"}, "us-west-2": {"state": "disabled"}},
|
||||
}, "every region should be disabled once, without bumping the version on later check-ins"
|
||||
|
||||
for _ in range(2):
|
||||
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=2)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"]["syncState"] == {"version": 2, "inSync": True, "regions": {}}, "acked removal should leave no regions"
|
||||
|
||||
@@ -21,6 +21,11 @@ AWS_ACCOUNT_SPEC = ProviderAccountSpec(
|
||||
updated_params={"deployment_region": "us-east-1", "regions": ["us-east-1", "us-west-2", "eu-west-1"]},
|
||||
build_config=lambda p: {"aws": {"deploymentRegion": p["deployment_region"], "regions": p["regions"]}},
|
||||
expected_config=lambda p: {"regions": p["regions"]},
|
||||
expected_sync_state=lambda p: {
|
||||
"version": 1,
|
||||
"inSync": True,
|
||||
"regions": {region: {"state": "enabled"} for region in p["regions"]},
|
||||
},
|
||||
)
|
||||
|
||||
GCP_ACCOUNT_SPEC = ProviderAccountSpec(
|
||||
@@ -128,6 +133,7 @@ def test_list_accounts_after_checkin(
|
||||
assert found["providerAccountId"] == provider_account_id, "providerAccountId should match"
|
||||
assert found["config"][spec.provider] == spec.expected_config(spec.initial_params), "config should match account config"
|
||||
assert found["agentReport"] is not None, "agentReport should be present after check-in"
|
||||
assert found["agentReport"]["syncState"] == spec.expected_sync_state(spec.initial_params), "syncState should be seeded from the account regions on first check-in"
|
||||
assert found["removedAt"] is None, "removedAt should be null for a live account"
|
||||
|
||||
|
||||
@@ -282,6 +288,7 @@ def test_update_account_after_checkin_preserves_connected_status(
|
||||
assert found_after is not None, "Account must still be listed after config update (account_id should not be reset)"
|
||||
assert found_after["providerAccountId"] == provider_account_id, "providerAccountId should be preserved after update"
|
||||
assert found_after["agentReport"] is not None, "agentReport should be preserved after update"
|
||||
assert found_after["agentReport"]["syncState"] == found_before["agentReport"]["syncState"], "config update must not change syncState"
|
||||
assert found_after["config"][spec.provider] == spec.expected_config(spec.updated_params), "Config should reflect the update"
|
||||
assert found_after["removedAt"] is None, "removedAt should still be null"
|
||||
|
||||
|
||||
@@ -2,8 +2,6 @@ from collections.abc import Callable
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
|
||||
import requests
|
||||
|
||||
from fixtures import querier, types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD, change_user_role, create_active_user
|
||||
from fixtures.querier import make_query_request
|
||||
@@ -152,27 +150,3 @@ def test_managed_viewer_meter_and_clickhouse_allowed_audit_denied(
|
||||
# audit-logs builder queries remain admin-only.
|
||||
audit = make_query_request(signoz, token, start, end, audit_query, request_type=querier.RequestType.RAW)
|
||||
assert audit.status_code == HTTPStatus.FORBIDDEN, audit.text
|
||||
|
||||
|
||||
def test_duplicate_signal_key_is_checked_on_the_bound_value(
|
||||
signoz: types.SigNoz,
|
||||
get_token: Callable[[str, str], str],
|
||||
) -> None:
|
||||
now = datetime.now(tz=UTC)
|
||||
start, end = int((now - timedelta(hours=1)).timestamp() * 1000), int(now.timestamp() * 1000)
|
||||
|
||||
# raw string: json= would collapse the duplicate "signal" key
|
||||
body = (
|
||||
f'{{"schemaVersion":"v1","start":{start},"end":{end},"requestType":"scalar",'
|
||||
'"compositeQuery":{"queries":[{"type":"builder_query","spec":{"name":"A","signal":"traces","signal":"logs",'
|
||||
'"disabled":false,"filter":{"expression":"signoz.workspace.key.id = \'key-a\'"},'
|
||||
'"aggregations":[{"expression":"count()"}]}}]},"noCache":true}'
|
||||
)
|
||||
|
||||
response = requests.post(
|
||||
signoz.self.host_configs["8080"].get("/api/v5/query_range"),
|
||||
timeout=querier.QUERY_TIMEOUT,
|
||||
headers={"authorization": f"Bearer {get_token(key_a_email, user_password)}", "content-type": "application/json"},
|
||||
data=body,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.FORBIDDEN, response.text
|
||||
|
||||
@@ -246,15 +246,6 @@ def test_attach_detach_dual_scoped(
|
||||
)
|
||||
assert resp.status_code == HTTPStatus.FORBIDDEN, f"assign viewer to target: expected 403, got {resp.status_code}: {resp.text}"
|
||||
|
||||
# duplicate roleId: the server binds the last one (viewer) -> forbidden. Raw string, json= would collapse the key.
|
||||
resp = requests.post(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/service_account_roles"),
|
||||
data=f'{{"serviceAccountId": "{target_id}", "roleId": "{editor_role_id}", "roleId": "{viewer_role_id}"}}',
|
||||
headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"},
|
||||
timeout=5,
|
||||
)
|
||||
assert resp.status_code == HTTPStatus.FORBIDDEN, f"assign duplicate roleId to target: expected 403, got {resp.status_code}: {resp.text}"
|
||||
|
||||
# Both SA-detach (target id) and role-detach (editor) present -> remove allowed.
|
||||
resp = requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/service_account_roles/{editor_entry_id}"),
|
||||
|
||||
280
tests/integration/tests/tracedetail/01_summary.py
Normal file
280
tests/integration/tests/tracedetail/01_summary.py
Normal file
@@ -0,0 +1,280 @@
|
||||
from collections.abc import Callable
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
|
||||
import pytest
|
||||
import requests
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.querierai import root_span
|
||||
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
|
||||
|
||||
WATERFALL_FIELDS = (
|
||||
"startTimestampMillis",
|
||||
"endTimestampMillis",
|
||||
"rootServiceName",
|
||||
"rootServiceEntryPoint",
|
||||
"totalSpansCount",
|
||||
"totalErrorSpansCount",
|
||||
"hasMissingSpans",
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("attribute_backend", ["map", "json"])
|
||||
def test_summary_ai_trace(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
use_attribute_backend: Callable[[str], None],
|
||||
attribute_backend: str,
|
||||
) -> None:
|
||||
"""The summary carries the waterfall's trace-level fields and, for a trace with gen_ai
|
||||
spans, token totals over every LLM span and the cost summed over the spans that carry it.
|
||||
Spans are written to one layout only, so a read from the wrong column sums to zero."""
|
||||
use_attribute_backend(attribute_backend)
|
||||
write_mode = "json_only" if attribute_backend == "json" else "legacy_only"
|
||||
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
|
||||
service = f"td-summary-{attribute_backend}"
|
||||
resources = {"service.name": service}
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
root_id = TraceIdGenerator.span_id()
|
||||
|
||||
insert_traces(
|
||||
[
|
||||
root_span(now=now, trace_id=trace_id, span_id=root_id, resources=resources, duration_s=4),
|
||||
Traces(
|
||||
timestamp=now - timedelta(seconds=4),
|
||||
duration=timedelta(seconds=1),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=root_id,
|
||||
name="chat gpt-4o-mini",
|
||||
kind=TracesKind.SPAN_KIND_CLIENT,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
attributes={
|
||||
"gen_ai.request.model": "gpt-4o-mini",
|
||||
"gen_ai.usage.input_tokens": 100,
|
||||
"gen_ai.usage.output_tokens": 20,
|
||||
"gen_ai.usage.cache_read.input_tokens": 7,
|
||||
"_signoz.gen_ai.total_cost": 0.01,
|
||||
},
|
||||
attribute_write_mode=write_mode,
|
||||
),
|
||||
# a failed LLM call: counted in tokens and errors, but priced by nobody
|
||||
Traces(
|
||||
timestamp=now - timedelta(seconds=3),
|
||||
duration=timedelta(seconds=0.5),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=root_id,
|
||||
name="chat gpt-4o-mini",
|
||||
kind=TracesKind.SPAN_KIND_CLIENT,
|
||||
status_code=TracesStatusCode.STATUS_CODE_ERROR,
|
||||
resources=resources,
|
||||
attributes={
|
||||
"gen_ai.request.model": "gpt-4o-mini",
|
||||
"gen_ai.usage.input_tokens": 50,
|
||||
"gen_ai.usage.output_tokens": 5,
|
||||
},
|
||||
attribute_write_mode=write_mode,
|
||||
),
|
||||
Traces(
|
||||
timestamp=now - timedelta(seconds=2),
|
||||
duration=timedelta(seconds=0.5),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=root_id,
|
||||
name="execute_tool",
|
||||
kind=TracesKind.SPAN_KIND_INTERNAL,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
attributes={"gen_ai.tool.name": "get_weather"},
|
||||
attribute_write_mode=write_mode,
|
||||
),
|
||||
]
|
||||
)
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
headers = {"authorization": f"Bearer {token}", "content-type": "application/json"}
|
||||
|
||||
summary = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/summary"), timeout=10, headers=headers)
|
||||
assert summary.status_code == HTTPStatus.OK, summary.text
|
||||
summary = summary.json()["data"]
|
||||
|
||||
waterfall = requests.post(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v4/traces/{trace_id}/waterfall"),
|
||||
timeout=10,
|
||||
headers=headers,
|
||||
json={"selectedSpanId": "", "uncollapsedSpans": []},
|
||||
)
|
||||
assert waterfall.status_code == HTTPStatus.OK, waterfall.text
|
||||
waterfall = waterfall.json()["data"]
|
||||
|
||||
assert {k: summary[k] for k in WATERFALL_FIELDS} == {k: waterfall[k] for k in WATERFALL_FIELDS}
|
||||
assert summary["rootServiceName"] == service
|
||||
assert summary["rootServiceEntryPoint"] == "POST /api/chat"
|
||||
assert summary["totalSpansCount"] == 4
|
||||
assert summary["totalErrorSpansCount"] == 1
|
||||
assert summary["hasMissingSpans"] is False
|
||||
|
||||
assert summary["ai"]["tokens"] == {"input": 150, "output": 25, "cacheRead": 7, "cacheWrite": 0, "reasoning": 0}
|
||||
assert summary["ai"]["totalCost"] == pytest.approx(0.01)
|
||||
|
||||
|
||||
def test_summary_ai_trace_across_json_rollout(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
seed_attribute_evolution: Callable[[str, datetime], None],
|
||||
) -> None:
|
||||
"""A trace that straddles the attribute JSON rollout has LLM spans written only to the legacy
|
||||
maps before it and to the JSON column after it. The summary window covers both, so the gen_ai
|
||||
reads must fall back across columns and sum every span."""
|
||||
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
|
||||
rollout = now - timedelta(minutes=30)
|
||||
seed_attribute_evolution("traces", rollout)
|
||||
|
||||
service = "td-summary-rollout"
|
||||
resources = {"service.name": service}
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
root_id = TraceIdGenerator.span_id()
|
||||
|
||||
insert_traces(
|
||||
[
|
||||
Traces(
|
||||
timestamp=rollout - timedelta(minutes=10),
|
||||
duration=timedelta(minutes=15),
|
||||
trace_id=trace_id,
|
||||
span_id=root_id,
|
||||
parent_span_id="",
|
||||
name="long agent run",
|
||||
kind=TracesKind.SPAN_KIND_SERVER,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
attribute_write_mode="legacy_only",
|
||||
),
|
||||
Traces(
|
||||
timestamp=rollout - timedelta(minutes=5),
|
||||
duration=timedelta(seconds=1),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=root_id,
|
||||
name="chat gpt-4o-mini",
|
||||
kind=TracesKind.SPAN_KIND_CLIENT,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
attributes={"gen_ai.request.model": "gpt-4o-mini", "gen_ai.usage.input_tokens": 100, "gen_ai.usage.output_tokens": 20, "_signoz.gen_ai.total_cost": 0.01},
|
||||
attribute_write_mode="legacy_only",
|
||||
),
|
||||
Traces(
|
||||
timestamp=rollout + timedelta(minutes=4),
|
||||
duration=timedelta(seconds=1),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=root_id,
|
||||
name="chat gpt-4o-mini",
|
||||
kind=TracesKind.SPAN_KIND_CLIENT,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
attributes={"gen_ai.request.model": "gpt-4o-mini", "gen_ai.usage.input_tokens": 50, "gen_ai.usage.output_tokens": 5, "_signoz.gen_ai.total_cost": 0.02},
|
||||
attribute_write_mode="json_only",
|
||||
),
|
||||
]
|
||||
)
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
summary = requests.get(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/summary"),
|
||||
timeout=10,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
)
|
||||
assert summary.status_code == HTTPStatus.OK, summary.text
|
||||
summary = summary.json()["data"]
|
||||
|
||||
assert summary["totalSpansCount"] == 3
|
||||
assert summary["rootServiceEntryPoint"] == "long agent run"
|
||||
assert summary["ai"]["tokens"] == {"input": 150, "output": 25, "cacheRead": 0, "cacheWrite": 0, "reasoning": 0}
|
||||
assert summary["ai"]["totalCost"] == pytest.approx(0.03)
|
||||
|
||||
|
||||
def test_summary_non_ai_trace_with_missing_root(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
) -> None:
|
||||
"""A trace whose recorded spans all hang off an unrecorded parent reports the synthetic
|
||||
"Missing Span" root exactly as the waterfall does, and a trace without gen_ai spans has
|
||||
no `ai` block."""
|
||||
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
|
||||
resources = {"service.name": "td-summary-orphan"}
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
missing_parent_id = TraceIdGenerator.span_id()
|
||||
|
||||
insert_traces(
|
||||
[
|
||||
Traces(
|
||||
timestamp=now - timedelta(seconds=5),
|
||||
duration=timedelta(seconds=2),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=missing_parent_id,
|
||||
name="SELECT users",
|
||||
kind=TracesKind.SPAN_KIND_CLIENT,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
),
|
||||
Traces(
|
||||
timestamp=now - timedelta(seconds=4),
|
||||
duration=timedelta(seconds=1),
|
||||
trace_id=trace_id,
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
parent_span_id=missing_parent_id,
|
||||
name="publish event",
|
||||
kind=TracesKind.SPAN_KIND_PRODUCER,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources=resources,
|
||||
),
|
||||
]
|
||||
)
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
headers = {"authorization": f"Bearer {token}", "content-type": "application/json"}
|
||||
|
||||
summary = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/summary"), timeout=10, headers=headers)
|
||||
assert summary.status_code == HTTPStatus.OK, summary.text
|
||||
summary = summary.json()["data"]
|
||||
|
||||
waterfall = requests.post(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v4/traces/{trace_id}/waterfall"),
|
||||
timeout=10,
|
||||
headers=headers,
|
||||
json={"selectedSpanId": "", "uncollapsedSpans": []},
|
||||
)
|
||||
assert waterfall.status_code == HTTPStatus.OK, waterfall.text
|
||||
waterfall = waterfall.json()["data"]
|
||||
|
||||
assert {k: summary[k] for k in WATERFALL_FIELDS} == {k: waterfall[k] for k in WATERFALL_FIELDS}
|
||||
assert summary["hasMissingSpans"] is True
|
||||
assert summary["rootServiceName"] == ""
|
||||
assert summary["rootServiceEntryPoint"] == "Missing Span"
|
||||
assert summary["totalSpansCount"] == 2
|
||||
assert "ai" not in summary
|
||||
|
||||
|
||||
def test_summary_unknown_trace(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
) -> None:
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v1/traces/{TraceIdGenerator.trace_id()}/summary"),
|
||||
timeout=10,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NOT_FOUND, response.text
|
||||
Reference in New Issue
Block a user