Compare commits

...

5 Commits

Author SHA1 Message Date
nityanandagohain
86878658ca Merge remote-tracking branch 'origin/main' into issue_6090 2026-09-30 22:34:50 +05:30
Vikrant Gupta
e374d03e54 revert(authz): restore gjson-based body extraction in resource middleware (#13024)
Some checks are pending
build-staging / go-build (push) Blocked by required conditions
cacheci / tests (push) Waiting to run
build-staging / staging (push) Blocked by required conditions
build-staging / prepare (push) Waiting to run
build-staging / js-build (push) Blocked by required conditions
Release Drafter / update_release_draft (push) Waiting to run
#### Description

- Reverts #13014 and #13015. The resource middleware goes back to
reading body-derived resource ids with `BodyJSONPath` / `BodyJSONArray`
over the raw body, and handlers decode their own request bodies again.
- Authz should not own request decoding; that ownership stays with the
handlers.

#### Additional Information

- Contributes to: https://github.com/SigNoz/keystone-pod/issues/37
2026-09-30 13:53:21 +00:00
Swapnil Nakade
d445b6c296 chore: bumping cloud integration agent version to v0.0.15 (#13023)
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description
Bumping the cloud integration agent's version to latest v0.0.15

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR
Contributes to https://github.com/SigNoz/keystone-pod/issues/101
2026-09-30 12:38:24 +00:00
Swapnil Nakade
fd8aaac300 feat: adding sync state in cloud integration (#12991)
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description
The agent only saw the current list of enabled regions, so it couldn’t
tell which regions had been removed. To find stacks to clean up, it
checked unrelated AWS regions, causing unnecessary calls and permission
errors. Sync state keeps track of regions sent to the agent and pending
removals until the agent acknowledges cleanup.

Please check
[comment](https://github.com/SigNoz/keystone-pod/issues/101#issuecomment-5832865898)
for approach

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR
Contributes to https://github.com/SigNoz/keystone-pod/issues/101

<!--Anything reviewers should keep in mind while reviewing -->
#### Additional Information
This PR should be merged before changes for cloud-integration repo.

<!--Please delete paragraphs that you did not use before submitting.-->
2026-09-30 11:39:46 +00:00
nityanandagohain
5dae0b975a feat: add trace summary endpoint 2026-09-18 17:55:43 +05:30
54 changed files with 1651 additions and 433 deletions

View File

@@ -68,6 +68,7 @@ jobs:
- semconvfamilies
- serviceaccount
- spanmapper
- tracedetail
- querier_json_body
- querier_skip_resource_fingerprint
- ttl

View File

@@ -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

View File

@@ -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).

View File

@@ -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
}

View File

@@ -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

View File

@@ -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

View File

@@ -24,6 +24,7 @@ const accountsResponse: ListAccounts200 = {
agentReport: {
timestampMillis: 1747114366214,
data: null,
syncState: null,
},
providerAccountId: PROVIDER_ACCOUNT_ID,
removedAt: null,

View File

@@ -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,

View File

@@ -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

View File

@@ -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 {

View File

@@ -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")),

View File

@@ -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),

View File

@@ -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
}

View File

@@ -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),

View File

@@ -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 {

View File

@@ -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{

View File

@@ -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

View File

@@ -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
}

View File

@@ -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
}

View File

@@ -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)

View File

@@ -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...)
}
}

View File

@@ -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 {

View File

@@ -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()
}

View File

@@ -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
}

View File

@@ -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",
},
}
}

View File

@@ -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
}

View File

@@ -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).

View File

@@ -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
}

View File

@@ -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 {

View File

@@ -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

View File

@@ -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

View File

@@ -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)

View File

@@ -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

View File

@@ -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

View File

@@ -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) {

View File

@@ -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,

View File

@@ -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),

View File

@@ -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"

View File

@@ -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}
}

View File

@@ -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,
}
}

View File

@@ -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),
}
}
}

View File

@@ -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

View File

@@ -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
}}
}

View File

@@ -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)

View 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
}

View File

@@ -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
}

View File

@@ -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),

View File

@@ -895,6 +895,7 @@ _TRACES_TABLES_TO_TRUNCATE = [
"span_attributes_keys",
"signoz_error_index_v2",
"top_level_operations",
"trace_summary",
]

View File

@@ -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"

View File

@@ -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"

View File

@@ -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

View File

@@ -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}"),

View 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