Compare commits

...

20 Commits

Author SHA1 Message Date
vikrantgupta25
2a3aef7b48 Merge remote-tracking branch 'origin/main' into keystone-pod/issues/38
# Conflicts:
#	pkg/http/handler/resourcedef.go
2026-09-30 19:47:03 +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
vikrantgupta25
d08cd99107 refactor(user): drop self-mutation guards from admin routes 2026-09-30 18:43:18 +05:30
vikrantgupta25
507ea382cf refactor(authz): read user route body ids with the gjson extractors 2026-09-30 18:43:16 +05:30
vikrantgupta25
58acd6d337 Merge remote-tracking branch 'origin/revert/authz-body-decoding' into keystone-pod/issues/38 2026-09-30 18:41:22 +05:30
vikrantgupta25
e4afc8832b revert(authz): decode the request body once in the resource middleware
This reverts commit 8d80f98710 (#13014).
Authz should not own request body decoding.
2026-09-30 18:31:04 +05:30
vikrantgupta25
585a3b6140 revert(authz): read every body-derived resource id from the decoded request
This reverts commit 8dba9a13ea (#13015).
Authz should not own request body decoding.
2026-09-30 18:31:03 +05:30
vikrantgupta25
6a1fac8118 chore(openapi): regenerate spec for reset password token read 2026-09-30 18:21:05 +05:30
vikrantgupta25
b8869e7005 fix(authz): check only factor-password list when reading a reset password token 2026-09-30 18:20:29 +05:30
vikrantgupta25
2a5647f6b8 Merge remote-tracking branch 'origin/main' into keystone-pod/issues/38
# Conflicts:
#	pkg/signoz/provider.go
2026-09-30 18:20:01 +05:30
vikrantgupta25
ed45d5054d chore(openapi): regenerate spec for users by role scope 2026-09-30 18:17:41 +05:30
vikrantgupta25
057b61fe68 fix(authz): check only role read when listing users by role 2026-09-30 18:17:39 +05:30
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
vikrantgupta25
643de2055d chore(openapi): regenerate spec for reset password token scope 2026-09-30 17:23:54 +05:30
vikrantgupta25
86d9dae524 fix(authz): check factor-password list when reading a reset password token 2026-09-30 17:23:52 +05:30
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
vikrantgupta25
7f3ec73dfd fix(authz): guard nil request in user role extractors 2026-09-30 17:06:58 +05:30
Naman Verma
e4cd8dbf10 chore: store v2 config for notification channels in db (#12984)
<!--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

1. add a config column to notification channels table where the channel
config goes without dealing with receiver at all. this helps in cleaning
up all round trip issues caused by dealing with receiver in the storage
layer. v2 apis treat receiver as a side effect now
1a. for applicable fields, defaults are filled in create/update api if
fields are omitted
1b. explicit [] and {} are no longer dropped, and omitted [] and {} are
returned as explicit null
2. add integration tests for all notification channel round trip issues
3. code cleanup of v2 channels types 
4. migration to fill the config column from receiver column, which logs
results like dashboards migration did
4a. it fails for receivers that cannot be modeled in v2. repair api can
be used for them
4b. for types that v2 supports, any fields that v1 supports but v2
doesn't, this migration drops those fields
5. make v1 API reject anything that v2 apis do not support, and also
fill the new config column added
6. add integration tests for v1<>v2 interaction to ensure that alert
manager doesn't break because of new changes added

What breaks/changes for v1:
1. types that v2 does not support can no longer be created
2. channels with multiple receivers cannot be created
3. channels with fields that v2 does not support cannot be created
4. existing channels of types that v2 does not support can no longer be
edited via v1 API. They can be deleted though. Also, they keep on
sending notifications as before (sigh).

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

Closes https://github.com/SigNoz/pulse-pod/issues/376
Closes https://github.com/SigNoz/pulse-pod/issues/378
2026-09-30 11:30:33 +00:00
vikrantgupta25
a39eec0f5c chore(openapi): regenerate spec for user routes 2026-09-30 16:43:17 +05:30
vikrantgupta25
fae6fc41e3 feat(authz): enable FGA for users and reset password tokens 2026-09-30 16:43:15 +05:30
81 changed files with 4928 additions and 1925 deletions

View File

@@ -184,6 +184,7 @@ components:
headers:
additionalProperties:
type: string
nullable: true
type: object
html:
type: string
@@ -217,6 +218,7 @@ components:
metadata:
additionalProperties:
type: string
nullable: true
type: object
sendResolved:
nullable: true
@@ -258,6 +260,7 @@ components:
type: string
customFields:
additionalProperties: {}
nullable: true
type: object
description:
type: string
@@ -268,6 +271,7 @@ components:
labels:
items:
type: string
nullable: true
type: array
priority:
type: string
@@ -346,6 +350,7 @@ components:
details:
additionalProperties:
type: string
nullable: true
type: object
message:
type: string
@@ -374,6 +379,7 @@ components:
details:
additionalProperties:
type: string
nullable: true
type: object
group:
type: string
@@ -451,6 +457,7 @@ components:
actions:
items:
$ref: '#/components/schemas/AlertmanagertypesChannelSlackAction'
nullable: true
type: array
apiUrl:
format: password
@@ -464,6 +471,7 @@ components:
fields:
items:
$ref: '#/components/schemas/AlertmanagertypesChannelSlackField'
nullable: true
type: array
footer:
type: string
@@ -1763,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:
@@ -2013,6 +2024,8 @@ components:
format: date-time
nullable: true
type: string
syncState:
$ref: '#/components/schemas/CloudintegrationtypesSyncState'
required:
- account_id
- cloud_account_id
@@ -2022,6 +2035,7 @@ components:
- providerAccountId
- integrationConfig
- removedAt
- syncState
type: object
CloudintegrationtypesGettableServicesMetadata:
properties:
@@ -2121,6 +2135,9 @@ components:
type: object
providerAccountId:
type: string
syncedVersion:
nullable: true
type: integer
required:
- data
type: object
@@ -2133,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:
@@ -2274,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:
@@ -21140,9 +21186,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- role:read
- tokenizer:
- ADMIN
- role:read
summary: Get users by role id
tags:
- users
@@ -25728,9 +25774,11 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:attach
- role:attach
- tokenizer:
- ADMIN
- user:attach
- role:attach
summary: Create user role
tags:
- users
@@ -25781,9 +25829,11 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:detach
- role:detach
- tokenizer:
- ADMIN
- user:detach
- role:detach
summary: Delete user role
tags:
- users
@@ -25845,9 +25895,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:read
- tokenizer:
- ADMIN
- user:read
summary: Get user role
tags:
- users
@@ -25894,9 +25944,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:list
- tokenizer:
- ADMIN
- user:list
summary: List users v2
tags:
- users
@@ -25957,9 +26007,13 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:create
- user:attach
- role:attach
- tokenizer:
- ADMIN
- user:create
- user:attach
- role:attach
summary: Create user
tags:
- users
@@ -26004,9 +26058,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:delete
- tokenizer:
- ADMIN
- user:delete
summary: Delete user
tags:
- users
@@ -26062,9 +26116,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:read
- tokenizer:
- ADMIN
- user:read
summary: Get user by user id
tags:
- users
@@ -26119,9 +26173,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:update
- tokenizer:
- ADMIN
- user:update
summary: Update user v2
tags:
- users
@@ -26178,9 +26232,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- factor-password:list
- tokenizer:
- ADMIN
- factor-password:list
summary: Get reset password token for a user
tags:
- users
@@ -26244,9 +26298,11 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- factor-password:create
- user:attach
- tokenizer:
- ADMIN
- factor-password:create
- user:attach
summary: Create or regenerate reset password token for a user
tags:
- users
@@ -26305,9 +26361,9 @@ paths:
description: Internal Server Error
security:
- api_key:
- ADMIN
- user:read
- tokenizer:
- ADMIN
- user:read
summary: Get user roles
tags:
- users
@@ -26621,10 +26677,7 @@ paths:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- ADMIN
- tokenizer:
- ADMIN
- tokenizer: []
summary: Updates my password
tags:
- users

View File

@@ -118,10 +118,10 @@ 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).
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). When a side's ids come from a list extractor (`BodyJSONArray` or a custom `ResourceIDsExtractor`) and that list resolves to nothing at request time, there is nothing to link: the def resolves to no resources, so no check runs and no audit event is emitted (e.g. inviting a user with an empty `userRoles`). A single-id side (`OneID`) always resolves to exactly one id and fails closed when it is empty. 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).
Prefer `CheckResources` with a `ResourceDef` for anything resource-shaped. The older coarse gates `ViewAccess`/`EditAccess`/`AdminAccess` only check "does the caller hold one of these roles" and give up per-resource granularity; `OpenAccess` performs no authorization (authentication still applies); `CheckWithoutClaims` serves anonymous routes such as public dashboards.

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

@@ -104,9 +104,9 @@ export interface AlertmanagertypesChannelSlackFieldDTO {
export interface AlertmanagertypesChannelSlackConfigDTO {
/**
* @type array
* @type array,null
*/
actions?: AlertmanagertypesChannelSlackActionDTO[];
actions?: AlertmanagertypesChannelSlackActionDTO[] | null;
/**
* @type string
* @format password
@@ -125,9 +125,9 @@ export interface AlertmanagertypesChannelSlackConfigDTO {
*/
fallback?: string;
/**
* @type array
* @type array,null
*/
fields?: AlertmanagertypesChannelSlackFieldDTO[];
fields?: AlertmanagertypesChannelSlackFieldDTO[] | null;
/**
* @type string
*/
@@ -166,13 +166,19 @@ export interface AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTy
export enum AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTypesAlertmanagertypesChannelEmailConfigDTOKind {
email = 'email',
}
export type AlertmanagertypesChannelEmailConfigDTOHeaders = {
export type AlertmanagertypesChannelEmailConfigDTOHeadersAnyOf = {
[key: string]: string;
};
/**
* @nullable
*/
export type AlertmanagertypesChannelEmailConfigDTOHeaders =
AlertmanagertypesChannelEmailConfigDTOHeadersAnyOf | null;
export interface AlertmanagertypesChannelEmailConfigDTO {
/**
* @type object
* @type object,null
*/
headers?: AlertmanagertypesChannelEmailConfigDTOHeaders;
/**
@@ -239,10 +245,16 @@ export interface AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTy
export enum AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTypesAlertmanagertypesChannelPagerdutyConfigDTOKind {
pagerduty = 'pagerduty',
}
export type AlertmanagertypesChannelPagerdutyConfigDTODetails = {
export type AlertmanagertypesChannelPagerdutyConfigDTODetailsAnyOf = {
[key: string]: string;
};
/**
* @nullable
*/
export type AlertmanagertypesChannelPagerdutyConfigDTODetails =
AlertmanagertypesChannelPagerdutyConfigDTODetailsAnyOf | null;
export interface AlertmanagertypesChannelPagerdutyConfigDTO {
/**
* @type string
@@ -265,7 +277,7 @@ export interface AlertmanagertypesChannelPagerdutyConfigDTO {
*/
description?: string;
/**
* @type object
* @type object,null
*/
details?: AlertmanagertypesChannelPagerdutyConfigDTODetails;
/**
@@ -307,10 +319,16 @@ export interface AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTy
export enum AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTypesAlertmanagertypesChannelOpsgenieConfigDTOKind {
opsgenie = 'opsgenie',
}
export type AlertmanagertypesChannelOpsgenieConfigDTODetails = {
export type AlertmanagertypesChannelOpsgenieConfigDTODetailsAnyOf = {
[key: string]: string;
};
/**
* @nullable
*/
export type AlertmanagertypesChannelOpsgenieConfigDTODetails =
AlertmanagertypesChannelOpsgenieConfigDTODetailsAnyOf | null;
export interface AlertmanagertypesChannelOpsgenieConfigDTO {
/**
* @type string
@@ -326,7 +344,7 @@ export interface AlertmanagertypesChannelOpsgenieConfigDTO {
*/
description?: string;
/**
* @type object
* @type object,null
*/
details?: AlertmanagertypesChannelOpsgenieConfigDTODetails;
/**
@@ -423,10 +441,16 @@ export interface AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTy
export enum AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTypesAlertmanagertypesChannelJiraConfigDTOKind {
jira = 'jira',
}
export type AlertmanagertypesChannelJiraConfigDTOCustomFields = {
export type AlertmanagertypesChannelJiraConfigDTOCustomFieldsAnyOf = {
[key: string]: unknown;
};
/**
* @nullable
*/
export type AlertmanagertypesChannelJiraConfigDTOCustomFields =
AlertmanagertypesChannelJiraConfigDTOCustomFieldsAnyOf | null;
export interface AlertmanagertypesChannelJiraConfigDTO {
/**
* @type string
@@ -434,7 +458,7 @@ export interface AlertmanagertypesChannelJiraConfigDTO {
*/
apiToken: string;
/**
* @type object
* @type object,null
*/
customFields?: AlertmanagertypesChannelJiraConfigDTOCustomFields;
/**
@@ -450,9 +474,9 @@ export interface AlertmanagertypesChannelJiraConfigDTO {
*/
issueType: string;
/**
* @type array
* @type array,null
*/
labels?: string[];
labels?: string[] | null;
/**
* @type string
*/
@@ -543,17 +567,23 @@ export interface AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTy
export enum AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTypesAlertmanagertypesChannelIncidentIOConfigDTOKind {
incidentio = 'incidentio',
}
export type AlertmanagertypesChannelIncidentIOConfigDTOMetadata = {
export type AlertmanagertypesChannelIncidentIOConfigDTOMetadataAnyOf = {
[key: string]: string;
};
/**
* @nullable
*/
export type AlertmanagertypesChannelIncidentIOConfigDTOMetadata =
AlertmanagertypesChannelIncidentIOConfigDTOMetadataAnyOf | null;
export interface AlertmanagertypesChannelIncidentIOConfigDTO {
/**
* @type string
*/
description?: string;
/**
* @type object
* @type object,null
*/
metadata?: AlertmanagertypesChannelIncidentIOConfigDTOMetadata;
/**
@@ -3366,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;
};
@@ -3384,6 +3445,7 @@ export type CloudintegrationtypesAgentReportDTO = {
* @type object,null
*/
data: CloudintegrationtypesAgentReportDTOData;
syncState: CloudintegrationtypesSyncStateDTO | null;
/**
* @type integer
* @format int64
@@ -3812,6 +3874,7 @@ export interface CloudintegrationtypesGettableAgentCheckInDTO {
* @format date-time
*/
removedAt: string | null;
syncState: CloudintegrationtypesSyncStateDTO | null;
}
export interface CloudintegrationtypesServiceMetadataDTO {
@@ -3882,6 +3945,10 @@ export interface CloudintegrationtypesPostableAgentCheckInDTO {
* @type string
*/
providerAccountId?: string;
/**
* @type integer,null
*/
syncedVersion?: number | null;
}
export interface CloudintegrationtypesStorableIntegrationDashboardDTO {

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

@@ -187,7 +187,7 @@ func (store *config) ListChannels(ctx context.Context, orgID string, params *ale
}
if !params.Kind.IsZero() {
q = q.Where("type = ?", params.Kind.ToStoredType())
q = q.Where("type = ?", params.Kind.StringValue())
}
q = q.

View File

@@ -98,11 +98,12 @@ func (handler *handler) ListChannels(rw http.ResponseWriter, req *http.Request)
}
// This ensures that the UI receives an empty array instead of null
if len(channels) == 0 {
channels = make([]*alertmanagertypes.Channel, 0)
v1Channels := make([]*alertmanagertypes.Channel, 0, len(channels))
for _, channel := range channels {
v1Channels = append(v1Channels, channel.ToV1Channel())
}
render.Success(rw, http.StatusOK, channels)
render.Success(rw, http.StatusOK, v1Channels)
}
func (handler *handler) ListAllChannels(rw http.ResponseWriter, req *http.Request) {
@@ -152,7 +153,7 @@ func (handler *handler) GetChannelByID(rw http.ResponseWriter, req *http.Request
return
}
render.Success(rw, http.StatusOK, channel)
render.Success(rw, http.StatusOK, channel.ToV1Channel())
}
func (handler *handler) UpdateChannelByID(rw http.ResponseWriter, req *http.Request) {
@@ -272,7 +273,7 @@ func (handler *handler) CreateChannel(rw http.ResponseWriter, req *http.Request)
return
}
render.Success(rw, http.StatusCreated, channel)
render.Success(rw, http.StatusCreated, channel.ToV1Channel())
}
func (handler *handler) CreateRoutePolicy(rw http.ResponseWriter, req *http.Request) {

View File

@@ -262,7 +262,7 @@ func (provider *provider) CreateChannel(ctx context.Context, orgID string, recei
}
func (provider *provider) CreateNotificationChannel(ctx context.Context, orgID string, postable alertmanagertypes.PostableNotificationChannel) (*alertmanagertypes.Channel, error) {
receiver, err := postable.ToReceiver()
channel, receiver, err := postable.ToChannel(orgID)
if err != nil {
return nil, err
}
@@ -280,11 +280,6 @@ func (provider *provider) CreateNotificationChannel(ctx context.Context, orgID s
return nil, err
}
channel, err := alertmanagertypes.NewChannelFromReceiverWithName(receiver, postable.Name, orgID)
if err != nil {
return nil, err
}
err = provider.configStore.CreateChannel(ctx, channel, alertmanagertypes.WithCb(func(ctx context.Context) error {
return provider.configStore.Set(ctx, config)
}))
@@ -304,15 +299,11 @@ func (provider *provider) UpdateNotificationChannel(ctx context.Context, orgID s
return nil, err
}
receiver, err := updatable.ToReceiver(channel.DisplayName)
receiver, err := channel.UpdateFromUpdatable(updatable)
if err != nil {
return nil, err
}
if err := channel.Update(receiver); err != nil {
return nil, err
}
config, err := provider.configStore.Get(ctx, orgID)
if err != nil {
return nil, err

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

@@ -85,6 +85,7 @@ type provider struct {
querierHandler querier.Handler
serviceAccountHandler serviceaccount.Handler
serviceAccountGetter serviceaccount.Getter
userGetter user.Getter
factoryHandler factory.Handler
cloudIntegrationHandler cloudintegration.Handler
ruleStateHistoryHandler rulestatehistory.Handler
@@ -129,6 +130,7 @@ func NewFactory(
querierHandler querier.Handler,
serviceAccountHandler serviceaccount.Handler,
serviceAccountGetter serviceaccount.Getter,
userGetter user.Getter,
factoryHandler factory.Handler,
cloudIntegrationHandler cloudintegration.Handler,
ruleStateHistoryHandler rulestatehistory.Handler,
@@ -181,6 +183,7 @@ func NewFactory(
querierHandler,
serviceAccountHandler,
serviceAccountGetter,
userGetter,
factoryHandler,
cloudIntegrationHandler,
ruleStateHistoryHandler,
@@ -235,6 +238,7 @@ func newProvider(
querierHandler querier.Handler,
serviceAccountHandler serviceaccount.Handler,
serviceAccountGetter serviceaccount.Getter,
userGetter user.Getter,
factoryHandler factory.Handler,
cloudIntegrationHandler cloudintegration.Handler,
ruleStateHistoryHandler rulestatehistory.Handler,
@@ -289,6 +293,7 @@ func newProvider(
querierHandler: querierHandler,
serviceAccountHandler: serviceAccountHandler,
serviceAccountGetter: serviceAccountGetter,
userGetter: userGetter,
factoryHandler: factoryHandler,
cloudIntegrationHandler: cloudIntegrationHandler,
ruleStateHistoryHandler: ruleStateHistoryHandler,

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

@@ -6,24 +6,35 @@ import (
"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"
)
func (provider *provider) addUserRoutes(router *mux.Router) error {
if err := router.Handle("/api/v2/users", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.ListUsers), handler.OpenAPIDef{
ID: "ListUsers",
Tags: []string{"users"},
Summary: "List users v2",
Description: "This endpoint lists all users for the organization",
Request: nil,
RequestContentType: "",
Response: make([]*types.User, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodGet).GetError(); err != nil {
if err := router.Handle("/api/v2/users", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.ListUsers, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "ListUsers",
Tags: []string{"users"},
Summary: "List users v2",
Description: "This endpoint lists all users for the organization",
Request: nil,
RequestContentType: "",
Response: make([]*types.User, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbList)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbList,
Category: coretypes.ActionCategoryAccessControl,
Selector: coretypes.WildcardSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
@@ -61,20 +72,46 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/users", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.CreateUser), handler.OpenAPIDef{
ID: "CreateUser",
Tags: []string{"users"},
Summary: "Create user",
Description: "This endpoint creates a user for the organization",
Request: new(authtypes.PostableUser),
RequestContentType: "application/json",
Response: new(types.Identifiable),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusConflict},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodPost).GetError(); err != nil {
if err := router.Handle("/api/v2/users", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.CreateUser, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "CreateUser",
Tags: []string{"users"},
Summary: "Create user",
Description: "This endpoint creates a user for the organization",
Request: new(authtypes.PostableUser),
RequestContentType: "application/json",
Response: new(types.Identifiable),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusConflict},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{
coretypes.ResourceUser.Scope(coretypes.VerbCreate),
coretypes.ResourceUser.Scope(coretypes.VerbAttach),
coretypes.ResourceRole.Scope(coretypes.VerbAttach),
}),
},
handler.WithResourceDefs(
handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbCreate,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.ResponseJSONPath("data.id"),
Selector: coretypes.WildcardSelector,
},
handler.AttachDetachSiblingResourceDef{
Verb: coretypes.VerbAttach,
Category: coretypes.ActionCategoryAccessControl,
SourceResource: coretypes.ResourceUser,
SourceIDs: coretypes.OneID(coretypes.ResponseJSONPath("data.id")),
SourceSelector: coretypes.WildcardSelector,
TargetResource: coretypes.ResourceRole,
TargetIDs: coretypes.BodyJSONArray("userRoles.#.id"),
TargetSelector: provider.roleSelector,
},
),
)).Methods(http.MethodPost).GetError(); err != nil {
return err
}
@@ -95,88 +132,152 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/users/{id}", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.GetUser), handler.OpenAPIDef{
ID: "GetUser",
Tags: []string{"users"},
Summary: "Get user by user id",
Description: "This endpoint returns the user by id",
Request: nil,
RequestContentType: "",
Response: new(authtypes.UserWithRoles),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodGet).GetError(); err != nil {
if err := router.Handle("/api/v2/users/{id}", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.GetUser, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "GetUser",
Tags: []string{"users"},
Summary: "Get user by user id",
Description: "This endpoint returns the user by id",
Request: nil,
RequestContentType: "",
Response: new(authtypes.UserWithRoles),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbRead)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbRead,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.PathParam("id"),
Selector: coretypes.IDSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/users/{id}", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.UpdateUser), handler.OpenAPIDef{
ID: "UpdateUser",
Tags: []string{"users"},
Summary: "Update user v2",
Description: "This endpoint updates the user by id",
Request: new(types.UpdatableUser),
RequestContentType: "application/json",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodPut).GetError(); err != nil {
if err := router.Handle("/api/v2/users/{id}", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.UpdateUser, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "UpdateUser",
Tags: []string{"users"},
Summary: "Update user v2",
Description: "This endpoint updates the user by id",
Request: new(types.UpdatableUser),
RequestContentType: "application/json",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbUpdate)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbUpdate,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.PathParam("id"),
Selector: coretypes.IDSelector,
}),
)).Methods(http.MethodPut).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/users/{id}", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.DeleteUser), handler.OpenAPIDef{
ID: "DeleteUser",
Tags: []string{"users"},
Summary: "Delete user",
Description: "This endpoint deletes the user by id",
Request: nil,
RequestContentType: "",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodDelete).GetError(); err != nil {
if err := router.Handle("/api/v2/users/{id}", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.DeleteUser, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "DeleteUser",
Tags: []string{"users"},
Summary: "Delete user",
Description: "This endpoint deletes the user by id",
Request: nil,
RequestContentType: "",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbDelete)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbDelete,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.PathParam("id"),
Selector: coretypes.IDSelector,
}),
)).Methods(http.MethodDelete).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/users/{id}/reset_password_tokens", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.GetResetPasswordToken), handler.OpenAPIDef{
ID: "GetResetPasswordToken",
Tags: []string{"users"},
Summary: "Get reset password token for a user",
Description: "This endpoint returns the existing reset password token for a user.",
Request: nil,
RequestContentType: "",
Response: new(types.ResetPasswordToken),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodGet).GetError(); err != nil {
if err := router.Handle("/api/v2/users/{id}/reset_password_tokens", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.GetResetPasswordToken, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "GetResetPasswordToken",
Tags: []string{"users"},
Summary: "Get reset password token for a user",
Description: "This endpoint returns the existing reset password token for a user.",
Request: nil,
RequestContentType: "",
Response: new(types.ResetPasswordToken),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceFactorPassword.Scope(coretypes.VerbList)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceMetaResourceFactorPassword,
Verb: coretypes.VerbList,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.ResponseJSONPath("data.id"),
Selector: coretypes.WildcardSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/users/{id}/reset_password_tokens", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.CreateResetPasswordToken), handler.OpenAPIDef{
ID: "CreateResetPasswordToken",
Tags: []string{"users"},
Summary: "Create or regenerate reset password token for a user",
Description: "This endpoint creates or regenerates a reset password token for a user. If a valid token exists, it is returned. If expired, a new one is created.",
Request: nil,
RequestContentType: "",
Response: new(types.ResetPasswordToken),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodPut).GetError(); err != nil {
if err := router.Handle("/api/v2/users/{id}/reset_password_tokens", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.CreateResetPasswordToken, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "CreateResetPasswordToken",
Tags: []string{"users"},
Summary: "Create or regenerate reset password token for a user",
Description: "This endpoint creates or regenerates a reset password token for a user. If a valid token exists, it is returned. If expired, a new one is created.",
Request: nil,
RequestContentType: "",
Response: new(types.ResetPasswordToken),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{
coretypes.ResourceMetaResourceFactorPassword.Scope(coretypes.VerbCreate),
coretypes.ResourceUser.Scope(coretypes.VerbAttach),
}),
},
handler.WithResourceDefs(
handler.BasicResourceDef{
Resource: coretypes.ResourceMetaResourceFactorPassword,
Verb: coretypes.VerbCreate,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.ResponseJSONPath("data.id"),
Selector: coretypes.WildcardSelector,
},
handler.AttachDetachParentChildResourceDef{
Verb: coretypes.VerbAttach,
Category: coretypes.ActionCategoryAccessControl,
ParentResource: coretypes.ResourceUser,
ParentID: coretypes.PathParam("id"),
ParentSelector: coretypes.IDSelector,
ChildResource: coretypes.ResourceMetaResourceFactorPassword,
ChildIDs: coretypes.OneID(coretypes.ResponseJSONPath("data.id")),
},
),
)).Methods(http.MethodPut).GetError(); err != nil {
return err
}
@@ -209,7 +310,7 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
SecuritySchemes: []handler.OpenAPISecurityScheme{{Name: authtypes.IdentNProviderTokenizer.StringValue()}},
})).Methods(http.MethodPut).GetError(); err != nil {
return err
}
@@ -248,90 +349,190 @@ func (provider *provider) addUserRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/users/{id}/roles", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.GetRolesByUserID), handler.OpenAPIDef{
ID: "GetRolesByUserID",
Tags: []string{"users"},
Summary: "Get user roles",
Description: "This endpoint returns the user roles by user id",
Request: nil,
RequestContentType: "",
Response: make([]*authtypes.Role, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodGet).GetError(); err != nil {
if err := router.Handle("/api/v2/users/{id}/roles", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.GetRolesByUserID, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "GetRolesByUserID",
Tags: []string{"users"},
Summary: "Get user roles",
Description: "This endpoint returns the user roles by user id",
Request: nil,
RequestContentType: "",
Response: make([]*authtypes.Role, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbRead)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbRead,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.PathParam("id"),
Selector: coretypes.IDSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/roles/{id}/users", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.GetUsersByRoleID), handler.OpenAPIDef{
ID: "GetUsersByRoleID",
Tags: []string{"users"},
Summary: "Get users by role id",
Description: "This endpoint returns the users having the role by role id",
Request: nil,
RequestContentType: "",
Response: make([]*types.User, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodGet).GetError(); err != nil {
if err := router.Handle("/api/v2/roles/{id}/users", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.GetUsersByRoleID, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "GetUsersByRoleID",
Tags: []string{"users"},
Summary: "Get users by role id",
Description: "This endpoint returns the users having the role by role id",
Request: nil,
RequestContentType: "",
Response: make([]*types.User, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceRole.Scope(coretypes.VerbRead)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceRole,
Verb: coretypes.VerbRead,
Category: coretypes.ActionCategoryAccessControl,
ID: coretypes.PathParam("id"),
Selector: provider.roleSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/user_roles", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.CreateUserRole), handler.OpenAPIDef{
ID: "CreateUserRole",
Tags: []string{"users"},
Summary: "Create user role",
Description: "This endpoint assigns a role to a user",
Request: new(authtypes.PostableUserRole),
RequestContentType: "",
Response: new(types.Identifiable),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodPost).GetError(); err != nil {
if err := router.Handle("/api/v2/user_roles", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.CreateUserRole, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "CreateUserRole",
Tags: []string{"users"},
Summary: "Create user role",
Description: "This endpoint assigns a role to a user",
Request: new(authtypes.PostableUserRole),
RequestContentType: "",
Response: new(types.Identifiable),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbAttach), coretypes.ResourceRole.Scope(coretypes.VerbAttach)}),
},
handler.WithResourceDefs(handler.AttachDetachSiblingResourceDef{
Verb: coretypes.VerbAttach,
Category: coretypes.ActionCategoryAccessControl,
SourceResource: coretypes.ResourceUser,
SourceIDs: coretypes.OneID(coretypes.BodyJSONPath("userId")),
SourceSelector: coretypes.IDSelector,
TargetResource: coretypes.ResourceRole,
TargetIDs: coretypes.OneID(coretypes.BodyJSONPath("roleId")),
TargetSelector: provider.roleSelector,
}),
)).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/user_roles/{id}", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.GetUserRole), handler.OpenAPIDef{
ID: "GetUserRole",
Tags: []string{"users"},
Summary: "Get user role",
Description: "This endpoint gets an existing user role",
Request: nil,
RequestContentType: "",
Response: new(authtypes.UserRole),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodGet).GetError(); err != nil {
if err := router.Handle("/api/v2/user_roles/{id}", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.GetUserRole, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "GetUserRole",
Tags: []string{"users"},
Summary: "Get user role",
Description: "This endpoint gets an existing user role",
Request: nil,
RequestContentType: "",
Response: new(authtypes.UserRole),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbRead)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceUser,
Verb: coretypes.VerbRead,
Category: coretypes.ActionCategoryAccessControl,
ID: provider.userRoleUserIDExtractor(),
Selector: coretypes.IDSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/user_roles/{id}", handler.New(provider.authzMiddleware.AdminAccess(provider.userHandler.DeleteUserRole), handler.OpenAPIDef{
ID: "DeleteUserRole",
Tags: []string{"users"},
Summary: "Delete user role",
Description: "This endpoint revokes a role from a user",
Request: nil,
RequestContentType: "",
Response: nil,
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleAdmin),
})).Methods(http.MethodDelete).GetError(); err != nil {
if err := router.Handle("/api/v2/user_roles/{id}", handler.New(
provider.authzMiddleware.CheckResources(provider.userHandler.DeleteUserRole, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "DeleteUserRole",
Tags: []string{"users"},
Summary: "Delete user role",
Description: "This endpoint revokes a role from a user",
Request: nil,
RequestContentType: "",
Response: nil,
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceUser.Scope(coretypes.VerbDetach), coretypes.ResourceRole.Scope(coretypes.VerbDetach)}),
},
handler.WithResourceDefs(handler.AttachDetachSiblingResourceDef{
Verb: coretypes.VerbDetach,
Category: coretypes.ActionCategoryAccessControl,
SourceResource: coretypes.ResourceUser,
SourceIDs: coretypes.OneID(provider.userRoleUserIDExtractor()),
SourceSelector: coretypes.IDSelector,
TargetResource: coretypes.ResourceRole,
TargetIDs: coretypes.OneID(provider.userRoleRoleIDExtractor()),
TargetSelector: provider.roleSelector,
}),
)).Methods(http.MethodDelete).GetError(); err != nil {
return err
}
return nil
}
func (provider *provider) userRoleUserIDExtractor() coretypes.ResourceIDExtractor {
return coretypes.NewResourceIDExtractor(coretypes.PhaseRequest, func(ec coretypes.ExtractorContext) (string, error) {
if ec.Request == nil {
return "", nil
}
userRole, err := provider.userRoleFromRequest(ec.Request)
if err != nil {
return "", err
}
return userRole.UserID.String(), nil
})
}
func (provider *provider) userRoleRoleIDExtractor() coretypes.ResourceIDExtractor {
return coretypes.NewResourceIDExtractor(coretypes.PhaseRequest, func(ec coretypes.ExtractorContext) (string, error) {
if ec.Request == nil {
return "", nil
}
userRole, err := provider.userRoleFromRequest(ec.Request)
if err != nil {
return "", err
}
return userRole.RoleID.String(), nil
})
}
func (provider *provider) userRoleFromRequest(req *http.Request) (*authtypes.UserRole, error) {
claims, err := authtypes.ClaimsFromContext(req.Context())
if err != nil {
return nil, err
}
userRoleID, err := valuer.NewUUID(mux.Vars(req)["id"])
if err != nil {
return nil, err
}
return provider.userGetter.GetUserRoleByOrgIDAndID(req.Context(), valuer.MustNewUUID(claims.OrgID), userRoleID)
}

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 {
@@ -72,24 +56,23 @@ type AttachDetachSiblingResourceDef struct {
}
func (def AttachDetachSiblingResourceDef) resolveRequest(ec coretypes.ExtractorContext) []coretypes.ResolvedResource {
return []coretypes.ResolvedResource{
coretypes.NewResolvedResourceWithTarget(
def.Verb,
def.Category,
def.SourceResource,
def.SourceIDs,
def.SourceSelector,
def.TargetResource,
def.TargetIDs,
def.TargetSelector,
false,
ec,
),
resolved := coretypes.NewResolvedResourceWithTarget(
def.Verb,
def.Category,
def.SourceResource,
def.SourceIDs,
def.SourceSelector,
def.TargetResource,
def.TargetIDs,
def.TargetSelector,
false,
ec,
)
if resolved.HasNoLinks() {
return nil
}
}
func (def AttachDetachSiblingResourceDef) requiresBody() bool {
return def.SourceIDs.RequiresBody || def.TargetIDs.RequiresBody
return []coretypes.ResolvedResource{resolved}
}
// AttachDetachParentChildResourceDef authz-checks only the parent; the child
@@ -121,20 +104,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

@@ -0,0 +1,71 @@
package handler
import (
"testing"
"github.com/SigNoz/signoz/pkg/types/coretypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestAttachDetachSiblingResourceDefResolvesNothingWithoutLinks(t *testing.T) {
def := AttachDetachSiblingResourceDef{
Verb: coretypes.VerbAttach,
Category: coretypes.ActionCategoryAccessControl,
SourceResource: coretypes.ResourceUser,
SourceIDs: coretypes.OneID(coretypes.ResponseJSONPath("data.id")),
SourceSelector: coretypes.WildcardSelector,
TargetResource: coretypes.ResourceRole,
TargetIDs: coretypes.BodyJSONArray("userRoles.#.id"),
TargetSelector: coretypes.IDSelector,
}
testCases := []struct {
name string
body string
expectedResolved int
expectedTargetIDs []string
}{
{name: "NoRolesKey_ResolvesNothing", body: `{"email":"jane@example.com"}`, expectedResolved: 0},
{name: "EmptyRoles_ResolvesNothing", body: `{"userRoles":[]}`, expectedResolved: 0},
{name: "OneRole_ResolvesOne", body: `{"userRoles":[{"id":"signoz-viewer"}]}`, expectedResolved: 1, expectedTargetIDs: []string{"signoz-viewer"}},
{name: "TwoRoles_ResolvesOneWithBothTargets", body: `{"userRoles":[{"id":"signoz-viewer"},{"id":"signoz-editor"}]}`, expectedResolved: 1, expectedTargetIDs: []string{"signoz-viewer", "signoz-editor"}},
{name: "EmptyRoleID_KeepsFailingClosed", body: `{"userRoles":[{"id":""}]}`, expectedResolved: 1, expectedTargetIDs: []string{""}},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
resolved := ResolveRequest([]ResourceDef{def}, coretypes.ExtractorContext{RequestBody: []byte(testCase.body)})
require.Len(t, resolved, testCase.expectedResolved)
if testCase.expectedResolved == 0 {
return
}
withTarget, ok := resolved[0].(coretypes.ResolvedResourceWithTargetResource)
require.True(t, ok)
assert.NoError(t, withTarget.Err())
assert.Equal(t, []string{""}, withTarget.SourceIDs())
assert.Equal(t, testCase.expectedTargetIDs, withTarget.TargetIDs())
})
}
}
func TestAttachDetachSiblingResourceDefKeepsEmptySingleID(t *testing.T) {
def := AttachDetachSiblingResourceDef{
Verb: coretypes.VerbAttach,
Category: coretypes.ActionCategoryAccessControl,
SourceResource: coretypes.ResourceUser,
SourceIDs: coretypes.OneID(coretypes.ResponseJSONPath("data.id")),
SourceSelector: coretypes.WildcardSelector,
TargetResource: coretypes.ResourceRole,
TargetIDs: coretypes.OneID(coretypes.BodyJSONPath("roleId")),
TargetSelector: coretypes.IDSelector,
}
resolved := ResolveRequest([]ResourceDef{def}, coretypes.ExtractorContext{RequestBody: []byte(`{"userId":"u1"}`)})
require.Len(t, resolved, 1)
withTarget, ok := resolved[0].(coretypes.ResolvedResourceWithTargetResource)
require.True(t, ok)
assert.Equal(t, []string{""}, withTarget.TargetIDs())
}

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

@@ -184,11 +184,6 @@ func (handler *handler) UpdateUser(w http.ResponseWriter, r *http.Request) {
return
}
if userID == claims.UserID {
render.Error(w, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "users cannot call this api on self"))
return
}
updatableUser := new(types.UpdatableUser)
if err := json.NewDecoder(r.Body).Decode(&updatableUser); err != nil {
render.Error(w, err)
@@ -431,11 +426,6 @@ func (handler *handler) CreateUserRole(w http.ResponseWriter, r *http.Request) {
return
}
if req.UserID.String() == claims.UserID {
render.Error(w, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "users cannot call this api on self"))
return
}
userRole, err := handler.setter.AddUserRoleByRoleID(ctx, valuer.MustNewUUID(claims.OrgID), req.UserID, req.RoleID)
if err != nil {
render.Error(w, err)
@@ -492,11 +482,6 @@ func (handler *handler) DeleteUserRole(w http.ResponseWriter, r *http.Request) {
return
}
if userRole.UserID.String() == claims.UserID {
render.Error(w, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "users cannot call this api on self"))
return
}
if err := handler.setter.RemoveUserRole(ctx, valuer.MustNewUUID(claims.OrgID), userRole.UserID, userRole.RoleID); err != nil {
render.Error(w, err)
return

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

@@ -94,6 +94,7 @@ func NewOpenAPI(ctx context.Context, instrumentation instrumentation.Instrumenta
struct{ querier.Handler }{},
struct{ serviceaccount.Handler }{},
struct{ serviceaccount.Getter }{},
struct{ user.Getter }{},
struct{ factory.Handler }{},
struct{ cloudintegration.Handler }{},
struct{ rulestatehistory.Handler }{},

View File

@@ -257,6 +257,8 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddCloudIntegrationTuplesFactory(sqlstore),
sqlmigration.NewAddNotificationChannelTuplesFactory(sqlstore),
sqlmigration.NewAddAIObservabilityQuickFiltersFactory(sqlstore),
sqlmigration.NewAddChannelSpecFactory(sqlschema),
sqlmigration.NewAddUserTuplesFactory(sqlstore),
)
}
@@ -352,6 +354,7 @@ func NewAPIServerProviderFactories(orgGetter organization.Getter, authz authz.Au
handlers.QuerierHandler,
handlers.ServiceAccountHandler,
modules.ServiceAccountGetter,
modules.UserGetter,
handlers.RegistryHandler,
handlers.CloudIntegrationHandler,
handlers.RuleStateHistory,

View File

@@ -0,0 +1,737 @@
package sqlmigration
import (
"context"
"encoding/json"
"log/slog"
"maps"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlschema"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
)
var channelSpecBackfillKinds = []channelSpecBackfillKind{
{kind: "slack", configsKey: "slack_configs", convert: convertSlackNotifierJSON},
{kind: "email", configsKey: "email_configs", convert: convertEmailNotifierJSON},
{kind: "webhook", configsKey: "webhook_configs", convert: convertWebhookNotifierJSON},
{kind: "pagerduty", configsKey: "pagerduty_configs", convert: convertPagerdutyNotifierJSON},
{kind: "opsgenie", configsKey: "opsgenie_configs", convert: convertOpsgenieNotifierJSON},
{kind: "msteams", configsKey: "msteamsv2_configs", convert: convertMSTeamsNotifierJSON},
{kind: "googlechat", configsKey: "googlechat_configs", convert: convertGoogleChatNotifierJSON},
{kind: "jira", configsKey: "jira_configs", convert: convertJiraNotifierJSON},
{kind: "jsmops", configsKey: "jsmops_configs", convert: convertJSMOpsNotifierJSON},
{kind: "incidentio", configsKey: "incidentio_configs", convert: convertIncidentIONotifierJSON},
}
func NewAddChannelSpecFactory(sqlschema sqlschema.SQLSchema) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(
factory.MustNewName("add_channel_spec"),
func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addChannelSpec{sqlschema: sqlschema, logger: ps.Logger}, nil
},
)
}
func (migration *addChannelSpec) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
// Up adds the column and fills it from each channel's receiver, as a write
// through a receiver does, pinning type to the kind the spec came from because
// a read decodes the spec under it. The msteams kind, which v1 stored as
// msteamsv2 after upstream's configs list, is renamed on every row. A receiver
// v2 cannot represent, such as one carrying several notifiers or a notifier
// kind v2 does not model, stays NULL and is logged; the repair endpoint is the
// remedy for those.
func (migration *addChannelSpec) Up(ctx context.Context, db *bun.DB) error {
table, uniqueConstraints, err := migration.sqlschema.GetTable(ctx, sqlschema.TableName("notification_channel"))
if err != nil {
return err
}
sqls := migration.sqlschema.Operator().AddColumn(table, uniqueConstraints, &sqlschema.Column{
Name: sqlschema.ColumnName("spec"),
DataType: sqlschema.DataTypeText,
Nullable: true,
}, nil)
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() {
_ = tx.Rollback()
}()
for _, sql := range sqls {
if _, err := tx.ExecContext(ctx, string(sql)); err != nil {
return err
}
}
if _, err := tx.NewUpdate().Model((*channelSpecBackfillRow)(nil)).Set("type = ?", "msteams").Where("type = ?", "msteamsv2").Exec(ctx); err != nil {
return err
}
rows := make([]*channelSpecBackfillRow, 0)
if err := tx.NewSelect().Model(&rows).Where("spec IS NULL").OrderExpr("org_id, id").Scan(ctx); err != nil {
return err
}
type orgStats struct{ total, filled, unrepresentable int }
statsByOrg := map[string]*orgStats{}
for _, row := range rows {
stats, ok := statsByOrg[row.OrgID]
if !ok {
stats = &orgStats{}
statsByOrg[row.OrgID] = stats
}
stats.total++
storedType, spec, err := channelSpecFromReceiverJSON(row.Data)
if err != nil {
stats.unrepresentable++
migration.logger.WarnContext(ctx, "leaving notification channel without a v2 spec", slog.String("org_id", row.OrgID), slog.String("channel_id", row.ID), errors.Attr(err))
continue
}
encoded, err := marshalUnescaped(spec)
if err != nil {
return err
}
if _, err := tx.NewUpdate().
Model((*channelSpecBackfillRow)(nil)).
Set("spec = ?", string(encoded)).
Set("type = ?", storedType).
Where("id = ?", row.ID).
Exec(ctx); err != nil {
return err
}
stats.filled++
}
for _, orgID := range slices.Sorted(maps.Keys(statsByOrg)) {
stats := statsByOrg[orgID]
migration.logger.InfoContext(ctx, "filled v2 spec on notification channels", slog.String("org_id", orgID), slog.Int("total", stats.total), slog.Int("filled", stats.filled), slog.Int("unrepresentable", stats.unrepresentable))
}
return tx.Commit()
}
func (migration *addChannelSpec) Down(context.Context, *bun.DB) error {
return nil
}
type addChannelSpec struct {
sqlschema sqlschema.SQLSchema
logger *slog.Logger
}
type channelSpecBackfillRow struct {
bun.BaseModel `bun:"table:notification_channel"`
ID string `bun:"id,pk"`
OrgID string `bun:"org_id"`
Data string `bun:"data"`
}
// notifierJSON is one entry of a receiver's *_configs list as stored in
// notification_channel.data.
type notifierJSON map[string]json.RawMessage
type channelSpecBackfillKind struct {
kind string
configsKey string
convert func(notifierJSON) (map[string]any, error)
}
// channelSpecFromReceiverJSON mirrors the v2 read of a stored receiver: one
// notifier of a modelled kind, with the receiver's field names renamed to the
// spec's and its unset templates left out. The type alongside is what a read
// decodes the spec under.
func channelSpecFromReceiverJSON(data string) (string, map[string]any, error) {
receiver := map[string]json.RawMessage{}
if err := json.Unmarshal([]byte(data), &receiver); err != nil {
return "", nil, err
}
total := 0
var found *channelSpecBackfillKind
var notifier notifierJSON
for key, raw := range receiver {
if !strings.HasSuffix(key, "_configs") {
continue
}
var list []notifierJSON
if err := json.Unmarshal(raw, &list); err != nil {
return "", nil, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "%s", key)
}
total += len(list)
if len(list) == 0 {
continue
}
for i := range channelSpecBackfillKinds {
if channelSpecBackfillKinds[i].configsKey == key {
found = &channelSpecBackfillKinds[i]
notifier = list[0]
}
}
}
if total > 1 {
return "", nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "carries %d notifier configurations; only one per channel is supported", total)
}
if found == nil {
return "", nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "carries no supported notifier configuration")
}
spec, err := found.convert(notifier)
if err != nil {
return "", nil, err
}
sendResolved, err := notifier.boolValue("send_resolved")
if err != nil {
return "", nil, err
}
spec["sendResolved"] = sendResolved
return found.kind, spec, nil
}
func convertSlackNotifierJSON(notifier notifierJSON) (map[string]any, error) {
if err := rejectAnyHTTPAuthJSON(notifier); err != nil {
return nil, err
}
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"api_url": "apiUrl", "channel": "channel"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"title": "title", "text": "text", "color": "color", "title_link": "titleLink", "pretext": "pretext", "fallback": "fallback", "footer": "footer"}); err != nil {
return nil, err
}
if spec["apiUrl"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "slack: api_url is required")
}
fields, err := notifier.objectList("fields")
if err != nil {
return nil, err
}
if len(fields) > 0 {
converted := make([]map[string]any, 0, len(fields))
for _, field := range fields {
item := map[string]any{}
if err := field.copyStrings(item, map[string]string{"title": "title", "value": "value"}); err != nil {
return nil, err
}
if field.has("short") {
short, err := field.boolValue("short")
if err != nil {
return nil, err
}
item["short"] = short
}
converted = append(converted, item)
}
spec["fields"] = converted
}
actions, err := notifier.objectList("actions")
if err != nil {
return nil, err
}
if len(actions) > 0 {
converted := make([]map[string]any, 0, len(actions))
for _, action := range actions {
item := map[string]any{}
if err := action.copyStrings(item, map[string]string{"type": "type", "text": "text", "url": "url", "style": "style", "name": "name", "value": "value"}); err != nil {
return nil, err
}
if action.has("confirm") {
confirm, err := action.object("confirm")
if err != nil {
return nil, err
}
confirmation := map[string]any{}
if err := confirm.copyStrings(confirmation, map[string]string{"text": "text", "title": "title", "ok_text": "okText", "dismiss_text": "dismissText"}); err != nil {
return nil, err
}
item["confirm"] = confirmation
}
converted = append(converted, item)
}
spec["actions"] = converted
}
return spec, nil
}
func convertEmailNotifierJSON(notifier notifierJSON) (map[string]any, error) {
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"to": "to"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"html": "html"}); err != nil {
return nil, err
}
if spec["to"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "email: to is required")
}
if err := notifier.copyNonEmptyObjects(spec, map[string]string{"headers": "headers"}); err != nil {
return nil, err
}
return spec, nil
}
func convertWebhookNotifierJSON(notifier notifierJSON) (map[string]any, error) {
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"url": "url"}); err != nil {
return nil, err
}
if spec["url"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "webhook: url is required")
}
if err := rejectUnsupportedHTTPConfigJSON(notifier); err != nil {
return nil, err
}
httpConfig, err := notifier.object("http_config")
if err != nil {
return nil, err
}
username, password, err := extractBasicAuthJSON(httpConfig)
if err != nil {
return nil, err
}
bearerToken, err := extractBearerTokenJSON(httpConfig)
if err != nil {
return nil, err
}
if (username != "" || password != "") && bearerToken != "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "webhook: basic auth and bearer token cannot be combined")
}
spec["username"], spec["password"], spec["bearerToken"] = username, password, bearerToken
return spec, nil
}
func convertPagerdutyNotifierJSON(notifier notifierJSON) (map[string]any, error) {
if err := rejectAnyHTTPAuthJSON(notifier); err != nil {
return nil, err
}
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"routing_key": "routingKey", "url": "url", "severity": "severity", "component": "component", "group": "group", "class": "class"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"source": "source", "client": "client", "client_url": "clientUrl", "description": "description"}); err != nil {
return nil, err
}
if spec["routingKey"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "pagerduty: routing_key is required")
}
details, err := notifier.object("details")
if err != nil {
return nil, err
}
if len(details) > 0 {
for key, raw := range details {
var value string
if err := json.Unmarshal(raw, &value); err != nil {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "pagerduty: details.%s is not a string", key)
}
}
spec["details"] = notifier["details"]
}
return spec, nil
}
func convertOpsgenieNotifierJSON(notifier notifierJSON) (map[string]any, error) {
if err := rejectAnyHTTPAuthJSON(notifier); err != nil {
return nil, err
}
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"api_key": "apiKey", "api_url": "apiUrl", "priority": "priority"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"message": "message", "description": "description", "source": "source"}); err != nil {
return nil, err
}
if spec["apiKey"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "opsgenie: api_key is required")
}
if err := notifier.copyNonEmptyObjects(spec, map[string]string{"details": "details"}); err != nil {
return nil, err
}
return spec, nil
}
func convertMSTeamsNotifierJSON(notifier notifierJSON) (map[string]any, error) {
return convertWebhookURLNotifierJSON("msteams", notifier)
}
func convertGoogleChatNotifierJSON(notifier notifierJSON) (map[string]any, error) {
return convertWebhookURLNotifierJSON("googlechat", notifier)
}
func convertWebhookURLNotifierJSON(name string, notifier notifierJSON) (map[string]any, error) {
if err := rejectAnyHTTPAuthJSON(notifier); err != nil {
return nil, err
}
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"webhook_url": "webhookUrl"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"title": "title", "text": "text"}); err != nil {
return nil, err
}
if spec["webhookUrl"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "%s: webhook_url is required", name)
}
return spec, nil
}
func convertJiraNotifierJSON(notifier notifierJSON) (map[string]any, error) {
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"site": "site", "project": "project", "issue_type": "issueType", "priority": "priority", "resolve_transition": "resolveTransition", "reopen_transition": "reopenTransition", "wont_fix_resolution": "wontFixResolution"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"summary": "summary", "description": "description", "reopen_duration": "reopenDuration"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyObjects(spec, map[string]string{"custom_fields": "customFields"}); err != nil {
return nil, err
}
labels, err := notifier.list("labels")
if err != nil {
return nil, err
}
if len(labels) > 0 {
spec["labels"] = notifier["labels"]
}
if err := rejectUnsupportedHTTPConfigJSON(notifier); err != nil {
return nil, err
}
httpConfig, err := notifier.object("http_config")
if err != nil {
return nil, err
}
if httpConfig.has("authorization") {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "jira: http_config.authorization is not supported")
}
email, apiToken, err := extractBasicAuthJSON(httpConfig)
if err != nil {
return nil, err
}
spec["email"], spec["apiToken"] = email, apiToken
for _, required := range []string{"site", "project", "issueType", "email", "apiToken"} {
if spec[required] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "jira: %s is required", required)
}
}
return spec, nil
}
func convertJSMOpsNotifierJSON(notifier notifierJSON) (map[string]any, error) {
if err := rejectAnyHTTPAuthJSON(notifier); err != nil {
return nil, err
}
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"api_key": "apiKey", "priority": "priority"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"message": "message", "description": "description", "tags": "tags"}); err != nil {
return nil, err
}
if spec["apiKey"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "jsmops: api_key is required")
}
return spec, nil
}
func convertIncidentIONotifierJSON(notifier notifierJSON) (map[string]any, error) {
if err := rejectAnyHTTPAuthJSON(notifier); err != nil {
return nil, err
}
spec := map[string]any{}
if err := notifier.copyStrings(spec, map[string]string{"url": "url", "token": "token"}); err != nil {
return nil, err
}
if err := notifier.copyNonEmptyStrings(spec, map[string]string{"title": "title", "description": "description"}); err != nil {
return nil, err
}
if spec["url"] == "" || spec["token"] == "" {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "incidentio: url and token are required")
}
if err := notifier.copyNonEmptyObjects(spec, map[string]string{"metadata": "metadata"}); err != nil {
return nil, err
}
return spec, nil
}
func rejectAnyHTTPAuthJSON(notifier notifierJSON) error {
httpConfig, err := notifier.object("http_config")
if err != nil {
return err
}
if httpConfig.has("basic_auth") {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.basic_auth is not supported")
}
if httpConfig.has("authorization") {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.authorization is not supported")
}
return rejectUnsupportedHTTPConfigJSON(notifier)
}
// rejectUnsupportedHTTPConfigJSON refuses every http_config setting the spec has
// no field for, since a config that dropped it would unauthenticate or reroute
// the channel on the next write. An absent follow_redirects or enable_http2 is
// the upstream default, true; only an explicit false is refused.
func rejectUnsupportedHTTPConfigJSON(notifier notifierJSON) error {
if !notifier.has("http_config") {
return nil
}
httpConfig, err := notifier.object("http_config")
if err != nil {
return err
}
for _, key := range []string{"oauth2", "http_headers"} {
if httpConfig.has(key) {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.%s is not supported", key)
}
}
for _, key := range []string{"bearer_token", "bearer_token_file", "proxy_url", "no_proxy"} {
value, err := httpConfig.stringValue(key)
if err != nil {
return err
}
if value != "" {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.%s is not supported", key)
}
}
if httpConfig.has("proxy_from_environment") {
fromEnvironment, err := httpConfig.boolValue("proxy_from_environment")
if err != nil {
return err
}
if fromEnvironment {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.proxy_from_environment is not supported")
}
}
tlsConfig, err := httpConfig.object("tls_config")
if err != nil {
return err
}
for key := range tlsConfig {
if key != "insecure_skip_verify" {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.tls_config is not supported")
}
}
if tlsConfig.has("insecure_skip_verify") {
insecure, err := tlsConfig.boolValue("insecure_skip_verify")
if err != nil {
return err
}
if insecure {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.tls_config is not supported")
}
}
for _, key := range []string{"follow_redirects", "enable_http2"} {
if !httpConfig.has(key) {
continue
}
enabled, err := httpConfig.boolValue(key)
if err != nil {
return err
}
if !enabled {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.%s cannot be disabled", key)
}
}
return nil
}
func extractBasicAuthJSON(httpConfig notifierJSON) (string, string, error) {
basicAuth, err := httpConfig.object("basic_auth")
if err != nil {
return "", "", err
}
if len(basicAuth) == 0 {
return "", "", nil
}
for key := range basicAuth {
if key != "username" && key != "password" {
return "", "", errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.basic_auth.%s is not supported", key)
}
}
username, err := basicAuth.stringValue("username")
if err != nil {
return "", "", err
}
password, err := basicAuth.stringValue("password")
if err != nil {
return "", "", err
}
return username, password, nil
}
func extractBearerTokenJSON(httpConfig notifierJSON) (string, error) {
if !httpConfig.has("authorization") {
return "", nil
}
authorization, err := httpConfig.object("authorization")
if err != nil {
return "", err
}
for key := range authorization {
if key != "type" && key != "credentials" {
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.authorization.%s is not supported", key)
}
}
scheme, err := authorization.stringValue("type")
if err != nil {
return "", err
}
if !strings.EqualFold(scheme, "Bearer") {
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "http_config.authorization.type %q is not supported", scheme)
}
return authorization.stringValue("credentials")
}
// has reports a key that is present and not null.
func (n notifierJSON) has(key string) bool {
raw, ok := n[key]
return ok && string(raw) != "null"
}
func (n notifierJSON) stringValue(key string) (string, error) {
if !n.has(key) {
return "", nil
}
var value string
if err := json.Unmarshal(n[key], &value); err != nil {
return "", errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "%s", key)
}
return value, nil
}
func (n notifierJSON) boolValue(key string) (bool, error) {
if !n.has(key) {
return false, nil
}
var value bool
if err := json.Unmarshal(n[key], &value); err != nil {
return false, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "%s", key)
}
return value, nil
}
// object reads an absent key as an empty object; a caller that must tell the
// two apart checks has first.
func (n notifierJSON) object(key string) (notifierJSON, error) {
if !n.has(key) {
return notifierJSON{}, nil
}
value := notifierJSON{}
if err := json.Unmarshal(n[key], &value); err != nil {
return nil, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "%s", key)
}
return value, nil
}
func (n notifierJSON) list(key string) ([]json.RawMessage, error) {
if !n.has(key) {
return nil, nil
}
var value []json.RawMessage
if err := json.Unmarshal(n[key], &value); err != nil {
return nil, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "%s", key)
}
return value, nil
}
func (n notifierJSON) objectList(key string) ([]notifierJSON, error) {
if !n.has(key) {
return nil, nil
}
var value []notifierJSON
if err := json.Unmarshal(n[key], &value); err != nil {
return nil, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "%s", key)
}
return value, nil
}
// copyStrings writes each field as the spec's plain string, "" when absent.
func (n notifierJSON) copyStrings(spec map[string]any, keys map[string]string) error {
for from, to := range keys {
value, err := n.stringValue(from)
if err != nil {
return err
}
spec[to] = value
}
return nil
}
// copyNonEmptyStrings leaves an empty field out, which is how the spec spells
// an unset template.
func (n notifierJSON) copyNonEmptyStrings(spec map[string]any, keys map[string]string) error {
for from, to := range keys {
value, err := n.stringValue(from)
if err != nil {
return err
}
if value != "" {
spec[to] = value
}
}
return nil
}
func (n notifierJSON) copyNonEmptyObjects(spec map[string]any, keys map[string]string) error {
for from, to := range keys {
value, err := n.object(from)
if err != nil {
return err
}
if len(value) > 0 {
spec[to] = n[from]
}
}
return nil
}

View File

@@ -0,0 +1,141 @@
package sqlmigration
import (
"context"
"database/sql"
"time"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/types/authtypes"
"github.com/oklog/ulid/v2"
"github.com/uptrace/bun"
"github.com/uptrace/bun/dialect"
"github.com/uptrace/bun/migrate"
)
type addUserTuples struct {
sqlstore sqlstore.SQLStore
}
func NewAddUserTuplesFactory(sqlstore sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(factory.MustNewName("add_user_tuples"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addUserTuples{sqlstore: sqlstore}, nil
})
}
func (migration *addUserTuples) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func (migration *addUserTuples) Up(ctx context.Context, db *bun.DB) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
var storeID string
err = tx.QueryRowContext(ctx, `SELECT id FROM store WHERE name = ? LIMIT 1`, "signoz").Scan(&storeID)
if err != nil {
return err
}
var orgIDs []string
err = tx.NewSelect().
Table("organizations").
Column("id").
Scan(ctx, &orgIDs)
if err != nil && err != sql.ErrNoRows {
return err
}
isPG := migration.sqlstore.BunDB().Dialect().Name() == dialect.PG
// user and factor-password moved from the legacy AdminAccess gate to
// CheckResources. Existing organizations never had these tuples written;
// only new organizations receive them from the managed-role registry at bootstrap.
tuples := []migrationTuple{
{authtypes.SigNozAdminRoleName, "user", "user", "create"},
{authtypes.SigNozAdminRoleName, "user", "user", "list"},
{authtypes.SigNozAdminRoleName, "user", "user", "read"},
{authtypes.SigNozAdminRoleName, "user", "user", "update"},
{authtypes.SigNozAdminRoleName, "user", "user", "delete"},
{authtypes.SigNozAdminRoleName, "user", "user", "attach"},
{authtypes.SigNozAdminRoleName, "user", "user", "detach"},
{authtypes.SigNozAdminRoleName, "metaresource", "factor-password", "read"},
{authtypes.SigNozAdminRoleName, "metaresource", "factor-password", "create"},
{authtypes.SigNozAdminRoleName, "metaresource", "factor-password", "list"},
}
for _, orgID := range orgIDs {
for _, tuple := range tuples {
entropy := ulid.DefaultEntropy()
now := time.Now().UTC()
tupleID := ulid.MustNew(ulid.Timestamp(now), entropy).String()
objectID := "organization/" + orgID + "/" + tuple.objectName + "/*"
roleSubject := "organization/" + orgID + "/role/" + tuple.roleName
if isPG {
user := "role:" + roleSubject + "#assignee"
result, err := tx.ExecContext(ctx, `
INSERT INTO tuple (store, object_type, object_id, relation, _user, user_type, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, object_type, object_id, relation, _user) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, user, "userset", tupleID, now,
)
if err != nil {
return err
}
rowsAffected, err := result.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
continue
}
_, err = tx.ExecContext(ctx, `
INSERT INTO changelog (store, object_type, object_id, relation, _user, operation, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, ulid, object_type) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, user, 0, tupleID, now,
)
if err != nil {
return err
}
} else {
result, err := tx.ExecContext(ctx, `
INSERT INTO tuple (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation, user_type, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, "role", roleSubject, "assignee", "userset", tupleID, now,
)
if err != nil {
return err
}
rowsAffected, err := result.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
continue
}
_, err = tx.ExecContext(ctx, `
INSERT INTO changelog (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation, operation, ulid, inserted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (store, ulid, object_type) DO NOTHING`,
storeID, tuple.objectType, objectID, tuple.relation, "role", roleSubject, "assignee", 0, tupleID, now,
)
if err != nil {
return err
}
}
}
}
return tx.Commit()
}
func (migration *addUserTuples) Down(context.Context, *bun.DB) error {
return nil
}

View File

@@ -1,18 +1,14 @@
package alertmanagertypes
import (
"context"
"crypto/rand"
"encoding/json"
"reflect"
"regexp"
"strings"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
"github.com/swaggest/jsonschema-go"
"github.com/uptrace/bun"
)
@@ -30,20 +26,6 @@ var (
type Channels = []*Channel
type GettableChannels = []*Channel
// TODO: the oneOf emitted by JSONSchema is not the shape OpenAPI wants for a
// discriminated union. OpenAPI's discriminator requires every oneOf branch to
// be a $ref to a named component and a sibling property whose value selects
// the variant. Our payload instead uses the *presence* of one of the 18
// *_configs arrays to imply the type, so no discriminator can be attached.
// Refactor PostableChannel into a {name, type, config} envelope (see
// ruletypes.RuleThresholdData for the pattern) so each notification kind
// becomes a named component and the discriminator can be wired up properly.
type PostableChannel struct {
Receiver
}
// Channel represents a single receiver of the alertmanager config.
type Channel struct {
bun.BaseModel `bun:"table:notification_channel"`
@@ -55,45 +37,51 @@ type Channel struct {
// reference, so it keeps the v1 wire tag and Name stays off the v1 contract.
Name string `json:"-" bun:"name"`
DisplayName string `json:"name" required:"true" bun:"display_name"`
Type string `json:"type" required:"true" bun:"type"`
Data string `json:"data" required:"true" bun:"data"`
OrgID string `json:"orgId" required:"true" bun:"org_id"`
// TODO: type this as ChannelKind once v1 is gone.
Type string `json:"type" required:"true" bun:"type"`
Data string `json:"data" required:"true" bun:"data"`
OrgID string `json:"orgId" required:"true" bun:"org_id"`
// Spec is the v2 spec a read returns, of the kind Type names. A v2 write
// stores it as the caller wrote it and a v1 write derives it from the
// defaulted receiver. Only a row the migration could not backfill has none.
// StoredSpec is the column, written by fillSpec and decoded by AfterScanRow.
Spec ChannelSpec `json:"-" bun:"-"`
StoredSpec string `json:"-" bun:"spec,type:text,nullzero"`
}
// NewChannelFromReceiver creates a new Channel from a Receiver.
// It can return nil if the receiver is the default receiver.
// A receiver carries no internal name, so one is generated from its name.
func NewChannelFromReceiver(receiver *Receiver, orgID string) (*Channel, error) {
if receiver.Name == DefaultReceiverName {
return nil, errors.Newf(errors.TypeInvalidInput, ErrCodeAlertmanagerChannelInvalid, "cannot use %s name as a channel name", receiver.Name)
var _ bun.AfterScanRowHook = (*Channel)(nil)
// AfterScanRow decodes the stored spec under Type, which bun cannot do column by
// column because the spec's Go type depends on it.
func (c *Channel) AfterScanRow(context.Context) error {
if c.StoredSpec == "" {
c.Spec = nil
return nil
}
// Initialize channel with common fields
channel := Channel{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
TimeAuditable: types.TimeAuditable{
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
},
Name: generateChannelName(receiver.Name),
DisplayName: receiver.Name,
OrgID: orgID,
channelKind, ok := parseChannelKind(c.Type)
if !ok {
return errors.NewInternalf(errors.CodeInternal, "channel %q stores a spec under unmodelled type %q", c.DisplayName, c.Type)
}
data, err := json.Marshal(receiver)
spec, _ := buildEmptyChannelSpecForKind(channelKind)
if err := json.Unmarshal([]byte(c.StoredSpec), spec); err != nil {
return errors.WrapInternalf(err, errors.CodeInternal, "unmarshal channel %q spec", c.DisplayName)
}
c.Spec = spec
return nil
}
func (c *Channel) fillSpec(spec ChannelSpec) error {
stored, err := json.Marshal(spec)
if err != nil {
return nil, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "marshal receiver")
return errors.WrapInternalf(err, errors.CodeInternal, "marshal channel %q spec", c.DisplayName)
}
channel.Data = string(data)
c.Spec, c.StoredSpec = spec, string(stored)
channel.Type = receiverChannelType(receiver)
if channel.Type == "" {
return nil, errors.Newf(errors.TypeInvalidInput, ErrCodeAlertmanagerChannelInvalid, "channel '%s' must have at least one notification configuration (e.g., email_configs, webhook_configs, slack_configs)", receiver.Name)
}
return &channel, nil
return nil
}
const channelNameSuffixLen = 8
@@ -137,57 +125,6 @@ func generateChannelName(displayName string) string {
return prefix + "-" + string(suffix)
}
// NewChannelFromReceiverWithName overrides the name that NewChannelFromReceiver
// generates.
func NewChannelFromReceiverWithName(receiver *Receiver, name string, orgID string) (*Channel, error) {
channel, err := NewChannelFromReceiver(receiver, orgID)
if err != nil {
return nil, err
}
channel.Name = name
return channel, nil
}
// receiverChannelType returns the channel.Type discriminator. Walks
// Receiver's own fields first (native), then the embed (upstream); first
// non-empty *_configs slice wins.
func receiverChannelType(receiver *Receiver) string {
if t := nonEmptyConfigsField(reflect.ValueOf(*receiver)); t != "" {
return t
}
if t := nonEmptyConfigsField(reflect.ValueOf(*receiver.Receiver)); t != "" {
return t
}
return ""
}
func nonEmptyConfigsField(v reflect.Value) string {
t := v.Type()
for i := 0; i < t.NumField(); i++ {
field := t.Field(i)
fieldVal := v.Field(i)
if fieldVal.Kind() != reflect.Slice || fieldVal.Len() == 0 {
continue
}
yamlTag := field.Tag.Get("yaml")
if yamlTag == "" {
continue
}
// Extract the base type name (e.g., "email_configs" -> "email").
matches := receiverTypeRegex.FindStringSubmatch(yamlTag)
if len(matches) != 2 {
continue
}
return matches[1]
}
return ""
}
func NewConfigFromChannels(globalConfig GlobalConfig, routeConfig RouteConfig, channels Channels, orgID string) (*Config, error) {
cfg, err := NewDefaultConfig(
globalConfig,
@@ -228,64 +165,3 @@ func NewStatsFromChannels(channels Channels) map[string]any {
stats["alertmanager.channel.count"] = int64(len(channels))
return stats
}
func (c *Channel) Update(receiver *Receiver) error {
channel, err := NewChannelFromReceiverWithName(receiver, c.Name, c.OrgID)
if err != nil {
return err
}
if c.DisplayName != channel.DisplayName {
return errors.Newf(errors.TypeInvalidInput, ErrCodeAlertmanagerChannelNameMismatch, "cannot update channel name")
}
// Unreachable while the name is passed in above rather than derived from the
// receiver, which is why this is internal rather than invalid input.
if c.Name != channel.Name {
return errors.NewInternalf(ErrCodeAlertmanagerChannelNameMismatch, "cannot update channel internal name")
}
c.Type = channel.Type
c.Data = channel.Data
c.UpdatedAt = time.Now()
return nil
}
func (PostableChannel) JSONSchema() (jsonschema.Schema, error) {
type alias PostableChannel
reflector := &jsonschema.Reflector{}
schema, err := reflector.Reflect(alias{}, jsonschema.DefinitionsPrefix("#/components/schemas/"))
if err != nil {
return jsonschema.Schema{}, err
}
schema.WithRequired("name")
var oneOf []jsonschema.SchemaOrBool
seen := map[string]struct{}{}
// Walk both halves: native fields on Receiver, upstream on the embed. A native
// field can shadow an upstream one with the same tag (e.g. jira_configs), so
// dedupe to avoid emitting two identical oneOf branches.
collect := func(t reflect.Type) {
for i := 0; i < t.NumField(); i++ {
jsonTag := strings.Split(t.Field(i).Tag.Get("json"), ",")[0]
if !strings.HasSuffix(jsonTag, "_configs") {
continue
}
if _, ok := seen[jsonTag]; ok {
continue
}
seen[jsonTag] = struct{}{}
branch := (&jsonschema.Schema{}).WithRequired(jsonTag)
oneOf = append(oneOf, branch.ToSchemaOrBool())
}
}
collect(reflect.TypeOf(Receiver{}))
collect(reflect.TypeOf(config.Receiver{}))
schema.WithOneOf(oneOf...)
return schema, nil
}

View File

@@ -16,8 +16,20 @@ import (
type ChannelEmailConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
To string `json:"to" required:"true"`
HTML valuer.UnsetOrNonEmptyString `json:"html"`
Headers map[string]string `json:"headers,omitempty"`
HTML valuer.UnsetOrNonEmptyString `json:"html,omitzero"`
Headers map[string]string `json:"headers,omitzero"`
}
func (c *ChannelEmailConfig) UnmarshalJSON(data []byte) error {
type alias ChannelEmailConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, config.DefaultEmailConfig.VSendResolved)
c.HTML.SetIfUnset(config.DefaultEmailConfig.HTML)
return c.Validate()
}
func (c ChannelEmailConfig) Validate() error {

View File

@@ -10,8 +10,21 @@ import (
type ChannelGoogleChatConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
Title valuer.UnsetOrNonEmptyString `json:"title,omitzero"`
Text valuer.UnsetOrNonEmptyString `json:"text,omitzero"`
}
func (c *ChannelGoogleChatConfig) UnmarshalJSON(data []byte) error {
type alias ChannelGoogleChatConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, DefaultGoogleChatReceiverConfig.VSendResolved)
c.Title.SetIfUnset(DefaultGoogleChatReceiverConfig.Title)
c.Text.SetIfUnset(DefaultGoogleChatReceiverConfig.Text)
return c.Validate()
}
func (c ChannelGoogleChatConfig) Validate() error {

View File

@@ -13,9 +13,22 @@ type ChannelIncidentIOConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
URL string `json:"url" required:"true"`
Token string `json:"token" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Metadata map[string]string `json:"metadata,omitempty"`
Title valuer.UnsetOrNonEmptyString `json:"title,omitzero"`
Description valuer.UnsetOrNonEmptyString `json:"description,omitzero"`
Metadata map[string]string `json:"metadata,omitzero"`
}
func (c *ChannelIncidentIOConfig) UnmarshalJSON(data []byte) error {
type alias ChannelIncidentIOConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, DefaultIncidentIOReceiverConfig.VSendResolved)
c.Title.SetIfUnset(DefaultIncidentIOReceiverConfig.Title)
c.Description.SetIfUnset(DefaultIncidentIOReceiverConfig.Description)
return c.Validate()
}
func (c ChannelIncidentIOConfig) Validate() error {

View File

@@ -19,20 +19,35 @@ type ChannelJiraConfig struct {
Site string `json:"site" required:"true"`
Project string `json:"project" required:"true"`
IssueType string `json:"issueType" required:"true"`
Summary valuer.UnsetOrNonEmptyString `json:"summary"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Summary valuer.UnsetOrNonEmptyString `json:"summary,omitzero"`
Description valuer.UnsetOrNonEmptyString `json:"description,omitzero"`
Priority string `json:"priority"`
Labels []string `json:"labels,omitempty"`
Labels []string `json:"labels,omitzero"`
ResolveTransition string `json:"resolveTransition"`
ReopenTransition string `json:"reopenTransition"`
ReopenDuration valuer.UnsetOrNonEmptyString `json:"reopenDuration"`
ReopenDuration valuer.UnsetOrNonEmptyString `json:"reopenDuration,omitzero"`
WontFixResolution string `json:"wontFixResolution"`
CustomFields map[string]any `json:"customFields,omitempty"`
CustomFields map[string]any `json:"customFields,omitzero"`
Email string `json:"email" required:"true"`
APIToken string `json:"apiToken" required:"true" format:"password"`
}
// UnmarshalJSON seeds send_resolved off, as JiraReceiverConfig does.
func (c *ChannelJiraConfig) UnmarshalJSON(data []byte) error {
type alias ChannelJiraConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, false)
c.Summary.SetIfUnset(DefaultJiraSummaryTemplate)
c.Description.SetIfUnset(DefaultJiraDescriptionTemplate)
c.ReopenDuration.SetIfUnset(defaultJiraReopenDuration.String())
return c.Validate()
}
func (c ChannelJiraConfig) Validate() error {
for _, required := range []struct {
value string

View File

@@ -12,11 +12,25 @@ import (
type ChannelJSMOpsConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIKey string `json:"apiKey" required:"true" format:"password"`
Message valuer.UnsetOrNonEmptyString `json:"message"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Message valuer.UnsetOrNonEmptyString `json:"message,omitzero"`
Description valuer.UnsetOrNonEmptyString `json:"description,omitzero"`
Priority string `json:"priority"`
// Tags is the comma-separated list JSM Ops attaches to the alert.
Tags valuer.UnsetOrNonEmptyString `json:"tags"`
Tags valuer.UnsetOrNonEmptyString `json:"tags,omitzero"`
}
func (c *ChannelJSMOpsConfig) UnmarshalJSON(data []byte) error {
type alias ChannelJSMOpsConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, DefaultJSMOpsReceiverConfig.VSendResolved)
c.Message.SetIfUnset(DefaultJSMOpsReceiverConfig.Message)
c.Description.SetIfUnset(DefaultJSMOpsReceiverConfig.Description)
c.Tags.SetIfUnset(DefaultJSMOpsReceiverConfig.Tags)
return c.Validate()
}
func (c ChannelJSMOpsConfig) Validate() error {

View File

@@ -9,8 +9,21 @@ import (
type ChannelMSTeamsConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
Title valuer.UnsetOrNonEmptyString `json:"title,omitzero"`
Text valuer.UnsetOrNonEmptyString `json:"text,omitzero"`
}
func (c *ChannelMSTeamsConfig) UnmarshalJSON(data []byte) error {
type alias ChannelMSTeamsConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, config.DefaultMSTeamsV2Config.VSendResolved)
c.Title.SetIfUnset(config.DefaultMSTeamsV2Config.Title)
c.Text.SetIfUnset(config.DefaultMSTeamsV2Config.Text)
return c.Validate()
}
func (c ChannelMSTeamsConfig) Validate() error {

View File

@@ -10,13 +10,27 @@ type ChannelOpsgenieConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIKey string `json:"apiKey" required:"true" format:"password"`
APIURL string `json:"apiUrl"`
Message valuer.UnsetOrNonEmptyString `json:"message"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Source valuer.UnsetOrNonEmptyString `json:"source"`
Details map[string]string `json:"details,omitempty"`
Message valuer.UnsetOrNonEmptyString `json:"message,omitzero"`
Description valuer.UnsetOrNonEmptyString `json:"description,omitzero"`
Source valuer.UnsetOrNonEmptyString `json:"source,omitzero"`
Details map[string]string `json:"details,omitzero"`
Priority string `json:"priority"`
}
func (c *ChannelOpsgenieConfig) UnmarshalJSON(data []byte) error {
type alias ChannelOpsgenieConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, config.DefaultOpsGenieConfig.VSendResolved)
c.Message.SetIfUnset(config.DefaultOpsGenieConfig.Message)
c.Description.SetIfUnset(config.DefaultOpsGenieConfig.Description)
c.Source.SetIfUnset(config.DefaultOpsGenieConfig.Source)
return c.Validate()
}
func (c ChannelOpsgenieConfig) Validate() error {
if c.APIKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for an opsgenie channel")

View File

@@ -10,15 +10,39 @@ type ChannelPagerdutyConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
RoutingKey string `json:"routingKey" required:"true" format:"password"`
URL string `json:"url"`
Source valuer.UnsetOrNonEmptyString `json:"source"`
Client valuer.UnsetOrNonEmptyString `json:"client"`
ClientURL valuer.UnsetOrNonEmptyString `json:"clientUrl"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Source valuer.UnsetOrNonEmptyString `json:"source,omitzero"`
Client valuer.UnsetOrNonEmptyString `json:"client,omitzero"`
ClientURL valuer.UnsetOrNonEmptyString `json:"clientUrl,omitzero"`
Description valuer.UnsetOrNonEmptyString `json:"description,omitzero"`
Severity string `json:"severity"`
Component string `json:"component"`
Group string `json:"group"`
Class string `json:"class"`
Details map[string]string `json:"details,omitempty"`
Details map[string]string `json:"details,omitzero"`
}
func (c *ChannelPagerdutyConfig) UnmarshalJSON(data []byte) error {
type alias ChannelPagerdutyConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, config.DefaultPagerdutyConfig.VSendResolved)
c.Description.SetIfUnset(config.DefaultPagerdutyConfig.Description)
c.Client.SetIfUnset(config.DefaultPagerdutyConfig.Client)
c.ClientURL.SetIfUnset(config.DefaultPagerdutyConfig.ClientURL)
c.Source.SetIfUnset(c.Client.StringValue())
if c.Details == nil {
c.Details = make(map[string]string, len(config.DefaultPagerdutyDetails))
for key, value := range config.DefaultPagerdutyDetails {
if template, ok := value.(string); ok {
c.Details[key] = template
}
}
}
return c.Validate()
}
func (c ChannelPagerdutyConfig) Validate() error {
@@ -90,3 +114,32 @@ func newChannelPagerdutyConfigFromReceiver(name string, receiver *Receiver) (Cha
Details: details,
}, nil
}
func newUpstreamDetails(details map[string]string) map[string]any {
if details == nil {
return nil
}
upstream := make(map[string]any, len(details))
for key, value := range details {
upstream[key] = value
}
return upstream
}
func extractStringDetails(name string, details map[string]any) (map[string]string, error) {
extracted := make(map[string]string, len(details))
for key, value := range details {
stringValue, ok := value.(string)
if !ok {
return nil, errors.NewInvalidInputf(
ErrCodeAlertmanagerChannelInvalid,
"channel %q sets a non-string value for details.%s, which is not supported", name, key,
)
}
extracted[key] = stringValue
}
return extracted, nil
}

View File

@@ -10,15 +10,15 @@ type ChannelSlackConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIURL string `json:"apiUrl" required:"true" format:"password"`
Channel string `json:"channel"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
Color valuer.UnsetOrNonEmptyString `json:"color"`
TitleLink valuer.UnsetOrNonEmptyString `json:"titleLink"`
Pretext valuer.UnsetOrNonEmptyString `json:"pretext"`
Fallback valuer.UnsetOrNonEmptyString `json:"fallback"`
Footer valuer.UnsetOrNonEmptyString `json:"footer"`
Fields []ChannelSlackField `json:"fields,omitempty"`
Actions []ChannelSlackAction `json:"actions,omitempty"`
Title valuer.UnsetOrNonEmptyString `json:"title,omitzero"`
Text valuer.UnsetOrNonEmptyString `json:"text,omitzero"`
Color valuer.UnsetOrNonEmptyString `json:"color,omitzero"`
TitleLink valuer.UnsetOrNonEmptyString `json:"titleLink,omitzero"`
Pretext valuer.UnsetOrNonEmptyString `json:"pretext,omitzero"`
Fallback valuer.UnsetOrNonEmptyString `json:"fallback,omitzero"`
Footer valuer.UnsetOrNonEmptyString `json:"footer,omitzero"`
Fields []ChannelSlackField `json:"fields,omitzero"`
Actions []ChannelSlackAction `json:"actions,omitzero"`
}
type ChannelSlackField struct {
@@ -46,6 +46,24 @@ type ChannelSlackConfirmation struct {
DismissText string `json:"dismissText"`
}
func (c *ChannelSlackConfig) UnmarshalJSON(data []byte) error {
type alias ChannelSlackConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, config.DefaultSlackConfig.VSendResolved)
c.Title.SetIfUnset(config.DefaultSlackConfig.Title)
c.Text.SetIfUnset(config.DefaultSlackConfig.Text)
c.Color.SetIfUnset(config.DefaultSlackConfig.Color)
c.TitleLink.SetIfUnset(config.DefaultSlackConfig.TitleLink)
c.Pretext.SetIfUnset(config.DefaultSlackConfig.Pretext)
c.Fallback.SetIfUnset(config.DefaultSlackConfig.Fallback)
c.Footer.SetIfUnset(config.DefaultSlackConfig.Footer)
return c.Validate()
}
func (c ChannelSlackConfig) Validate() error {
if c.APIURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiUrl is required for a slack channel")

View File

@@ -1,11 +1,16 @@
package alertmanagertypes
import (
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
// bearerAuthorizationType is the scheme SigNoz writes for token auth.
const bearerAuthorizationType = "Bearer"
// ChannelWebhookConfig splits apart the two authentication modes the legacy API
// overloaded onto one password field, where an empty username meant the password
// was really a bearer token. Username or Password may be set without the other,
@@ -18,6 +23,17 @@ type ChannelWebhookConfig struct {
BearerToken string `json:"bearerToken" format:"password"`
}
func (c *ChannelWebhookConfig) UnmarshalJSON(data []byte) error {
type alias ChannelWebhookConfig
if err := decodeStrict(data, (*alias)(c)); err != nil {
return err
}
fillSendResolved(&c.SendResolved, config.DefaultWebhookConfig.VSendResolved)
return c.Validate()
}
func (c ChannelWebhookConfig) Validate() error {
if c.URL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for a webhook channel")
@@ -96,3 +112,16 @@ func newChannelWebhookConfigFromReceiver(name string, receiver *Receiver) (Chann
return webhook, nil
}
func rejectHTTPAuthorizationBeyondBearer(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil || httpConfig.Authorization == nil {
return nil
}
authorization := httpConfig.Authorization
if !strings.EqualFold(authorization.Type, bearerAuthorizationType) || *authorization != (commoncfg.Authorization{Type: authorization.Type, Credentials: authorization.Credentials}) {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization with fields other than a bearer token, which is not supported", channelName)
}
return nil
}

View File

@@ -1,104 +1,50 @@
package alertmanagertypes
import (
"bytes"
"encoding/json"
"net/url"
"reflect"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
"github.com/swaggest/jsonschema-go"
)
var (
ErrCodeChannelUnsupportedKind = errors.MustNewCode("channel_unsupported_kind")
)
// ════════════════════════════════════════════════════════════════════════
// Kind
// ════════════════════════════════════════════════════════════════════════
// ChannelKind selects which ChannelSpec a channel carries and which notifier
// integration is built for it.
type ChannelKind struct {
valuer.String
}
var (
ChannelKindSlack = ChannelKind{valuer.NewString("slack")}
ChannelKindEmail = ChannelKind{valuer.NewString("email")}
ChannelKindWebhook = ChannelKind{valuer.NewString("webhook")}
ChannelKindPagerduty = ChannelKind{valuer.NewString("pagerduty")}
ChannelKindOpsgenie = ChannelKind{valuer.NewString("opsgenie")}
ChannelKindMSTeams = ChannelKind{valuer.NewString("msteams")}
ChannelKindGoogleChat = ChannelKind{valuer.NewString("googlechat")}
ChannelKindJira = ChannelKind{valuer.NewString("jira")}
ChannelKindJSMOps = ChannelKind{valuer.NewString("jsmops")}
ChannelKindIncidentIO = ChannelKind{valuer.NewString("incidentio")}
)
func (ChannelKind) Enum() []any {
kinds := make([]any, 0, len(channelKinds))
for _, channelKind := range channelKinds {
kinds = append(kinds, channelKind.kind)
}
return kinds
}
func (t ChannelKind) IsValid() bool {
return slices.ContainsFunc(t.Enum(), func(v any) bool { return v == t })
}
// ToStoredType returns the Channel.Type a channel of this kind is stored under,
// which matches the kind for all but msteams.
func (t ChannelKind) ToStoredType() string {
if t == ChannelKindMSTeams {
return "msteamsv2"
}
return t.StringValue()
}
// parseStoredChannelType inverts ToStoredType. It reports false for the notifier
// kinds v1 accepted but v2 does not model.
func parseStoredChannelType(stored string) (ChannelKind, bool) {
for _, channelKind := range channelKinds {
if channelKind.kind.ToStoredType() == stored {
return channelKind.kind, true
}
}
return ChannelKind{}, false
}
func ErrUnsupportedChannelKind(s string) error {
return errors.Newf(
errors.TypeInvalidInput,
ErrCodeChannelUnsupportedKind,
"unknown notification channel kind %q; allowed values: %s",
s, allowedValuesForChannelKind(),
)
}
// ════════════════════════════════════════════════════════════════════════
// Union
// ════════════════════════════════════════════════════════════════════════
// ChannelConfig is the discriminated union of per-kind configurations. The
// envelope sits on config rather than the resource root, so clients narrow on
// config.kind instead of every request and response flavor becoming a oneOf.
type ChannelConfig struct {
Kind ChannelKind `json:"kind" required:"true"`
Spec any `json:"spec" required:"true"`
}
func (c *ChannelConfig) UnmarshalJSON(data []byte) error {
var envelope struct {
Kind ChannelKind `json:"kind"`
// json.RawMessage keeps spec bytes unparsed until Kind is known. Once
// Kind is known, this Spec can be decoded into the correct type.
Spec json.RawMessage `json:"spec"`
}
if err := decodeStrict(data, &envelope); err != nil {
return err
}
if len(envelope.Spec) == 0 {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec is required")
}
spec, ok := buildEmptyChannelSpecForKind(envelope.Kind)
if !ok {
return ErrUnsupportedChannelKind(envelope.Kind.StringValue())
}
if err := json.Unmarshal(envelope.Spec, spec); err != nil {
return err
}
c.Kind = envelope.Kind
c.Spec = spec
return nil
}
func (c ChannelConfig) Validate() error {
newSpec, ok := newChannelSpec(c.Kind)
expected, ok := buildEmptyChannelSpecForKind(c.Kind)
if !ok {
return ErrUnsupportedChannelKind(c.Kind.StringValue())
}
@@ -107,44 +53,21 @@ func (c ChannelConfig) Validate() error {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec is required")
}
spec, ok := c.Spec.(ChannelSpec)
if !ok {
return errors.NewInternalf(errors.CodeInternal, "config.spec was not decoded into a known type")
spec, err := c.toChannelSpec()
if err != nil {
return err
}
// A decoded config cannot disagree, because UnmarshalJSON builds the spec
// from the kind. A caller assembling the struct can, and the conversion to a
// receiver dispatches on the spec, so a mismatch would silently outrank the
// declared kind.
if reflect.TypeOf(spec) != reflect.TypeOf(newSpec()) {
// Only a config assembled in Go gets here, and one can pair a spec with the
// wrong kind. The conversion to a receiver dispatches on the spec, so the
// mismatch would silently outrank the declared kind.
if reflect.TypeOf(spec) != reflect.TypeOf(expected) {
return errors.NewInternalf(errors.CodeInternal, "config.spec does not match kind %q", c.Kind.StringValue())
}
return spec.Validate()
}
func (c *ChannelConfig) UnmarshalJSON(data []byte) error {
channelKindString, specJSON, err := extractKindAndSpec(data)
if err != nil {
return err
}
factory, ok := newChannelSpec(ChannelKind{valuer.NewString(channelKindString)})
if !ok {
return ErrUnsupportedChannelKind(channelKindString)
}
spec, err := decodeChannelSpec(specJSON, factory(), channelKindString)
if err != nil {
return err
}
c.Kind = ChannelKind{valuer.NewString(channelKindString)}
c.Spec = *spec
return nil
}
// ChannelConfigVariant names one branch of the union. Each instantiation becomes
// its own OpenAPI component with kind pinned to the one value it accepts.
type ChannelConfigVariant[S any] struct {
@@ -193,311 +116,6 @@ func (ChannelConfig) PrepareJSONSchema(s *jsonschema.Schema) error {
})
}
// ════════════════════════════════════════════════════════════════════════
// Specs
// ════════════════════════════════════════════════════════════════════════
type ChannelSpec interface {
Validate() error
toUndefaultedReceiver(displayName string) (*Receiver, error)
}
// ════════════════════════════════════════════════════════════════════════
// Helpers
// ════════════════════════════════════════════════════════════════════════
// bearerAuthorizationType is the scheme SigNoz writes for token auth.
const bearerAuthorizationType = "Bearer"
// parseSecretURL and parseUpstreamURL wrap the two URL types upstream uses for
// notifier endpoints. Callers holding an optional URL skip the call on an empty
// string, so the field stays nil and is omitted rather than stored as an empty URL.
func parseSecretURL(raw string) (*config.SecretURL, error) {
parsed, err := parseUpstreamURL(raw)
if err != nil {
return nil, err
}
return (*config.SecretURL)(parsed), nil
}
func parseUpstreamURL(raw string) (*config.URL, error) {
parsed, err := url.Parse(raw)
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "parse url %q", raw)
}
return &config.URL{URL: parsed}, nil
}
func formatSecretURL(secretURL *config.SecretURL) string {
if secretURL == nil {
return ""
}
return formatUpstreamURL((*config.URL)(secretURL))
}
func formatUpstreamURL(upstreamURL *config.URL) string {
if upstreamURL == nil || upstreamURL.URL == nil {
return ""
}
return upstreamURL.String()
}
// PagerDuty is the one notifier whose details upstream types as map[string]any.
func newUpstreamDetails(details map[string]string) map[string]any {
if details == nil {
return nil
}
upstream := make(map[string]any, len(details))
for key, value := range details {
upstream[key] = value
}
return upstream
}
func extractStringDetails(name string, details map[string]any) (map[string]string, error) {
extracted := make(map[string]string, len(details))
for key, value := range details {
stringValue, ok := value.(string)
if !ok {
return nil, errors.NewInvalidInputf(
ErrCodeAlertmanagerChannelInvalid,
"channel %q sets a non-string value for details.%s, which is not supported", name, key,
)
}
extracted[key] = stringValue
}
return extracted, nil
}
func rejectAnyHTTPAuth(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil {
return nil
}
if httpConfig.BasicAuth != nil {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.basic_auth, which is not supported", channelName)
}
if httpConfig.Authorization != nil {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization, which is not supported", channelName)
}
return rejectUnsupportedHTTPConfig(channelName, httpConfig)
}
func rejectUnsupportedHTTPConfig(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil {
return nil
}
for _, field := range []struct {
fieldName string
isFieldConfigured bool
}{
{"oauth2", httpConfig.OAuth2 != nil},
{"bearer_token", httpConfig.BearerToken != ""},
{"bearer_token_file", httpConfig.BearerTokenFile != ""},
{"proxy_url", httpConfig.ProxyURL.URL != nil && httpConfig.ProxyURL.String() != ""},
{"no_proxy", httpConfig.NoProxy != ""},
{"proxy_from_environment", httpConfig.ProxyFromEnvironment},
{"http_headers", httpConfig.HTTPHeaders != nil},
{"tls_config", httpConfig.TLSConfig != (commoncfg.TLSConfig{})},
{"follow_redirects", !httpConfig.FollowRedirects},
{"enable_http2", !httpConfig.EnableHTTP2},
} {
if field.isFieldConfigured {
return errors.NewInvalidInputf(
ErrCodeAlertmanagerChannelInvalid,
"channel %q sets http_config.%s, which is not supported", channelName, field.fieldName,
)
}
}
return nil
}
func rejectHTTPBasicAuthBeyondPassword(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil || httpConfig.BasicAuth == nil {
return nil
}
basicAuth := httpConfig.BasicAuth
if *basicAuth != (commoncfg.BasicAuth{Username: basicAuth.Username, Password: basicAuth.Password}) {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.basic_auth with fields other than username and password, which is not supported", channelName)
}
return nil
}
func rejectHTTPAuthorizationBeyondBearer(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil || httpConfig.Authorization == nil {
return nil
}
authorization := httpConfig.Authorization
if !strings.EqualFold(authorization.Type, bearerAuthorizationType) || *authorization != (commoncfg.Authorization{Type: authorization.Type, Credentials: authorization.Credentials}) {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization with fields other than a bearer token, which is not supported", channelName)
}
return nil
}
// channelKinds registers each notification kind with the spec constructor
// UnmarshalJSON picks by kind and the extractor that reads a stored receiver
// back. The ChannelKind enum derives from it; the JSON schema hooks stay
// literal lists so each branch reads as one line.
var channelKinds = []channelKindEntry{
{
kind: ChannelKindSlack,
newSpec: func() ChannelSpec { return new(ChannelSlackConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.SlackConfigs) },
extractSpec: newChannelSlackConfigFromReceiver,
},
{
kind: ChannelKindEmail,
newSpec: func() ChannelSpec { return new(ChannelEmailConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.EmailConfigs) },
extractSpec: newChannelEmailConfigFromReceiver,
},
{
kind: ChannelKindWebhook,
newSpec: func() ChannelSpec { return new(ChannelWebhookConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.WebhookConfigs) },
extractSpec: newChannelWebhookConfigFromReceiver,
},
{
kind: ChannelKindPagerduty,
newSpec: func() ChannelSpec { return new(ChannelPagerdutyConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.PagerdutyConfigs) },
extractSpec: newChannelPagerdutyConfigFromReceiver,
},
{
kind: ChannelKindOpsgenie,
newSpec: func() ChannelSpec { return new(ChannelOpsgenieConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.OpsGenieConfigs) },
extractSpec: newChannelOpsgenieConfigFromReceiver,
},
{
kind: ChannelKindMSTeams,
newSpec: func() ChannelSpec { return new(ChannelMSTeamsConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.MSTeamsV2Configs) },
extractSpec: newChannelMSTeamsConfigFromReceiver,
},
{
kind: ChannelKindGoogleChat,
newSpec: func() ChannelSpec { return new(ChannelGoogleChatConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.GoogleChatConfigs) },
extractSpec: newChannelGoogleChatConfigFromReceiver,
},
{
kind: ChannelKindJira,
newSpec: func() ChannelSpec { return new(ChannelJiraConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.JiraConfigs) },
extractSpec: newChannelJiraConfigFromReceiver,
},
{
kind: ChannelKindJSMOps,
newSpec: func() ChannelSpec { return new(ChannelJSMOpsConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.JSMOpsConfigs) },
extractSpec: newChannelJSMOpsConfigFromReceiver,
},
{
kind: ChannelKindIncidentIO,
newSpec: func() ChannelSpec { return new(ChannelIncidentIOConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.IncidentIOConfigs) },
extractSpec: newChannelIncidentIOConfigFromReceiver,
},
}
type channelKindEntry struct {
kind ChannelKind
newSpec func() ChannelSpec
// countConfigs guards extractSpec, which reads the receiver's first config of
// this kind and so must not be called when there is none.
countConfigs func(receiver *Receiver) int
extractSpec func(name string, receiver *Receiver) (ChannelSpec, error)
}
func newChannelSpec(kind ChannelKind) (func() ChannelSpec, bool) {
for _, channelKind := range channelKinds {
if channelKind.kind == kind {
return channelKind.newSpec, true
}
}
return nil, false
}
// resolveSendResolved falls back to the notifier's own upstream default, because
// send_resolved has no omitempty: a zero value would marshal as an explicit false
// and overwrite the default rather than leave it in place.
func resolveSendResolved(sendResolved *bool, upstreamDefault bool) bool {
if sendResolved == nil {
return upstreamDefault
}
return *sendResolved
}
func allowedValuesForChannelKind() string {
return formatAllowedValues((ChannelKind{}).Enum())
}
func formatAllowedValues(enum []any) string {
values := make([]string, 0, len(enum))
for _, value := range enum {
stringValuer, ok := value.(interface{ StringValue() string })
if !ok {
continue
}
values = append(values, "`"+stringValuer.StringValue()+"`")
}
slices.Sort(values)
return strings.Join(values, ", ")
}
// extractKindAndSpec parses a {"kind": "...", "spec": {...}} envelope. Unknown
// keys are rejected here rather than by the caller's decoder: a custom
// UnmarshalJSON receives raw bytes, so DisallowUnknownFields on the request body
// does not reach inside config.
func extractKindAndSpec(data []byte) (string, []byte, error) {
var head struct {
Kind string `json:"kind"`
Spec json.RawMessage `json:"spec"`
}
dec := json.NewDecoder(bytes.NewReader(data))
dec.DisallowUnknownFields()
if err := dec.Decode(&head); err != nil {
return "", nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "invalid channel config envelope")
}
return head.Kind, head.Spec, nil
}
// decodeChannelSpec rejects unknown fields so a spec meant for another kind is an
// error rather than a silently empty struct, and validates before returning.
func decodeChannelSpec[T ChannelSpec](specJSON []byte, target T, channelType string) (*T, error) {
if len(specJSON) == 0 {
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "type %q: spec is required", channelType)
}
dec := json.NewDecoder(bytes.NewReader(specJSON))
dec.DisallowUnknownFields()
if err := dec.Decode(target); err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "type %q: invalid spec JSON", channelType)
}
if err := target.Validate(); err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "type %q: %s", channelType, err.Error())
}
return &target, nil
}
// signozDiscriminatorKey is the extension key that signoz.attachDiscriminators
// promotes into a native OpenAPI 3 discriminator after reflection.
const signozDiscriminatorKey = "x-signoz-discriminator"
@@ -514,6 +132,18 @@ func channelVariantRef(spec string) string {
return schemaRef("AlertmanagertypesChannelConfigVariantGithubComSigNozSignozPkgTypesAlertmanagertypes" + spec)
}
// toChannelSpec asserts what UnmarshalJSON decoded. Spec is any rather than
// ChannelSpec because the OpenAPI reflector turns an interface field into an
// empty component.
func (c ChannelConfig) toChannelSpec() (ChannelSpec, error) {
spec, ok := c.Spec.(ChannelSpec)
if !ok {
return nil, errors.NewInternalf(errors.CodeInternal, "config.spec was not decoded into a known type")
}
return spec, nil
}
// markDiscriminator tags a oneOf schema with x-signoz-discriminator, keyed on
// propertyName with the given value -> schema-ref mapping, so generated clients
// get a discriminated DTO instead of an intersection.

View File

@@ -3,8 +3,11 @@ package alertmanagertypes
import (
"encoding/json"
"reflect"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
)
@@ -12,13 +15,47 @@ import (
// API -> storage
// ════════════════════════════════════════════════════════════════════════
// ToChannel returns the receiver alongside because the alertmanager config is
// updated from it, not from the channel.
func (p *PostableNotificationChannel) ToChannel(orgID string) (*Channel, *Receiver, error) {
receiver, err := p.ToReceiver()
if err != nil {
return nil, nil, err
}
data, err := json.Marshal(receiver)
if err != nil {
return nil, nil, errors.WrapInternalf(err, errors.CodeInternal, "marshal receiver")
}
spec, err := p.Config.toChannelSpec()
if err != nil {
return nil, nil, err
}
channel := &Channel{
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
TimeAuditable: types.TimeAuditable{CreatedAt: time.Now(), UpdatedAt: time.Now()},
Name: p.Name,
DisplayName: p.DisplayName,
Type: p.Config.Kind.StringValue(),
Data: string(data),
OrgID: orgID,
}
if err := channel.fillSpec(spec); err != nil {
return nil, nil, err
}
return channel, receiver, nil
}
// ToReceiver hands the assembled receiver to newDefaultedReceiver, which is the
// only place upstream applies a notifier's defaults and validation — several
// integrations panic without them.
func (p *PostableNotificationChannel) ToReceiver() (*Receiver, error) {
spec, ok := p.Config.Spec.(ChannelSpec)
if !ok {
return nil, errors.NewInternalf(errors.CodeInternal, "config.spec was not decoded into a known type")
spec, err := p.Config.toChannelSpec()
if err != nil {
return nil, err
}
receiver, err := spec.toUndefaultedReceiver(p.DisplayName)
@@ -45,22 +82,63 @@ func (t *TestableNotificationChannel) ToReceiver() (*Receiver, error) {
return postable.ToReceiver()
}
func (c *Channel) UpdateFromUpdatable(updatable UpdatableNotificationChannel) (*Receiver, error) {
receiver, err := updatable.ToReceiver(c.DisplayName)
if err != nil {
return nil, err
}
data, err := json.Marshal(receiver)
if err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "marshal receiver")
}
spec, err := updatable.Config.toChannelSpec()
if err != nil {
return nil, err
}
if err := c.fillSpec(spec); err != nil {
return nil, err
}
c.Type = updatable.Config.Kind.StringValue()
c.Data = string(data)
c.UpdatedAt = time.Now()
return receiver, nil
}
// ════════════════════════════════════════════════════════════════════════
// Storage -> API
// ════════════════════════════════════════════════════════════════════════
// toPostableNotificationChannel derives the kind from the config the receiver
// actually carries rather than from Channel.Type, so a row written with several
// notifier kinds is rejected instead of reported under whichever one
// receiverChannelType happened to pick.
func (c *Channel) toPostableNotificationChannel() (*PostableNotificationChannel, error) {
// toChannelConfig pairs the scanned spec with the kind Type names. Only a row
// the migration could not backfill has none, so deriving it here reports why.
func (c *Channel) toChannelConfig() (ChannelConfig, error) {
if c.Spec == nil {
return c.deriveChannelConfig()
}
channelKind, ok := parseChannelKind(c.Type)
if !ok {
return ChannelConfig{}, errors.NewInternalf(errors.CodeInternal, "channel %q carries a spec under unmodelled type %q", c.DisplayName, c.Type)
}
return ChannelConfig{Kind: channelKind, Spec: c.Spec}, nil
}
// deriveChannelConfig derives the kind from the config the receiver actually
// carries rather than from Channel.Type, so a row written with several notifier
// kinds is rejected instead of reported under whichever one receiverChannelType
// happened to pick.
func (c *Channel) deriveChannelConfig() (ChannelConfig, error) {
receiver := &Receiver{Receiver: &config.Receiver{}}
if err := json.Unmarshal([]byte(c.Data), receiver); err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "unmarshal channel %q", c.DisplayName)
return ChannelConfig{}, errors.WrapInternalf(err, errors.CodeInternal, "unmarshal channel %q", c.DisplayName)
}
if total := countNotifierConfigs(receiver); total > 1 {
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q carries %d notifier configurations; only one per channel is supported", c.DisplayName, total)
return ChannelConfig{}, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q carries %d notifier configurations; only one per channel is supported", c.DisplayName, total)
}
for _, channelKind := range channelKinds {
@@ -70,17 +148,21 @@ func (c *Channel) toPostableNotificationChannel() (*PostableNotificationChannel,
spec, err := channelKind.extractSpec(c.DisplayName, receiver)
if err != nil {
return nil, err
return ChannelConfig{}, err
}
return &PostableNotificationChannel{
Name: c.Name,
DisplayName: c.DisplayName,
Config: ChannelConfig{Kind: channelKind.kind, Spec: spec},
}, nil
// The derived config is stored and decoded back through the same
// validation a request goes through, so one that would not decode is
// unrepresentable rather than stored.
channelConfig := ChannelConfig{Kind: channelKind.kind, Spec: spec}
if err := channelConfig.Validate(); err != nil {
return ChannelConfig{}, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "channel %q: %s", c.DisplayName, err.Error())
}
return channelConfig, nil
}
return nil, errors.NewInvalidInputf(ErrCodeChannelUnsupportedKind, "channel %q carries no supported notifier configuration", c.DisplayName)
return ChannelConfig{}, errors.NewInvalidInputf(ErrCodeChannelUnsupportedKind, "channel %q carries no supported notifier configuration", c.DisplayName)
}
// countNotifierConfigs totals every *_configs entry on the receiver, including
@@ -111,15 +193,15 @@ func countConfigsFields(v reflect.Value) int {
}
func (c *Channel) ToGettableNotificationChannel() (*GettableNotificationChannel, error) {
postable, err := c.toPostableNotificationChannel()
channelConfig, err := c.toChannelConfig()
if err != nil {
return nil, err
}
return &GettableNotificationChannel{
Name: postable.Name,
DisplayName: postable.DisplayName,
Config: postable.Config,
Name: c.Name,
DisplayName: c.DisplayName,
Config: channelConfig,
ID: c.ID,
CreatedAt: c.CreatedAt,
UpdatedAt: c.UpdatedAt,
@@ -130,7 +212,7 @@ func (c *Channel) ToGettableNotificationChannel() (*GettableNotificationChannel,
// no ChannelKind models, which v1 allowed because it accepted every upstream
// notifier kind. One such row must not fail the whole page.
func (c *Channel) ToListedNotificationChannel() *ListedNotificationChannel {
channelKind, _ := parseStoredChannelType(c.Type)
channelKind, _ := parseChannelKind(c.Type)
return &ListedNotificationChannel{
ID: c.ID,

View File

@@ -18,7 +18,7 @@ import (
// field fails rather than going unasserted. Webhook is covered by
// TestPostableChannelToReceiverRoundTripsWebhookAuthModes, whose auth modes are
// mutually exclusive and so cannot all be set at once.
func TestChannelToPostableChannelRoundTripsEveryFieldOfEveryKind(t *testing.T) {
func TestDeriveChannelConfigRoundTripsEveryFieldOfEveryKind(t *testing.T) {
sendResolved := true
short := true
@@ -269,18 +269,14 @@ func TestChannelToPostableChannelRoundTripsEveryFieldOfEveryKind(t *testing.T) {
}
require.NoError(t, postable.Validate())
receiver, err := postable.ToReceiver()
channel, _, err := postable.ToChannel("org-1")
require.NoError(t, err)
channel, err := NewChannelFromReceiverWithName(receiver, postable.Name, "org-1")
derived, err := channel.deriveChannelConfig()
require.NoError(t, err)
roundTripped, err := channel.toPostableNotificationChannel()
require.NoError(t, err)
assert.Equal(t, postable.Name, roundTripped.Name)
assert.Equal(t, testCase.kind, roundTripped.Config.Kind)
assert.Equal(t, testCase.expectedRoundTrip, roundTripped.Config.Spec)
assert.Equal(t, testCase.kind, derived.Kind)
assert.Equal(t, testCase.expectedRoundTrip, derived.Spec)
})
}
}
@@ -298,10 +294,7 @@ func TestPostableChannelToReceiverOmitsEmailTransportCredentials(t *testing.T) {
},
}
receiver, err := postable.ToReceiver()
require.NoError(t, err)
channel, err := NewChannelFromReceiverWithName(receiver, postable.Name, "org-1")
channel, _, err := postable.ToChannel("org-1")
require.NoError(t, err)
for _, credentialKey := range []string{"auth_username", "auth_password", "auth_secret", "tls_config"} {
@@ -353,16 +346,13 @@ func TestPostableChannelToReceiverRoundTripsWebhookAuthModes(t *testing.T) {
}
require.NoError(t, postable.Validate())
receiver, err := postable.ToReceiver()
require.NoError(t, err)
channel, err := NewChannelFromReceiverWithName(receiver, postable.Name, "org-1")
channel, _, err := postable.ToChannel("org-1")
require.NoError(t, err)
assert.Contains(t, channel.Data, testCase.expectedInData)
roundTripped, err := channel.toPostableNotificationChannel()
derived, err := channel.deriveChannelConfig()
require.NoError(t, err)
assert.Equal(t, testCase.expectedRoundTrip, roundTripped.Config.Spec)
assert.Equal(t, testCase.expectedRoundTrip, derived.Spec)
})
}
}
@@ -375,7 +365,7 @@ func TestRejectUnrepresentableHTTPConfigCoversEveryUpstreamMember(t *testing.T)
assert.Equal(t, 5, reflect.TypeFor[commoncfg.ProxyConfig]().NumField())
}
func TestChannelToPostableChannelRejectsUnrepresentableChannels(t *testing.T) {
func TestDeriveChannelConfigRejectsUnrepresentableChannels(t *testing.T) {
testCases := []struct {
description string
channel Channel
@@ -537,7 +527,7 @@ func TestChannelToPostableChannelRejectsUnrepresentableChannels(t *testing.T) {
for _, testCase := range testCases {
t.Run(testCase.description, func(t *testing.T) {
_, err := testCase.channel.toPostableNotificationChannel()
_, err := testCase.channel.deriveChannelConfig()
assert.Error(t, err)
})
}
@@ -545,7 +535,7 @@ func TestChannelToPostableChannelRejectsUnrepresentableChannels(t *testing.T) {
// The HTTP auth scheme is case-insensitive (RFC 7235) and Alertmanager sends
// the stored spelling verbatim, so a hand-written receiver may carry any casing.
func TestChannelToPostableChannelReadsWebhookBearerSchemeCaseInsensitively(t *testing.T) {
func TestDeriveChannelConfigReadsWebhookBearerSchemeCaseInsensitively(t *testing.T) {
sendResolved := config.DefaultWebhookConfig.VSendResolved
testCases := []struct {
@@ -574,10 +564,10 @@ func TestChannelToPostableChannelReadsWebhookBearerSchemeCaseInsensitively(t *te
t.Run(testCase.name, func(t *testing.T) {
channel := Channel{DisplayName: "hook", Data: testCase.storedChannelData}
postable, err := channel.toPostableNotificationChannel()
derived, err := channel.deriveChannelConfig()
require.NoError(t, err)
assert.Equal(t, ChannelKindWebhook, postable.Config.Kind)
assert.Equal(t, testCase.expectedWebhookSpec, postable.Config.Spec)
assert.Equal(t, ChannelKindWebhook, derived.Kind)
assert.Equal(t, testCase.expectedWebhookSpec, derived.Spec)
})
}
}

View File

@@ -0,0 +1,78 @@
package alertmanagertypes
import (
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
)
var (
ErrCodeChannelUnsupportedKind = errors.MustNewCode("channel_unsupported_kind")
)
// ChannelKind selects which ChannelSpec a channel carries and which notifier
// integration is built for it.
type ChannelKind struct {
valuer.String
}
var (
ChannelKindSlack = ChannelKind{valuer.NewString("slack")}
ChannelKindEmail = ChannelKind{valuer.NewString("email")}
ChannelKindWebhook = ChannelKind{valuer.NewString("webhook")}
ChannelKindPagerduty = ChannelKind{valuer.NewString("pagerduty")}
ChannelKindOpsgenie = ChannelKind{valuer.NewString("opsgenie")}
ChannelKindMSTeams = ChannelKind{valuer.NewString("msteams")}
ChannelKindGoogleChat = ChannelKind{valuer.NewString("googlechat")}
ChannelKindJira = ChannelKind{valuer.NewString("jira")}
ChannelKindJSMOps = ChannelKind{valuer.NewString("jsmops")}
ChannelKindIncidentIO = ChannelKind{valuer.NewString("incidentio")}
)
func (ChannelKind) Enum() []any {
kinds := make([]any, 0, len(channelKinds))
for _, channelKind := range channelKinds {
kinds = append(kinds, channelKind.kind)
}
return kinds
}
func (t ChannelKind) IsValid() bool {
return slices.ContainsFunc(t.Enum(), func(v any) bool { return v == t })
}
func ErrUnsupportedChannelKind(s string) error {
return errors.Newf(errors.TypeInvalidInput, ErrCodeChannelUnsupportedKind, "unknown notification channel kind %q; allowed values: %s", s, allowedValuesForChannelKind())
}
// parseChannelKind reads a stored type. It reports false for the notifier
// kinds v1 accepted but v2 does not model.
func parseChannelKind(storedType string) (ChannelKind, bool) {
channelKind := ChannelKind{valuer.NewString(storedType)}
if !channelKind.IsValid() {
return ChannelKind{}, false
}
return channelKind, true
}
func allowedValuesForChannelKind() string {
return formatAllowedValues((ChannelKind{}).Enum())
}
func formatAllowedValues(enum []any) string {
values := make([]string, 0, len(enum))
for _, value := range enum {
stringValuer, ok := value.(interface{ StringValue() string })
if !ok {
continue
}
values = append(values, "`"+stringValuer.StringValue()+"`")
}
slices.Sort(values)
return strings.Join(values, ", ")
}

View File

@@ -0,0 +1,86 @@
package alertmanagertypes
type channelKindEntry struct {
kind ChannelKind
newEmptySpec func() ChannelSpec
// countConfigs guards extractSpec, which reads the receiver's first config of
// this kind and so must not be called when there is none.
countConfigs func(receiver *Receiver) int
extractSpec func(name string, receiver *Receiver) (ChannelSpec, error)
}
// channelKinds registers each notification kind with the spec constructor
// UnmarshalJSON picks by kind and the extractor that reads a stored receiver
// back. The ChannelKind enum derives from it; the JSON schema hooks stay
// literal lists so each branch reads as one line.
var channelKinds = []channelKindEntry{
{
kind: ChannelKindSlack,
newEmptySpec: func() ChannelSpec { return new(ChannelSlackConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.SlackConfigs) },
extractSpec: newChannelSlackConfigFromReceiver,
},
{
kind: ChannelKindEmail,
newEmptySpec: func() ChannelSpec { return new(ChannelEmailConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.EmailConfigs) },
extractSpec: newChannelEmailConfigFromReceiver,
},
{
kind: ChannelKindWebhook,
newEmptySpec: func() ChannelSpec { return new(ChannelWebhookConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.WebhookConfigs) },
extractSpec: newChannelWebhookConfigFromReceiver,
},
{
kind: ChannelKindPagerduty,
newEmptySpec: func() ChannelSpec { return new(ChannelPagerdutyConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.PagerdutyConfigs) },
extractSpec: newChannelPagerdutyConfigFromReceiver,
},
{
kind: ChannelKindOpsgenie,
newEmptySpec: func() ChannelSpec { return new(ChannelOpsgenieConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.OpsGenieConfigs) },
extractSpec: newChannelOpsgenieConfigFromReceiver,
},
{
kind: ChannelKindMSTeams,
newEmptySpec: func() ChannelSpec { return new(ChannelMSTeamsConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.MSTeamsV2Configs) },
extractSpec: newChannelMSTeamsConfigFromReceiver,
},
{
kind: ChannelKindGoogleChat,
newEmptySpec: func() ChannelSpec { return new(ChannelGoogleChatConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.GoogleChatConfigs) },
extractSpec: newChannelGoogleChatConfigFromReceiver,
},
{
kind: ChannelKindJira,
newEmptySpec: func() ChannelSpec { return new(ChannelJiraConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.JiraConfigs) },
extractSpec: newChannelJiraConfigFromReceiver,
},
{
kind: ChannelKindJSMOps,
newEmptySpec: func() ChannelSpec { return new(ChannelJSMOpsConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.JSMOpsConfigs) },
extractSpec: newChannelJSMOpsConfigFromReceiver,
},
{
kind: ChannelKindIncidentIO,
newEmptySpec: func() ChannelSpec { return new(ChannelIncidentIOConfig) },
countConfigs: func(receiver *Receiver) int { return len(receiver.IncidentIOConfigs) },
extractSpec: newChannelIncidentIOConfigFromReceiver,
},
}
func buildEmptyChannelSpecForKind(kind ChannelKind) (ChannelSpec, bool) {
for _, channelKind := range channelKinds {
if channelKind.kind == kind {
return channelKind.newEmptySpec(), true
}
}
return nil, false
}

View File

@@ -53,23 +53,3 @@ func TestChannelToListedChannelLeavesUnmodelledKindsEmpty(t *testing.T) {
assert.Equal(t, "tg", listed.Name)
assert.True(t, listed.Kind.IsZero())
}
// msteams is the only kind whose stored Channel.Type differs from the api kind,
// so ToStoredType has to agree with the type the write path derives.
func TestChannelKindMSTeamsIsStoredAsMSTeamsV2(t *testing.T) {
assert.Equal(t, "msteamsv2", ChannelKindMSTeams.ToStoredType())
postable := PostableNotificationChannel{
Name: "channel",
DisplayName: "channel",
Config: ChannelConfig{Kind: ChannelKindMSTeams, Spec: &ChannelMSTeamsConfig{WebhookURL: "https://a"}},
}
receiver, err := postable.ToReceiver()
require.NoError(t, err)
channel, err := NewChannelFromReceiverWithName(receiver, postable.Name, "org-1")
require.NoError(t, err)
assert.Equal(t, "msteamsv2", channel.Type)
}

View File

@@ -93,7 +93,7 @@ func (c *Channel) Diagnose() *ChannelRepair {
return repair
}
if _, err := c.toPostableNotificationChannel(); err != nil {
if _, err := c.toChannelConfig(); err != nil {
repair.Defect, repair.Detail = ChannelDefectUnrepresentable, err.Error()
return repair
}
@@ -117,7 +117,10 @@ func (c *Channel) Retype() error {
// SplitByNotifier turns a receiver carrying several notifier configurations into
// one channel per configuration. The first keeps this channel's identity so
// references to it stay valid; the rest are new channels numbered after it.
// references to it stay valid; the rest are new channels numbered after it. A
// part of a kind v2 does not model is kept without a spec, as the migration
// leaves such a row, so v1 still delivers through it until its own repair
// deletes it.
func (c *Channel) SplitByNotifier() ([]*Channel, error) {
receiver, err := NewReceiver(c.Data)
if err != nil {
@@ -129,21 +132,28 @@ func (c *Channel) SplitByNotifier() ([]*Channel, error) {
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q carries %d notifier configuration; nothing to split", c.DisplayName, len(singles))
}
if err := c.Update(singles[0]); err != nil {
return nil, err
}
parts := make([]*Channel, 0, len(singles))
for i, single := range singles {
if i > 0 {
single.Name = fmt.Sprintf("%s (%d)", c.DisplayName, i+1)
}
channels := []*Channel{c}
for i, single := range singles[1:] {
single.Name = fmt.Sprintf("%s (%d)", c.DisplayName, i+2)
channel, err := NewChannelFromReceiver(single, c.OrgID)
var part *Channel
if hasModelledNotifier(single) {
part, err = NewChannelFromReceiver(single, c.OrgID)
} else {
part, err = newChannelWithoutSpec(single, c.OrgID)
}
if err != nil {
return nil, err
}
channels = append(channels, channel)
parts = append(parts, part)
}
return channels, nil
c.Type, c.Data, c.Spec, c.StoredSpec = parts[0].Type, parts[0].Data, parts[0].Spec, parts[0].StoredSpec
c.UpdatedAt = time.Now()
return append([]*Channel{c}, parts[1:]...), nil
}
func hasModelledNotifier(receiver *Receiver) bool {

View File

@@ -0,0 +1,131 @@
package alertmanagertypes
import (
"net/url"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
type ChannelSpec interface {
Validate() error
toUndefaultedReceiver(displayName string) (*Receiver, error)
}
// parseSecretURL and parseUpstreamURL wrap the two URL types upstream uses for
// notifier endpoints. Callers holding an optional URL skip the call on an empty
// string, so the field stays nil and is omitted rather than stored as an empty URL.
func parseSecretURL(raw string) (*config.SecretURL, error) {
parsed, err := parseUpstreamURL(raw)
if err != nil {
return nil, err
}
return (*config.SecretURL)(parsed), nil
}
func parseUpstreamURL(raw string) (*config.URL, error) {
parsed, err := url.Parse(raw)
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "parse url %q", raw)
}
return &config.URL{URL: parsed}, nil
}
func formatSecretURL(secretURL *config.SecretURL) string {
if secretURL == nil {
return ""
}
return formatUpstreamURL((*config.URL)(secretURL))
}
func formatUpstreamURL(upstreamURL *config.URL) string {
if upstreamURL == nil || upstreamURL.URL == nil {
return ""
}
return upstreamURL.String()
}
func rejectAnyHTTPAuth(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil {
return nil
}
if httpConfig.BasicAuth != nil {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.basic_auth, which is not supported", channelName)
}
if httpConfig.Authorization != nil {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization, which is not supported", channelName)
}
return rejectUnsupportedHTTPConfig(channelName, httpConfig)
}
func rejectUnsupportedHTTPConfig(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil {
return nil
}
for _, field := range []struct {
fieldName string
isFieldConfigured bool
}{
{"oauth2", httpConfig.OAuth2 != nil},
{"bearer_token", httpConfig.BearerToken != ""},
{"bearer_token_file", httpConfig.BearerTokenFile != ""},
{"proxy_url", httpConfig.ProxyURL.URL != nil && httpConfig.ProxyURL.String() != ""},
{"no_proxy", httpConfig.NoProxy != ""},
{"proxy_from_environment", httpConfig.ProxyFromEnvironment},
{"http_headers", httpConfig.HTTPHeaders != nil},
{"tls_config", httpConfig.TLSConfig != (commoncfg.TLSConfig{})},
{"follow_redirects", !httpConfig.FollowRedirects},
{"enable_http2", !httpConfig.EnableHTTP2},
} {
if field.isFieldConfigured {
return errors.NewInvalidInputf(
ErrCodeAlertmanagerChannelInvalid,
"channel %q sets http_config.%s, which is not supported", channelName, field.fieldName,
)
}
}
return nil
}
func rejectHTTPBasicAuthBeyondPassword(channelName string, httpConfig *commoncfg.HTTPClientConfig) error {
if httpConfig == nil || httpConfig.BasicAuth == nil {
return nil
}
basicAuth := httpConfig.BasicAuth
if *basicAuth != (commoncfg.BasicAuth{Username: basicAuth.Username, Password: basicAuth.Password}) {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.basic_auth with fields other than username and password, which is not supported", channelName)
}
return nil
}
// fillSendResolved gives an omitted field its notifier default, so what is
// stored and read back is what takes effect.
func fillSendResolved(sendResolved **bool, upstreamDefault bool) {
if *sendResolved == nil {
value := upstreamDefault
*sendResolved = &value
}
}
// resolveSendResolved covers a spec assembled in code rather than decoded, whose
// defaults were never filled. send_resolved has no omitempty, so a zero value
// would marshal as an explicit false and overwrite the default rather than leave it.
func resolveSendResolved(sendResolved *bool, upstreamDefault bool) bool {
if sendResolved == nil {
return upstreamDefault
}
return *sendResolved
}

View File

@@ -0,0 +1,181 @@
package alertmanagertypes
import (
"encoding/json"
"reflect"
"strings"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
"github.com/swaggest/jsonschema-go"
)
type PostableChannel struct {
Receiver
}
func (PostableChannel) JSONSchema() (jsonschema.Schema, error) {
type alias PostableChannel
reflector := &jsonschema.Reflector{}
schema, err := reflector.Reflect(alias{}, jsonschema.DefinitionsPrefix("#/components/schemas/"))
if err != nil {
return jsonschema.Schema{}, err
}
schema.WithRequired("name")
var oneOf []jsonschema.SchemaOrBool
seen := map[string]struct{}{}
// Walk both halves: native fields on Receiver, upstream on the embed. A native
// field can shadow an upstream one with the same tag (e.g. jira_configs), so
// dedupe to avoid emitting two identical oneOf branches.
collect := func(t reflect.Type) {
for i := 0; i < t.NumField(); i++ {
jsonTag := strings.Split(t.Field(i).Tag.Get("json"), ",")[0]
if !strings.HasSuffix(jsonTag, "_configs") {
continue
}
if _, ok := seen[jsonTag]; ok {
continue
}
seen[jsonTag] = struct{}{}
branch := (&jsonschema.Schema{}).WithRequired(jsonTag)
oneOf = append(oneOf, branch.ToSchemaOrBool())
}
}
collect(reflect.TypeOf(Receiver{}))
collect(reflect.TypeOf(config.Receiver{}))
schema.WithOneOf(oneOf...)
return schema, nil
}
// NewChannelFromReceiver builds the channel a v1 write carries. The receiver is
// all there is, so the name is generated from its display name and the type and
// spec derived from it.
func NewChannelFromReceiver(receiver *Receiver, orgID string) (*Channel, error) {
channel, err := newChannelWithoutSpec(receiver, orgID)
if err != nil {
return nil, err
}
// A receiver v2 cannot represent is refused rather than stored, so every
// row written from here on reads through v2.
channelConfig, err := channel.deriveChannelConfig()
if err != nil {
return nil, err
}
spec, err := channelConfig.toChannelSpec()
if err != nil {
return nil, err
}
if err := channel.fillSpec(spec); err != nil {
return nil, err
}
return channel, nil
}
func (c *Channel) Update(receiver *Receiver) error {
channel, err := NewChannelFromReceiver(receiver, c.OrgID)
if err != nil {
return err
}
if c.DisplayName != channel.DisplayName {
return errors.Newf(errors.TypeInvalidInput, ErrCodeAlertmanagerChannelNameMismatch, "cannot update channel name")
}
c.Type = channel.Type
c.Data = channel.Data
c.Spec, c.StoredSpec = channel.Spec, channel.StoredSpec
c.UpdatedAt = time.Now()
return nil
}
// ToV1Channel copies the channel with its type named as v1 always has, after
// upstream's msteamsv2_configs list rather than the kind. Goes with the v1 API.
func (c *Channel) ToV1Channel() *Channel {
v1 := *c
if c.Type == ChannelKindMSTeams.StringValue() {
v1.Type = "msteamsv2"
}
return &v1
}
func newChannelWithoutSpec(receiver *Receiver, orgID string) (*Channel, error) {
if receiver.Name == DefaultReceiverName {
return nil, errors.Newf(errors.TypeInvalidInput, ErrCodeAlertmanagerChannelInvalid, "cannot use %s name as a channel name", receiver.Name)
}
channelType := receiverChannelType(receiver)
if channelType == "" {
return nil, errors.Newf(errors.TypeInvalidInput, ErrCodeAlertmanagerChannelInvalid, "channel '%s' must have at least one notification configuration (e.g., email_configs, webhook_configs, slack_configs)", receiver.Name)
}
data, err := json.Marshal(receiver)
if err != nil {
return nil, errors.WrapInvalidInputf(err, errors.CodeInvalidInput, "cannot save channel %q notification configuration", receiver.Name)
}
return &Channel{
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
TimeAuditable: types.TimeAuditable{CreatedAt: time.Now(), UpdatedAt: time.Now()},
Name: generateChannelName(receiver.Name),
DisplayName: receiver.Name,
Type: channelType,
Data: string(data),
OrgID: orgID,
}, nil
}
// receiverChannelType returns the type a channel is stored under: the kind for
// a notifier v2 models, else the notifier's own *_configs name. For the latter
// it walks Receiver's own fields first (native), then the embed (upstream);
// first non-empty *_configs slice wins.
func receiverChannelType(receiver *Receiver) string {
for _, channelKind := range channelKinds {
if channelKind.countConfigs(receiver) > 0 {
return channelKind.kind.StringValue()
}
}
if t := nonEmptyConfigsField(reflect.ValueOf(*receiver)); t != "" {
return t
}
if t := nonEmptyConfigsField(reflect.ValueOf(*receiver.Receiver)); t != "" {
return t
}
return ""
}
func nonEmptyConfigsField(v reflect.Value) string {
t := v.Type()
for i := 0; i < t.NumField(); i++ {
field := t.Field(i)
fieldVal := v.Field(i)
if fieldVal.Kind() != reflect.Slice || fieldVal.Len() == 0 {
continue
}
yamlTag := field.Tag.Get("yaml")
if yamlTag == "" {
continue
}
// Extract the base type name (e.g., "email_configs" -> "email").
matches := receiverTypeRegex.FindStringSubmatch(yamlTag)
if len(matches) != 2 {
continue
}
return matches[1]
}
return ""
}

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,9 +50,9 @@ 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 == "" {
if err != nil {
return nil, err
}
return []string{id}, nil
@@ -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

@@ -31,6 +31,8 @@ type ResolvedResourceWithTargetResource interface {
// IsParentChild true: the target is a child audited along but not authz-checked
// (only the source is); false: a sibling peer that is also authz-checked.
IsParentChild() bool
// HasNoLinks true: a request-phase ids extractor ran and found nothing to link.
HasNoLinks() bool
}
func NewContextWithResolvedResources(ctx context.Context, resolved []ResolvedResource) context.Context {

View File

@@ -12,6 +12,7 @@ type resolvedResourceWithTarget struct {
targetExtractor ResourceIDsExtractor
targetIDs []string
parentChild bool
noLinks bool
err error
}
@@ -44,29 +45,40 @@ func NewResolvedResourceWithTarget(
}
func (resolved *resolvedResourceWithTarget) fill(phase ExtractPhase, ec ExtractorContext) {
if resolved.sourceExtractor.IsPhase(phase) {
ids, err := resolved.sourceExtractor.Fn(ec)
if err != nil && phase == PhaseRequest {
resolved.err = err
return
}
if len(ids) > 0 {
resolved.sourceIDs = ids
}
if err := resolved.fillSide(phase, ec, resolved.sourceExtractor, &resolved.sourceIDs); err != nil {
resolved.err = err
return
}
if resolved.targetExtractor.IsPhase(phase) {
ids, err := resolved.targetExtractor.Fn(ec)
if err != nil && phase == PhaseRequest {
resolved.err = err
return
if err := resolved.fillSide(phase, ec, resolved.targetExtractor, &resolved.targetIDs); err != nil {
resolved.err = err
}
}
func (resolved *resolvedResourceWithTarget) fillSide(phase ExtractPhase, ec ExtractorContext, extractor ResourceIDsExtractor, ids *[]string) error {
if !extractor.IsPhase(phase) {
return nil
}
extracted, err := extractor.Fn(ec)
if err != nil {
if phase == PhaseRequest {
return err
}
if len(ids) > 0 {
resolved.targetIDs = ids
}
return nil
}
if len(extracted) == 0 {
if phase == PhaseRequest {
resolved.noLinks = true
}
return nil
}
*ids = extracted
return nil
}
func (resolved *resolvedResourceWithTarget) Err() error {
@@ -117,6 +129,10 @@ func (resolved *resolvedResourceWithTarget) IsParentChild() bool {
return resolved.parentChild
}
func (resolved *resolvedResourceWithTarget) HasNoLinks() bool {
return resolved.noLinks
}
func (resolved *resolvedResourceWithTarget) ResolveResponse(ec ExtractorContext) {
resolved.fill(PhaseResponse, ec)
}

View File

@@ -44,6 +44,12 @@ func (enum UnsetOrNonEmptyString) IsZero() bool {
return enum.val == ""
}
func (enum *UnsetOrNonEmptyString) SetIfUnset(val string) {
if enum.IsZero() {
enum.val = val
}
}
func (enum UnsetOrNonEmptyString) StringValue() string {
return enum.val
}

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

@@ -1,4 +1,5 @@
# pylint: disable=line-too-long
import hashlib
import json
import time
from collections.abc import Callable
@@ -8,6 +9,7 @@ import docker
import docker.errors
import pytest
import requests
from sqlalchemy import sql
from testcontainers.core.container import Network
from wiremock.testing.testcontainer import WireMockContainer
@@ -46,6 +48,34 @@ def assert_email_channel_payload_clean(payload: str) -> None:
assert SMTP_TEST_FROM not in payload
def rewrite_channel_as_legacy_receiver(signoz: types.SigNoz, channel_id: str, receiver: dict) -> None:
"""Overwrite a channel row, and its receiver in the org's alertmanager config,
the way the spec migration leaves a row it cannot fill. Neither API writes
such rows any more, so tests that need one seed it here. The receiver's name
must be the channel's display name. The alertmanager picks the swapped
receiver up on its next poll of the stored config."""
configs_key = next(key for key in receiver if key.endswith("_configs"))
# Storage names the kind; only msteams differs from its upstream configs list.
notifier_type = "msteams" if configs_key == "msteamsv2_configs" else configs_key.removesuffix("_configs")
with signoz.sqlstore.conn.connect() as conn:
conn.execute(
sql.text("UPDATE notification_channel SET type = :type, data = :data, spec = NULL WHERE id = :id"),
{"id": channel_id, "type": notifier_type, "data": json.dumps(receiver)},
)
org_id, stored = conn.execute(
sql.text("SELECT c.org_id, c.config FROM alertmanager_config c JOIN notification_channel n ON n.org_id = c.org_id WHERE n.id = :id"),
{"id": channel_id},
).one()
config = json.loads(stored)
config["receivers"] = [receiver if existing["name"] == receiver["name"] else existing for existing in config["receivers"]]
raw = json.dumps(config)
conn.execute(
sql.text("UPDATE alertmanager_config SET config = :config, hash = :hash WHERE org_id = :org_id"),
{"config": raw, "hash": hashlib.md5(raw.encode()).hexdigest(), "org_id": org_id},
)
conn.commit()
"""
Default notification channel configs shared across alertmanager tests.
"""

View File

@@ -9,6 +9,7 @@ from fixtures.auth import (
USER_ADMIN_EMAIL,
USER_ADMIN_PASSWORD,
)
from fixtures.notification_channel import rewrite_channel_as_legacy_receiver
TIMEOUT = 10
@@ -49,24 +50,45 @@ def test_repair_reports_nothing_for_a_readable_channel(
assert [channel["id"] for channel in repair["channels"]] == [channel_id]
def test_repair_deletes_a_v1_channel_of_an_unmodelled_kind(
def test_repair_deletes_a_legacy_channel_of_an_unmodelled_kind(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
name = f"v1-telegram-{uuid.uuid4().hex[:8]}"
name = f"legacy-telegram-{uuid.uuid4().hex[:8]}"
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v1/channels"),
json={"name": name, "telegram_configs": [{"chat": 12345, "token": "telegram-bot-token"}]},
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json={"name": name, "config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
rewrite_channel_as_legacy_receiver(signoz, channel_id, {"name": name, "telegram_configs": [{"chat": 12345, "token": "telegram-bot-token"}]})
# v2 lists the row with an empty kind and refuses to read it.
response = requests.get(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
params={"query": name},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
listed = response.json()["data"]
assert listed["total"] == 1
assert listed["channels"][0]["displayName"] == name
assert listed["channels"][0]["kind"] == ""
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
response = requests.post(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}/repair"),
@@ -111,29 +133,33 @@ def test_repair_deletes_a_v1_channel_of_an_unmodelled_kind(
assert response.json()["data"]["total"] == 0
def test_repair_splits_a_v1_channel_carrying_several_notifiers(
def test_repair_splits_a_legacy_channel_carrying_several_notifiers(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
name = f"v1-fanout-{uuid.uuid4().hex[:8]}"
name = f"legacy-fanout-{uuid.uuid4().hex[:8]}"
# Only v1 accepts a receiver with more than one notifier configuration.
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v1/channels"),
json={
"name": name,
"slack_configs": [{"api_url": "https://hooks.slack.test/services/T/B/X", "channel": "#alerts"}],
"webhook_configs": [{"url": "https://webhook.test/hook"}],
},
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json={"name": name, "config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X", "channel": "#alerts"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
rewrite_channel_as_legacy_receiver(
signoz,
channel_id,
{
"name": name,
"slack_configs": [{"api_url": "https://hooks.slack.test/services/T/B/X", "channel": "#alerts"}],
"webhook_configs": [{"url": "https://webhook.test/hook"}],
},
)
response = requests.post(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}/repair"),
@@ -189,3 +215,65 @@ def test_repair_splits_a_v1_channel_carrying_several_notifiers(
)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["config"]["kind"] == "slack"
def test_repair_splits_off_an_unmodelled_notifier_for_its_own_delete(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
name = f"legacy-mixed-{uuid.uuid4().hex[:8]}"
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json={"name": name, "config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X", "channel": "#alerts"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
rewrite_channel_as_legacy_receiver(
signoz,
channel_id,
{
"name": name,
"slack_configs": [{"api_url": "https://hooks.slack.test/services/T/B/X", "channel": "#alerts"}],
"telegram_configs": [{"chat": 12345, "token": "telegram-bot-token"}],
},
)
response = requests.post(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}/repair"),
params={"apply": "true"},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
repair = response.json()["data"]
assert repair["defect"] == "multiple_notifiers"
assert repair["applied"] is True
assert [channel["displayName"] for channel in repair["channels"]] == [name, f"{name} (2)"]
assert [channel["kind"] for channel in repair["channels"]] == ["slack", ""]
telegram_id = repair["channels"][1]["id"]
cleanup_notification_channels.append(telegram_id)
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["config"]["kind"] == "slack"
response = requests.post(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{telegram_id}/repair"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
repair = response.json()["data"]
assert repair["defect"] == "unsupported_notifier"
assert repair["action"] == "delete"

View File

@@ -0,0 +1,555 @@
import uuid
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,
)
TIMEOUT = 10
V2_BASE_URL = "/api/v2/notification_channels"
@pytest.mark.parametrize(
"clashing_field,message_fragment",
[
pytest.param("name", "with name", id="name"),
pytest.param("displayName", "with display name", id="display_name"),
],
)
def test_create_rejects_a_duplicate_with_conflict( # pylint: disable=too-many-arguments,too-many-positional-arguments
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
clashing_field: str,
message_fragment: str,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
shared = f"v2-dup-{uuid.uuid4().hex[:8]}"
first = {"name": f"{shared}-first", "displayName": f"{shared} first", "config": {"kind": "email", "spec": {"to": "first@integration.test", "html": "<p>body</p>"}}}
first[clashing_field] = shared
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json=first,
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
cleanup_notification_channels.append(response.json()["data"]["id"])
second = {"name": f"{shared}-second", "displayName": f"{shared} second", "config": {"kind": "email", "spec": {"to": "second@integration.test", "html": "<p>body</p>"}}}
second[clashing_field] = shared
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json=second,
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CONFLICT, response.text
# Both v2 conflicts share a status and an error code, so only the message
# separates a clashing display name from a clashing name.
assert message_fragment in response.text
@pytest.mark.parametrize(
"body",
[
pytest.param(
{
"name": "Not_A_Label",
"config": {"kind": "email", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}},
},
id="name_not_dns1123_label",
),
pytest.param(
{"config": {"kind": "email", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}}},
id="no_name_and_no_generate_name",
),
pytest.param(
{
"name": "explicit",
"generateName": True,
"displayName": "Explicit",
"config": {"kind": "email", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}},
},
id="name_with_generate_name",
),
pytest.param(
{
"generateName": True,
"config": {"kind": "email", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}},
},
id="generate_name_without_display_name",
),
pytest.param(
{
"name": "default-receiver",
"config": {"kind": "email", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}},
},
id="reserved_receiver_name",
),
pytest.param({"name": "rejected"}, id="no_config"),
pytest.param(
{"name": "rejected", "config": {"kind": "telegram", "spec": {"chatId": 1}}},
id="unmodelled_kind",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "slack",
"spec": {
"apiUrl": "https://hooks.slack.test/services/T/B/X",
"channel": "#a",
"text": "body",
"iconEmoji": ":tada:",
},
},
},
id="unknown_spec_field",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "slack", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}},
},
id="spec_of_another_kind",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "slack", "spec": {"channel": "#alerts", "title": "Alert", "text": "body"}},
},
id="slack_without_api_url",
),
pytest.param(
{"name": "rejected", "config": {"kind": "email", "spec": {"html": "<p>body</p>"}}},
id="email_without_to",
),
pytest.param(
{"name": "rejected", "config": {"kind": "webhook", "spec": {}}},
id="webhook_without_url",
),
pytest.param(
{"name": "rejected", "config": {"kind": "pagerduty", "spec": {"description": "body"}}},
id="pagerduty_without_routing_key",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "opsgenie", "spec": {"message": "subject", "description": "body"}},
},
id="opsgenie_without_api_key",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "msteams", "spec": {"title": "Alert", "text": "body"}},
},
id="msteams_without_webhook_url",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "googlechat", "spec": {"title": "Alert", "text": "body"}},
},
id="googlechat_without_webhook_url",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"project": "OPS",
"issueType": "Bug",
"email": "oncall@integration.test",
"apiToken": "jira-api-token",
},
},
},
id="jira_without_site",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"site": "https://acme.atlassian.net",
"issueType": "Bug",
"email": "oncall@integration.test",
"apiToken": "jira-api-token",
},
},
},
id="jira_without_project",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"site": "https://acme.atlassian.net",
"project": "OPS",
"email": "oncall@integration.test",
"apiToken": "jira-api-token",
},
},
},
id="jira_without_issue_type",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"site": "https://acme.atlassian.net",
"project": "OPS",
"issueType": "Bug",
"apiToken": "jira-api-token",
},
},
},
id="jira_without_email",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"site": "https://acme.atlassian.net",
"project": "OPS",
"issueType": "Bug",
"email": "oncall@integration.test",
},
},
},
id="jira_without_api_token",
),
pytest.param(
{"name": "rejected", "config": {"kind": "jsmops", "spec": {}}},
id="jsmops_without_api_key",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "incidentio", "spec": {"token": "incidentio-token"}},
},
id="incidentio_without_url",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "incidentio",
"spec": {"url": "https://api.incident.io/v2/alert_events/http/01ABCDEF"},
},
},
id="incidentio_without_token",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "slack",
"spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X", "fields": [{"title": "Severity"}]},
},
},
id="slack_field_without_value",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "slack",
"spec": {
"apiUrl": "https://hooks.slack.test/services/T/B/X",
"actions": [{"type": "button", "url": "https://signoz.test"}],
},
},
},
id="slack_action_without_text",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "slack",
"spec": {
"apiUrl": "https://hooks.slack.test/services/T/B/X",
"actions": [{"type": "button", "text": "Open"}],
},
},
},
id="slack_action_without_url_or_name",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "slack",
"spec": {
"apiUrl": "https://hooks.slack.test/services/T/B/X",
"actions": [{"type": "button", "text": "Ack", "name": "ack", "confirm": {"title": "Sure?"}}],
},
},
},
id="slack_action_confirm_without_text",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "email", "spec": {"to": "a@integration.test", "html": "<p>body</p>"}},
"type": "this key is not a valid",
},
id="unknown_envelope_field",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "webhook",
"spec": {"url": "https://webhook.test/hook", "username": "u", "password": "p", "bearerToken": "t"},
},
},
id="webhook_basic_auth_with_bearer_token",
),
# The next two break a rule of the notifier rather than of the request
# shape, and still surface as a 400.
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"site": "https://acme.atlassian.net",
"project": "OPS",
"issueType": "Bug",
"email": "a@integration.test",
"apiToken": "t",
"summary": "Alert",
"description": "body",
"reopenDuration": "30s",
},
},
},
id="jira_reopen_duration_below_a_minute",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "incidentio",
"spec": {
"url": "https://api.incident.io/v2/alert_events/http/01ABCDEF",
"token": "Bearer incidentio-token",
"title": "Alert",
"description": "body",
},
},
},
id="incidentio_token_with_bearer_prefix",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "slack",
"spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X", "title": ""},
},
},
id="slack_title_empty_instead_of_omitted",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "jsmops", "spec": {"apiKey": "jsm-api-key", "tags": ""}},
},
id="jsmops_tags_empty_instead_of_omitted",
),
pytest.param(
{
"name": "rejected",
"config": {
"kind": "jira",
"spec": {
"site": "https://acme.atlassian.net",
"project": "OPS",
"issueType": "Bug",
"email": "a@integration.test",
"apiToken": "t",
"reopenDuration": "72h",
},
},
},
id="jira_reopen_duration_not_as_reported",
),
pytest.param(
{
"name": "rejected",
"config": {"kind": "email", "spec": {"to": "a@integration.test", "headers": {"subject": "must be written in canonical form, Subject"}}},
},
id="email_header_name_not_canonical",
),
],
)
def test_create_rejects_invalid_bodies(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
body: dict,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json=body,
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
@pytest.mark.parametrize(
"params",
[
pytest.param({"sort": "data"}, id="sort_outside_the_enum"),
pytest.param({"order": "sideways"}, id="order_outside_the_enum"),
pytest.param({"kind": "telegram"}, id="kind_outside_the_enum"),
pytest.param({"limit": -1}, id="negative_limit"),
pytest.param({"offset": -1}, id="negative_offset"),
],
)
def test_list_rejects_invalid_params(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
params: dict,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
params=params,
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
def test_get_unknown_id(
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"{V2_BASE_URL}/0199a1b2-c3d4-7000-8000-000000000000"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.NOT_FOUND, response.text
@pytest.mark.parametrize(
"body",
[
pytest.param(
{
"name": "renamed",
"config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X"}},
},
id="name_in_body",
),
pytest.param(
{
"displayName": "Renamed",
"config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X"}},
},
id="display_name_in_body",
),
pytest.param(
{"config": {"kind": "slack", "spec": {}}},
id="spec_missing_required_field",
),
pytest.param({}, id="no_config"),
],
)
def test_update_rejects_invalid_bodies( # pylint: disable=too-many-arguments,too-many-positional-arguments
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
body: dict,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
name = f"v2-badupdate-{uuid.uuid4().hex[:8]}"
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json={"name": name, "config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
response = requests.put(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
json=body,
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
@pytest.mark.parametrize(
"body",
[
pytest.param(
{
"name": "test-send",
"config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X"}},
},
id="name_in_body",
),
pytest.param(
{"config": {"kind": "slack", "spec": {}}},
id="spec_missing_required_field",
),
pytest.param(
{"config": {"kind": "telegram", "spec": {"chatId": 1}}},
id="unmodelled_kind",
),
pytest.param({}, id="no_config"),
],
)
def test_test_rejects_invalid_bodies(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
body: dict,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.post(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/test"),
json=body,
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text

View File

@@ -0,0 +1,239 @@
import json
import time
import uuid
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
import pytest
import requests
from wiremock.client import HttpMethods, Mapping, MappingRequest, MappingResponse
from fixtures import types
from fixtures.alerts import update_rule_channel_name, verify_notification_expectation
from fixtures.auth import (
USER_ADMIN_EMAIL,
USER_ADMIN_PASSWORD,
)
from fixtures.fs import get_testdata_file_path
from fixtures.notification_channel import rewrite_channel_as_legacy_receiver
TIMEOUT = 10
V2_BASE_URL = "/api/v2/notification_channels"
def test_get_reflects_a_v2_update_after_a_v1_create(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
display_name = f"V1 then V2 {uuid.uuid4().hex[:8]}"
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v1/channels"),
json={"name": display_name, "slack_configs": [{"api_url": "https://hooks.slack.test/services/T/B/V1", "channel": "#from-v1"}]},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["config"]["spec"]["channel"] == "#from-v1"
response = requests.put(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
json={"config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/V2", "channel": "#from-v2"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["config"]["spec"]["channel"] == "#from-v2"
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
fetched = response.json()["data"]
assert fetched["displayName"] == display_name
assert fetched["config"]["spec"]["apiUrl"] == "https://hooks.slack.test/services/T/B/V2"
assert fetched["config"]["spec"]["channel"] == "#from-v2"
@pytest.mark.parametrize(
"receiver",
[
pytest.param(
{"telegram_configs": [{"chat": 12345, "token": "telegram-bot-token"}]},
id="kind_v2_does_not_model",
),
pytest.param(
{
"slack_configs": [{"api_url": "https://hooks.slack.test/services/T/B/X", "channel": "#alerts"}],
"webhook_configs": [{"url": "https://webhook.test/hook"}],
},
id="several_notifiers",
),
pytest.param(
{"webhook_configs": [{"url": "https://webhook.test/hook", "http_config": {"proxy_url": "http://proxy.test:3128"}}]},
id="unsupported_http_config",
),
],
)
def test_v1_rejects_a_receiver_v2_cannot_represent(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
receiver: dict,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v1/channels"),
json={"name": f"v1-rejected-{uuid.uuid4().hex[:8]}", **receiver},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
def test_update_retypes_a_legacy_channel_of_an_unmodelled_kind(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
cleanup_notification_channels: list[str],
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
name = f"legacy-telegram-{uuid.uuid4().hex[:8]}"
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json={"name": name, "config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
rewrite_channel_as_legacy_receiver(signoz, channel_id, {"name": name, "telegram_configs": [{"chat": 12345, "token": "telegram-bot-token"}]})
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
# A v2 update needs nothing from the stored config, so it can rewrite a
# channel v2 cannot read.
response = requests.put(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
json={"config": {"kind": "slack", "spec": {"apiUrl": "https://hooks.slack.test/services/T/B/X", "channel": "#retyped"}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["config"]["kind"] == "slack"
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
fetched = response.json()["data"]
assert fetched["displayName"] == name
assert fetched["config"]["kind"] == "slack"
assert fetched["config"]["spec"]["channel"] == "#retyped"
response = requests.get(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
params={"query": name},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["channels"][0]["kind"] == "slack"
def test_alerts_still_reach_a_legacy_channel_v2_cannot_read( # pylint: disable=too-many-arguments,too-many-positional-arguments
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
notification_channel: types.TestContainerDocker,
make_http_mocks: Callable[[types.TestContainerDocker, list[Mapping]], None],
cleanup_notification_channels: list[str],
create_alert_rule: Callable[[dict], str],
insert_alert_data: Callable[[list[types.AlertData], datetime], None],
maildev: types.TestContainerDocker,
) -> None:
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
name = f"legacy-delivery-{uuid.uuid4().hex[:8]}"
slack_path = f"/services/T/B/{name}"
webhook_path = f"/webhook/{name}"
make_http_mocks(
notification_channel,
[Mapping(request=MappingRequest(method=HttpMethods.POST, url=path), response=MappingResponse(status=200, json_body={}), persistent=False) for path in (slack_path, webhook_path)],
)
response = requests.post(
signoz.self.host_configs["8080"].get(V2_BASE_URL),
json={"name": name, "config": {"kind": "webhook", "spec": {"url": notification_channel.container_configs["8080"].get(webhook_path)}}},
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.CREATED, response.text
channel_id = response.json()["data"]["id"]
cleanup_notification_channels.append(channel_id)
rewrite_channel_as_legacy_receiver(
signoz,
channel_id,
{
"name": name,
"slack_configs": [{"api_url": notification_channel.container_configs["8080"].get(slack_path), "channel": "#legacy"}],
"webhook_configs": [{"url": notification_channel.container_configs["8080"].get(webhook_path)}],
},
)
response = requests.get(
signoz.self.host_configs["8080"].get(f"{V2_BASE_URL}/{channel_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=TIMEOUT,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
# The rule must not fire before the alertmanager has polled the swapped receiver.
time.sleep(12)
insert_alert_data(
[types.AlertData(type="metrics", data_path="ruler/test_scenarios/threshold_above_at_least_once/alert_data.jsonl")],
base_time=datetime.now(tz=UTC) - timedelta(minutes=5),
)
with open(get_testdata_file_path("ruler/test_scenarios/threshold_above_at_least_once/rule.json"), encoding="utf-8") as f:
rule_data = json.load(f)
update_rule_channel_name(rule_data, name)
create_alert_rule(rule_data)
verify_notification_expectation(
notification_channel,
maildev,
types.AMNotificationExpectation(
should_notify=True,
wait_time_seconds=120,
notification_validations=[
types.NotificationValidation(destination_type="webhook", validation_data={"path": slack_path, "json_body": {"channel": "#legacy"}}),
types.NotificationValidation(destination_type="webhook", validation_data={"path": webhook_path, "json_body": {"status": "firing", "receiver": name}}),
],
),
)

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

@@ -275,11 +275,51 @@ def test_user_with_roles_reflects_change(
assert "signoz-admin" in role_names
def test_admin_cannot_assign_role_to_self(
def test_admin_can_change_own_roles(
signoz: types.SigNoz,
get_token: Callable[[str, str], str],
):
"""Verify POST /api/v2/user_roles for the caller's own user is rejected (self-mutation guard)."""
"""Verify a non-root admin can assign a role to and remove a role from their own user."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
self_token = get_token(ROLECHANGE_USER_EMAIL, ROLECHANGE_USER_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
headers={"Authorization": f"Bearer {self_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK
self_id = response.json()["data"]["id"]
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/user_roles"),
json={"userId": self_id, "roleId": find_role_by_name(signoz, admin_token, "signoz-editor")},
headers={"Authorization": f"Bearer {self_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.CREATED, response.text
editor_entry_id = response.json()["data"]["id"]
response = requests.delete(
signoz.self.host_configs["8080"].get(f"/api/v2/user_roles/{editor_entry_id}"),
headers={"Authorization": f"Bearer {self_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
headers={"Authorization": f"Bearer {self_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK
assert {ur["role"]["name"] for ur in response.json()["data"]["userRoles"]} == {"signoz-admin"}
def test_root_roles_cannot_be_changed(
signoz: types.SigNoz,
get_token: Callable[[str, str], str],
):
"""Verify role assignment and removal on the root user are rejected by the root user guard."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
@@ -295,22 +335,7 @@ def test_admin_cannot_assign_role_to_self(
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.BAD_REQUEST
def test_admin_cannot_remove_own_role(
signoz: types.SigNoz,
get_token: Callable[[str, str], str],
):
"""Verify DELETE /api/v2/user_roles/{id} for the caller's own assignment is rejected (self-mutation guard)."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK
admin_data = response.json()["data"]
assert response.status_code == HTTPStatus.NOT_IMPLEMENTED
admin_entry_id = next((ur["id"] for ur in admin_data["userRoles"] if ur["role"]["name"] == "signoz-admin"), None)
assert admin_entry_id is not None
@@ -320,7 +345,7 @@ def test_admin_cannot_remove_own_role(
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.BAD_REQUEST
assert response.status_code == HTTPStatus.NOT_IMPLEMENTED
def test_editor_cannot_manage_roles(

View File

@@ -112,8 +112,8 @@ def test_update_my_user(signoz: types.SigNoz, get_token: Callable[[str, str], st
assert response.json()["data"]["displayName"] == "self updated editor"
def test_admin_cannot_update_self_via_id(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
"""Verify PUT /api/v2/users/{own_id} is rejected (self-mutation guard)."""
def test_root_cannot_be_updated_via_id(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
"""Verify PUT /api/v2/users/{root_id} is rejected by the root user guard."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
@@ -130,7 +130,37 @@ def test_admin_cannot_update_self_via_id(signoz: types.SigNoz, get_token: Callab
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.BAD_REQUEST
assert response.status_code == HTTPStatus.NOT_IMPLEMENTED
def test_admin_can_update_self_via_id(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
"""Verify a non-root admin can call PUT /api/v2/users/{own_id}."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
self_admin_id = create_active_user(
signoz,
admin_token,
email="admin+selfupdate@integration.test",
role="signoz-admin",
password="password123Z$",
name="self admin",
)
self_admin_token = get_token("admin+selfupdate@integration.test", "password123Z$")
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v2/users/{self_admin_id}"),
json={"displayName": "self admin updated"},
headers={"Authorization": f"Bearer {self_admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
headers={"Authorization": f"Bearer {self_admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK
assert response.json()["data"]["displayName"] == "self admin updated"
def test_editor_cannot_list_users(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:

View File

@@ -0,0 +1,295 @@
from collections.abc import Callable
from http import HTTPStatus
import requests
from fixtures import types
from fixtures.auth import (
USER_ADMIN_EMAIL,
USER_ADMIN_PASSWORD,
USER_EDITOR_EMAIL,
USER_EDITOR_PASSWORD,
USER_ROLES_BASE,
USERS_BASE,
change_user_role,
create_active_user,
find_user_by_email,
)
from fixtures.role import find_role_by_name, transaction_group
_ACTOR_ROLE_NAME = "user-fga-actor"
_ACTOR_EMAIL = "customrole+userfga@integration.test"
_ACTOR_PASSWORD = "password123Z$"
_VIEWER_EMAIL = "viewer+userfga@integration.test"
_VIEWER_PASSWORD = "password123Z$"
# Instance verbs are granted on _TARGET_EMAIL's id only; _OTHER_EMAIL must stay forbidden.
_TARGET_EMAIL = "target+userfga@integration.test"
_OTHER_EMAIL = "other+userfga@integration.test"
_TARGET_PASSWORD = "password123Z$"
_INVITED_VIEWER_EMAIL = "invited-viewer+userfga@integration.test"
_INVITED_EDITOR_EMAIL = "invited-editor+userfga@integration.test"
_INVITED_NO_ROLE_EMAIL = "invited-norole+userfga@integration.test"
def _set_actor_role(signoz: types.SigNoz, admin_token: str, transaction_groups: list[dict]) -> None:
role_id = find_role_by_name(signoz, admin_token, _ACTOR_ROLE_NAME)
resp = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/roles/{role_id}"),
json={"description": "", "transactionGroups": transaction_groups},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.NO_CONTENT, resp.text
def test_setup_actor_and_targets(
signoz: types.SigNoz,
get_token: Callable[[str, str], str],
create_role: Callable[..., str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
create_active_user(signoz, admin_token, email=_VIEWER_EMAIL, role="signoz-viewer", password=_VIEWER_PASSWORD, name="user-fga-viewer")
create_active_user(signoz, admin_token, email=_TARGET_EMAIL, role="signoz-viewer", password=_TARGET_PASSWORD, name="user-fga-target")
create_active_user(signoz, admin_token, email=_OTHER_EMAIL, role="signoz-viewer", password=_TARGET_PASSWORD, name="user-fga-other")
target_id = find_user_by_email(signoz, admin_token, _TARGET_EMAIL)["id"]
create_role(
admin_token,
_ACTOR_ROLE_NAME,
[
transaction_group("read", "user", "user", [target_id]),
transaction_group("list", "user", "user", ["*"]),
],
)
actor_id = create_active_user(signoz, admin_token, email=_ACTOR_EMAIL, role="signoz-viewer", password=_ACTOR_PASSWORD, name="user-fga-actor")
change_user_role(signoz, admin_token, actor_id, "signoz-viewer", _ACTOR_ROLE_NAME)
def test_managed_roles_matrix(signoz: types.SigNoz, get_token: Callable[[str, str], str]):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
target_id = find_user_by_email(signoz, admin_token, _TARGET_EMAIL)["id"]
viewer_role_id = find_role_by_name(signoz, admin_token, "signoz-viewer")
for email, password in ((USER_EDITOR_EMAIL, USER_EDITOR_PASSWORD), (_VIEWER_EMAIL, _VIEWER_PASSWORD)):
token = get_token(email, password)
headers = {"Authorization": f"Bearer {token}"}
resp = requests.get(signoz.self.host_configs["8080"].get(USERS_BASE), headers=headers, timeout=5)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"{email} list users: expected 403, got {resp.status_code}: {resp.text}"
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}"), headers=headers, timeout=5)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"{email} get user: expected 403, got {resp.status_code}: {resp.text}"
resp = requests.put(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/reset_password_tokens"), headers=headers, timeout=5)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"{email} reset token: expected 403, got {resp.status_code}: {resp.text}"
resp = requests.post(
signoz.self.host_configs["8080"].get(USER_ROLES_BASE),
json={"userId": target_id, "roleId": viewer_role_id},
headers=headers,
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"{email} assign role: expected 403, got {resp.status_code}: {resp.text}"
admin_headers = {"Authorization": f"Bearer {admin_token}"}
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}"), headers=admin_headers, timeout=5)
assert resp.status_code == HTTPStatus.OK, resp.text
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/reset_password_tokens"), headers=admin_headers, timeout=5)
assert resp.status_code in (HTTPStatus.OK, HTTPStatus.NOT_FOUND), resp.text
resp = requests.put(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/reset_password_tokens"), headers=admin_headers, timeout=5)
assert resp.status_code == HTTPStatus.CREATED, resp.text
def test_read_scoped_to_granted_user(signoz: types.SigNoz, get_token: Callable[[str, str], str]):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
token = get_token(_ACTOR_EMAIL, _ACTOR_PASSWORD)
target_id = find_user_by_email(signoz, admin_token, _TARGET_EMAIL)["id"]
other_id = find_user_by_email(signoz, admin_token, _OTHER_EMAIL)["id"]
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.OK, f"get granted user: {resp.text}"
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/roles"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.OK, f"get granted user roles: {resp.text}"
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{other_id}"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"get other user: expected 403, got {resp.status_code}: {resp.text}"
# list is collection-scoped: list on "*" returns every user, including the one
# the actor cannot read individually.
resp = requests.get(signoz.self.host_configs["8080"].get(USERS_BASE), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.OK, resp.text
ids = {user["id"] for user in resp.json()["data"]}
assert {target_id, other_id} <= ids
resp = requests.put(
signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}"),
json={"displayName": "user-fga-target-renamed"},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"update user without grant: expected 403, got {resp.status_code}: {resp.text}"
def test_reset_password_token_scoped_to_granted_user(signoz: types.SigNoz, get_token: Callable[[str, str], str]):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
target_id = find_user_by_email(signoz, admin_token, _TARGET_EMAIL)["id"]
other_id = find_user_by_email(signoz, admin_token, _OTHER_EMAIL)["id"]
_set_actor_role(
signoz,
admin_token,
[
transaction_group("read", "user", "user", [target_id]),
transaction_group("attach", "user", "user", [target_id]),
transaction_group("create", "metaresource", "factor-password", ["*"]),
transaction_group("list", "metaresource", "factor-password", ["*"]),
],
)
token = get_token(_ACTOR_EMAIL, _ACTOR_PASSWORD)
resp = requests.put(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/reset_password_tokens"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.CREATED, f"create reset token for granted user: {resp.text}"
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/reset_password_tokens"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.OK, f"get reset token for granted user: {resp.text}"
resp = requests.put(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{other_id}/reset_password_tokens"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"create reset token for other user: expected 403, got {resp.status_code}: {resp.text}"
# user:attach alone is not enough: factor-password:create is checked too.
_set_actor_role(signoz, admin_token, [transaction_group("attach", "user", "user", [target_id])])
token = get_token(_ACTOR_EMAIL, _ACTOR_PASSWORD)
resp = requests.put(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{target_id}/reset_password_tokens"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"create reset token without factor-password:create: expected 403, got {resp.status_code}: {resp.text}"
def test_user_role_attach_detach_dual_scoped(signoz: types.SigNoz, get_token: Callable[[str, str], str]):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
target_id = find_user_by_email(signoz, admin_token, _TARGET_EMAIL)["id"]
other_id = find_user_by_email(signoz, admin_token, _OTHER_EMAIL)["id"]
editor_role_id = find_role_by_name(signoz, admin_token, "signoz-editor")
admin_role_id = find_role_by_name(signoz, admin_token, "signoz-admin")
# attach/detach granted on the target user id AND the signoz-editor role name only.
_set_actor_role(
signoz,
admin_token,
[
transaction_group("read", "user", "user", [target_id]),
transaction_group("attach", "user", "user", [target_id]),
transaction_group("detach", "user", "user", [target_id]),
transaction_group("attach", "role", "role", ["signoz-editor"]),
transaction_group("detach", "role", "role", ["signoz-editor"]),
],
)
token = get_token(_ACTOR_EMAIL, _ACTOR_PASSWORD)
resp = requests.post(
signoz.self.host_configs["8080"].get(USER_ROLES_BASE),
json={"userId": target_id, "roleId": editor_role_id},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.CREATED, f"assign editor to target: {resp.text}"
editor_entry_id = resp.json()["data"]["id"]
resp = requests.get(signoz.self.host_configs["8080"].get(f"{USER_ROLES_BASE}/{editor_entry_id}"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.OK, f"get user role of granted user: {resp.text}"
resp = requests.post(
signoz.self.host_configs["8080"].get(USER_ROLES_BASE),
json={"userId": other_id, "roleId": editor_role_id},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"assign editor to other user: expected 403, got {resp.status_code}: {resp.text}"
resp = requests.post(
signoz.self.host_configs["8080"].get(USER_ROLES_BASE),
json={"userId": target_id, "roleId": admin_role_id},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"assign admin to target: expected 403, got {resp.status_code}: {resp.text}"
resp = requests.delete(signoz.self.host_configs["8080"].get(f"{USER_ROLES_BASE}/{editor_entry_id}"), headers={"Authorization": f"Bearer {token}"}, timeout=5)
assert resp.status_code == HTTPStatus.NO_CONTENT, f"remove editor from target: {resp.text}"
def test_invite_checks_role_attach_per_role(signoz: types.SigNoz, get_token: Callable[[str, str], str]):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
viewer_role_id = find_role_by_name(signoz, admin_token, "signoz-viewer")
editor_role_id = find_role_by_name(signoz, admin_token, "signoz-editor")
_set_actor_role(
signoz,
admin_token,
[
transaction_group("create", "user", "user", ["*"]),
transaction_group("attach", "user", "user", ["*"]),
transaction_group("attach", "role", "role", ["signoz-viewer"]),
],
)
token = get_token(_ACTOR_EMAIL, _ACTOR_PASSWORD)
resp = requests.post(
signoz.self.host_configs["8080"].get(USERS_BASE),
json={"email": _INVITED_VIEWER_EMAIL, "displayName": "invited-viewer", "userRoles": [{"id": viewer_role_id}]},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.CREATED, f"invite with viewer: {resp.text}"
resp = requests.post(
signoz.self.host_configs["8080"].get(USERS_BASE),
json={"email": _INVITED_EDITOR_EMAIL, "displayName": "invited-editor", "userRoles": [{"id": editor_role_id}]},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"invite with editor: expected 403, got {resp.status_code}: {resp.text}"
# No roles in the body: nothing to attach, so neither user:attach nor role:attach is checked.
_set_actor_role(signoz, admin_token, [transaction_group("create", "user", "user", ["*"])])
token = get_token(_ACTOR_EMAIL, _ACTOR_PASSWORD)
resp = requests.post(
signoz.self.host_configs["8080"].get(USERS_BASE),
json={"email": _INVITED_NO_ROLE_EMAIL, "displayName": "invited-norole", "userRoles": []},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.CREATED, f"invite without roles: {resp.text}"
resp = requests.post(
signoz.self.host_configs["8080"].get(USERS_BASE),
json={"email": "invited-viewer-denied+userfga@integration.test", "displayName": "invited-viewer-denied", "userRoles": [{"id": viewer_role_id}]},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert resp.status_code == HTTPStatus.FORBIDDEN, f"invite with viewer without role:attach: expected 403, got {resp.status_code}: {resp.text}"
def test_cleanup(signoz: types.SigNoz, get_token: Callable[[str, str], str]):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
headers = {"Authorization": f"Bearer {admin_token}"}
actor_id = find_user_by_email(signoz, admin_token, _ACTOR_EMAIL)["id"]
change_user_role(signoz, admin_token, actor_id, _ACTOR_ROLE_NAME, "signoz-viewer")
role_id = find_role_by_name(signoz, admin_token, _ACTOR_ROLE_NAME)
resp = requests.delete(signoz.self.host_configs["8080"].get(f"/api/v1/roles/{role_id}"), headers=headers, timeout=5)
assert resp.status_code == HTTPStatus.NO_CONTENT, resp.text
for email in (_ACTOR_EMAIL, _VIEWER_EMAIL, _TARGET_EMAIL, _OTHER_EMAIL, _INVITED_VIEWER_EMAIL, _INVITED_NO_ROLE_EMAIL):
user_id = find_user_by_email(signoz, admin_token, email)["id"]
resp = requests.delete(signoz.self.host_configs["8080"].get(f"{USERS_BASE}/{user_id}"), headers=headers, timeout=5)
assert resp.status_code == HTTPStatus.NO_CONTENT, f"delete {email}: {resp.text}"

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