Compare commits

...

7 Commits

Author SHA1 Message Date
Nityananda Gohain
a75442f31e fix: add quick filters v2 api to support TelemetryFieldKey (#12698)
Some checks are pending
build-staging / staging (push) Blocked by required conditions
build-staging / prepare (push) Waiting to run
build-staging / js-build (push) Blocked by required conditions
build-staging / go-build (push) Blocked by required conditions
cacheci / tests (push) Waiting to run
Release Drafter / update_release_draft (push) Waiting to run
#### Description
The old API didn't support telemetryFieldKey, so adding a new v2 API to
support it.

This PR
* Migrates old data to the new one.
* Existing API's now internally stores it in the new struct so that they
don't break the UI.

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

## Additional details
* the old api is safe with new field as it is just a subset of it.
2026-09-03 07:49:31 +00:00
Nityananda Gohain
52588c4582 fix: support for related values in ai field values (#12716)
#### Description

Adds support for related values in ai observability field values.

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

Closes https://github.com/SigNoz/engineering-pod/issues/5975
2026-09-03 07:12:22 +00:00
Nityananda Gohain
e485f50221 feat: support ai trace aggregate filtering in ai span list (#12122)
## Pull Request

---

### 📄 Summary
* Span list (raw) queries can now leverage trace-level `trace.` filter
conditions leading to the trace-level component being qualified with
`__trace_scope`.
* An error is now raised if the span list is ordered by the trace level
key.


The following constraint is known: in case of ordering by `timestamp`,
the querier processes the span list by time bucket thus performing
trace-level aggregates calculation per bucket not within a time window.
The records are kept separately.


#### Issues closed by this PR

Fixes https://github.com/SigNoz/engineering-pod/issues/5976

---

###  Change Type
_Select all that apply_

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

---

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

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

---

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

- Blast radius: None
- Potential regressions:
- Rollback plan:
2026-09-03 07:10:24 +00:00
Vikrant Gupta
b3b547f34a feat(gateway): add first-class ingestion limit APIs (#12625)
#### Description

- Adds first-class ingestion limit APIs under
`/api/v2/gateway/ingestion_limits`: create (`keyId` in body), get,
update, and delete by `{limitId}`. Get proxies the new upstream `GET
/v1/workspaces/me/limits/{limitID}`.
- Adds key read APIs: `GET /api/v2/gateway/ingestion_keys/{keyId}` (key
by id — upstream does not embed limits here) and `GET
/api/v2/gateway/ingestion_keys/{keyId}/limits` (limits for a key, with
current-period usage metrics).
- Marks the existing limit routes (`POST
/ingestion_keys/{keyId}/limits`, `PATCH/DELETE
/ingestion_keys/limits/{limitId}`) as deprecated; they keep working
unchanged.
- Renames the old create body to `DeprecatedPostableIngestionKeyLimit`;
`PostableIngestionKeyLimit` is now the first-class body carrying
`keyId`. Handlers decode via `binding.JSON` and the create response is
`types.Identifiable`.

Part of SigNoz/platform-pod#2651.

#### Additional Information

- OpenAPI spec and the generated frontend client are regenerated; the UI
stays on the deprecated routes for now.
- Requires the upstream get-by-id endpoints from
SigNoz/opentelemetry-gateway#96 (merged and deployed).
2026-09-03 07:07:10 +00:00
Naman Verma
be9a045455 chore: add internal name to data layer of notification channels (#12754)
<!--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

Convert the existing `name` to store an immutable DNS1123 internal name
of a notification_channel, and add a display name column where the data
from the existing `name` column will go.

This is just a database level change. The internal name is not being
used by any consumer, be it the API or rules or route policies. All that
will come in subsequent PRs

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

Part of https://github.com/SigNoz/pulse-pod/issues/296

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

Eventually references (rules, routing policies) migrate onto the
internal name, freeing the display name to become a user-editable. But
that will happen post rules migration so that all rules are on v2.


<!--Please delete paragraphs that you did not use before submitting.-->
2026-09-03 06:52:58 +00:00
Nikhil Mantri
297d3dd44b feat(alert-channel-integrations): incidentio frontend (#12645)
## Description

Frontend for the **incident.io** alert channel (backend in #12644 — this
PR is stacked on it).

- Adds **incident.io** to the channel-type dropdown with a settings
form: alert source **URL** + **token** (both required), and
**title/description** template fields prefilled with the backend
defaults — same UX as Jira/JSM.
- A tip above the form links to the setup docs (create an HTTP alert
source in incident.io, copy URL + token).
- Client-side validation mirrors the backend for a nicer error
experience: both fields required, URL must be an alert events URL
(`…/v2/alert_events/http/<source_config_id>`).
- `send_resolved` is seeded **on** so incident.io alerts resolve with
the rule (backend can't default it — same reasoning as JSM).
- **Additional metadata** section: key-value rows merged into every
alert's metadata on top of the alert's labels (channel wins on clash);
values may use templates.
- Create, edit and test-channel flows all wired; editing prefills from
the stored `incidentio_configs`.

Notes for the reviewer:

- Follows the JSM Ops form/handler pattern file-for-file; no new
patterns introduced.
- Tests: 4 create-flow cases (fields render, required-field error, URL
validation error, payload shape with defaults) + 1 edit-flow payload
case.

## Issues closed by this PR

Closes SigNoz/pulse-pod#173

---------

Co-authored-by: Naman Verma <naman.verma@signoz.io>
2026-09-03 06:20:47 +00:00
Nikhil Mantri
e84a61d43f feat(alert-channel-integrations): incidentio channel integration (#12644)
## Description

Adds **incident.io** as a native alert notification channel, using
incident.io's HTTP alert source (Alert Events V2 API)

- A channel is configured with the alert source's **URL + token**; title
and description templates are prefilled with the same defaults as
Jira/JSM.
- Alerts fire and auto-resolve in incident.io; the description is
markdown (incident.io renders it natively) and carries the usual deep
links — **View in SigNoz, related logs, related traces**.
- All rule labels (severity, team, custom labels) are sent as
**metadata**, so users can map them to incident.io attributes and
route/escalate on them.

Notes and decisions for the reviewer (full details in the [discussion
ticket and doc](https://github.com/SigNoz/pulse-pod/issues/171)):

- **Dedup:** one incident.io alert per notification group, keyed by the
group key hash (same identity Jira uses). A resolve targets the same
key; re-fires after resolve correctly open a fresh alert — no key
rotation needed.
- **Repeat notifications are no-ops on incident.io** (it drops duplicate
firing events) — unlike Jira, we cannot append updated values to an open
alert; operators click through to SigNoz for current values.
- **Limits:** description capped client-side under incident.io's
documented 512 KB payload limit; retries only on 429/5xx (documented
limit: 120 events/min per source).
- **Channel-level metadata:** optional key-value pairs on the channel
config, merged into every event's metadata on top of the alert's labels
(channel wins on key clash — Opsgenie precedent). Values are
template-expanded; a value that fails to expand is sent raw with a
warning logged, so delivery never breaks on a bad template.
- Upstream alertmanager ships its own basic incident.io notifier — our
config **shadows it** so the SigNoz notifier (templates, dedup,
metadata) handles delivery.
- Frontend (channel form) follows in a stacked PR.

## Issues closed by this PR

Closes SigNoz/pulse-pod#172

---------

Co-authored-by: Naman Verma <naman.verma@signoz.io>
2026-09-03 05:32:51 +00:00
70 changed files with 5602 additions and 514 deletions

View File

@@ -50,6 +50,7 @@ jobs:
- logspipelines
- passwordauthn
- preference
- quickfilter
- querierlogs
- queriertraces
- queriermetrics

View File

@@ -109,6 +109,25 @@ components:
webhook_url:
$ref: '#/components/schemas/ConfigSecretURL'
type: object
AlertmanagertypesIncidentIOReceiverConfig:
properties:
description:
type: string
http_config:
$ref: '#/components/schemas/ConfigHTTPClientConfig'
metadata:
additionalProperties:
type: string
type: object
send_resolved:
type: boolean
title:
type: string
token:
type: string
url:
type: string
type: object
AlertmanagertypesJSMOpsReceiverConfig:
properties:
api_key:
@@ -217,12 +236,12 @@ components:
- jira_configs
- required:
- jsmops_configs
- required:
- incidentio_configs
- required:
- discord_configs
- required:
- email_configs
- required:
- incidentio_configs
- required:
- pagerduty_configs
- required:
@@ -266,7 +285,7 @@ components:
type: array
incidentio_configs:
items:
$ref: '#/components/schemas/ConfigIncidentioConfig'
$ref: '#/components/schemas/AlertmanagertypesIncidentIOReceiverConfig'
type: array
jira_configs:
items:
@@ -397,7 +416,7 @@ components:
type: array
incidentio_configs:
items:
$ref: '#/components/schemas/ConfigIncidentioConfig'
$ref: '#/components/schemas/AlertmanagertypesIncidentIOReceiverConfig'
type: array
jira_configs:
items:
@@ -4071,6 +4090,18 @@ components:
nullable: true
type: object
type: object
GatewaytypesDeprecatedPostableIngestionKeyLimit:
properties:
config:
$ref: '#/components/schemas/GatewaytypesLimitConfig'
signal:
type: string
tags:
items:
type: string
nullable: true
type: array
type: object
GatewaytypesGettableCreatedIngestionKey:
properties:
id:
@@ -4214,6 +4245,8 @@ components:
properties:
config:
$ref: '#/components/schemas/GatewaytypesLimitConfig'
keyId:
type: string
signal:
type: string
tags:
@@ -4221,6 +4254,8 @@ components:
type: string
nullable: true
type: array
required:
- keyId
type: object
GatewaytypesUpdatableIngestionKeyLimit:
properties:
@@ -7838,6 +7873,48 @@ components:
- custom
- text
type: string
QuickfiltertypesSource:
enum:
- traces
- logs
- api_monitoring
- exceptions
- meter
- ai_observability
type: string
QuickfiltertypesSourceFilters:
properties:
createdAt:
format: date-time
type: string
filters:
items:
$ref: '#/components/schemas/TelemetrytypesTelemetryFieldKey'
type: array
id:
type: string
orgId:
type: string
source:
$ref: '#/components/schemas/QuickfiltertypesSource'
updatedAt:
format: date-time
type: string
required:
- id
- orgId
- source
- filters
type: object
QuickfiltertypesUpdatableQuickFilters:
properties:
filters:
items:
$ref: '#/components/schemas/TelemetrytypesTelemetryFieldKey'
type: array
required:
- filters
type: object
RenderErrorResponse:
properties:
error:
@@ -9678,6 +9755,10 @@ paths:
name: name
schema:
type: string
- in: query
name: existingQuery
schema:
type: string
responses:
"200":
content:
@@ -16142,6 +16223,63 @@ paths:
summary: Delete ingestion key for workspace
tags:
- gateway
get:
deprecated: false
description: This endpoint returns an ingestion key for the workspace
operationId: GetIngestionKey
parameters:
- in: path
name: keyId
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/GatewaytypesIngestionKey'
status:
type: string
required:
- status
- data
type: object
description: OK
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Get ingestion key for workspace
tags:
- gateway
patch:
deprecated: false
description: This endpoint updates an ingestion key for the workspace
@@ -16187,8 +16325,68 @@ paths:
tags:
- gateway
/api/v2/gateway/ingestion_keys/{keyId}/limits:
post:
get:
deprecated: false
description: This endpoint returns the ingestion limits for an ingestion key
operationId: GetIngestionKeyLimits
parameters:
- in: path
name: keyId
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
items:
$ref: '#/components/schemas/GatewaytypesLimit'
nullable: true
type: array
status:
type: string
required:
- status
- data
type: object
description: OK
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Get limits for the ingestion key
tags:
- gateway
post:
deprecated: true
description: This endpoint creates an ingestion key limit
operationId: CreateIngestionKeyLimit
parameters:
@@ -16201,7 +16399,7 @@ paths:
content:
application/json:
schema:
$ref: '#/components/schemas/GatewaytypesPostableIngestionKeyLimit'
$ref: '#/components/schemas/GatewaytypesDeprecatedPostableIngestionKeyLimit'
responses:
"201":
content:
@@ -16245,7 +16443,7 @@ paths:
- gateway
/api/v2/gateway/ingestion_keys/limits/{limitId}:
delete:
deprecated: false
deprecated: true
description: This endpoint deletes an ingestion key limit
operationId: DeleteIngestionKeyLimit
parameters:
@@ -16284,7 +16482,7 @@ paths:
tags:
- gateway
patch:
deprecated: false
deprecated: true
description: This endpoint updates an ingestion key limit
operationId: UpdateIngestionKeyLimit
parameters:
@@ -16387,6 +16585,205 @@ paths:
summary: Search ingestion keys for workspace
tags:
- gateway
/api/v2/gateway/ingestion_limits:
post:
deprecated: false
description: This endpoint creates an ingestion limit for the ingestion key
referenced by keyId
operationId: CreateIngestionLimit
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/GatewaytypesPostableIngestionKeyLimit'
responses:
"201":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/TypesIdentifiable'
status:
type: string
required:
- status
- data
type: object
description: Created
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Create ingestion limit
tags:
- gateway
/api/v2/gateway/ingestion_limits/{limitId}:
delete:
deprecated: false
description: This endpoint deletes an ingestion limit
operationId: DeleteIngestionLimit
parameters:
- in: path
name: limitId
required: true
schema:
type: string
responses:
"204":
description: No Content
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Delete ingestion limit
tags:
- gateway
get:
deprecated: false
description: This endpoint returns an ingestion limit
operationId: GetIngestionLimit
parameters:
- in: path
name: limitId
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/GatewaytypesLimit'
status:
type: string
required:
- status
- data
type: object
description: OK
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Get ingestion limit
tags:
- gateway
patch:
deprecated: false
description: This endpoint updates an ingestion limit
operationId: UpdateIngestionLimit
parameters:
- in: path
name: limitId
required: true
schema:
type: string
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/GatewaytypesUpdatableIngestionKeyLimit'
responses:
"204":
description: No Content
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Update ingestion limit
tags:
- gateway
/api/v2/healthz:
get:
operationId: Healthz
@@ -18841,6 +19238,170 @@ paths:
summary: Get query range result (v2)
tags:
- dashboard
/api/v2/quick_filters:
get:
deprecated: false
description: Returns the org's quick filters for every source, each filter as
a telemetry field key.
operationId: ListQuickFilters
responses:
"200":
content:
application/json:
schema:
properties:
data:
items:
$ref: '#/components/schemas/QuickfiltertypesSourceFilters'
type: array
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- quick-filter:list
- tokenizer:
- quick-filter:list
summary: List quick filters
tags:
- quick_filter
/api/v2/quick_filters/{source}:
get:
deprecated: false
description: Returns the org's quick filters for one source, each filter as
a telemetry field key.
operationId: GetQuickFilters
parameters:
- in: path
name: source
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/QuickfiltertypesSourceFilters'
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- quick-filter:read
- tokenizer:
- quick-filter:read
summary: Get a source's quick filters
tags:
- quick_filter
put:
deprecated: false
description: Replaces the org's quick filters for the source named in the path.
operationId: UpdateQuickFilters
parameters:
- in: path
name: source
required: true
schema:
type: string
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/QuickfiltertypesUpdatableQuickFilters'
responses:
"204":
description: No Content
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- quick-filter:update
- tokenizer:
- quick-filter:update
summary: Update quick filters
tags:
- quick_filter
/api/v2/readyz:
get:
operationId: Readyz

View File

@@ -108,6 +108,20 @@ func (provider *Provider) SearchIngestionKeysByName(ctx context.Context, orgID v
}, nil
}
func (provider *Provider) GetIngestionKey(ctx context.Context, orgID valuer.UUID, keyID string) (*gatewaytypes.IngestionKey, error) {
responseBody, err := provider.do(ctx, orgID, http.MethodGet, "/v1/workspaces/me/keys/"+keyID, nil, nil)
if err != nil {
return nil, err
}
var ingestionKey gatewaytypes.IngestionKey
if err := json.Unmarshal([]byte(gjson.GetBytes(responseBody, "data").String()), &ingestionKey); err != nil {
return nil, err
}
return &ingestionKey, nil
}
func (provider *Provider) CreateIngestionKey(ctx context.Context, orgID valuer.UUID, name string, tags []string, expiresAt time.Time) (*gatewaytypes.GettableCreatedIngestionKey, error) {
requestBody := gatewaytypes.PostableIngestionKey{
Name: name,
@@ -161,7 +175,7 @@ func (provider *Provider) DeleteIngestionKey(ctx context.Context, orgID valuer.U
}
func (provider *Provider) CreateIngestionKeyLimit(ctx context.Context, orgID valuer.UUID, keyID string, signal string, limitConfig gatewaytypes.LimitConfig, tags []string) (*gatewaytypes.GettableCreatedIngestionKeyLimit, error) {
requestBody := gatewaytypes.PostableIngestionKeyLimit{
requestBody := gatewaytypes.DeprecatedPostableIngestionKeyLimit{
Signal: signal,
Config: limitConfig,
Tags: tags,
@@ -184,6 +198,34 @@ func (provider *Provider) CreateIngestionKeyLimit(ctx context.Context, orgID val
return &createdIngestionKeyLimitResponse, nil
}
func (provider *Provider) GetIngestionKeyLimit(ctx context.Context, orgID valuer.UUID, limitID string) (*gatewaytypes.Limit, error) {
responseBody, err := provider.do(ctx, orgID, http.MethodGet, "/v1/workspaces/me/limits/"+limitID, nil, nil)
if err != nil {
return nil, err
}
var limit gatewaytypes.Limit
if err := json.Unmarshal([]byte(gjson.GetBytes(responseBody, "data").String()), &limit); err != nil {
return nil, err
}
return &limit, nil
}
func (provider *Provider) GetIngestionKeyLimits(ctx context.Context, orgID valuer.UUID, keyID string) ([]gatewaytypes.Limit, error) {
responseBody, err := provider.do(ctx, orgID, http.MethodGet, "/v1/workspaces/me/keys/"+keyID+"/limits", nil, nil)
if err != nil {
return nil, err
}
var limits []gatewaytypes.Limit
if err := json.Unmarshal([]byte(gjson.GetBytes(responseBody, "data").String()), &limits); err != nil {
return nil, err
}
return limits, nil
}
func (provider *Provider) UpdateIngestionKeyLimit(ctx context.Context, orgID valuer.UUID, limitID string, limitConfig gatewaytypes.LimitConfig, tags []string) error {
requestBody := gatewaytypes.UpdatableIngestionKeyLimit{
Config: limitConfig,

View File

@@ -75,6 +75,23 @@
"field_jsmops_tags": "Tags",
"placeholder_jsmops_tags": "Type a tag and press Enter",
"help_jsmops_tags": "Tags added to every alert.",
"incidentio_tip": "Create an HTTP alert source in incident.io (On-call \u2192 Alert routing \u2192 Sources) and paste its URL and token below.",
"incidentio_tip_link": "Learn how",
"field_incidentio_url": "Alert source URL",
"help_incidentio_url": "The alert events URL from the source's setup page, e.g. https://api.incident.io/v2/alert_events/http/<source_config_id>.",
"field_incidentio_token": "Token",
"help_incidentio_token": "The alert source's secret token, from the same setup page.",
"field_incidentio_title": "Title",
"help_incidentio_title": "Template for the alert title. Kept stable while the alert fires \u2014 incident.io ignores content updates on repeat events.",
"field_incidentio_description": "Description",
"help_incidentio_description": "Template for the alert description. Markdown, rendered natively by incident.io.",
"incidentio_required_fields": "Alert source URL and token are required",
"incidentio_url_invalid": "URL must be an incident.io alert events URL (https://api.incident.io/v2/alert_events/http/<source_config_id>)",
"field_incidentio_metadata": "Additional metadata",
"help_incidentio_metadata": "Key-value pairs added to every alert's metadata, on top of the alert's labels (these win on a key clash). Values may use templates, e.g. {{ .CommonLabels.severity }}.",
"placeholder_incidentio_metadata_key": "Key",
"placeholder_incidentio_metadata_value": "Value",
"button_incidentio_add_metadata": "Add metadata",
"field_slack_recipient": "Recipient",
"field_slack_title": "Title",

View File

@@ -75,6 +75,23 @@
"field_jsmops_tags": "Tags",
"placeholder_jsmops_tags": "Type a tag and press Enter",
"help_jsmops_tags": "Tags added to every alert.",
"incidentio_tip": "Create an HTTP alert source in incident.io (On-call \u2192 Alert routing \u2192 Sources) and paste its URL and token below.",
"incidentio_tip_link": "Learn how",
"field_incidentio_url": "Alert source URL",
"help_incidentio_url": "The alert events URL from the source's setup page, e.g. https://api.incident.io/v2/alert_events/http/<source_config_id>.",
"field_incidentio_token": "Token",
"help_incidentio_token": "The alert source's secret token, from the same setup page.",
"field_incidentio_title": "Title",
"help_incidentio_title": "Template for the alert title. Kept stable while the alert fires \u2014 incident.io ignores content updates on repeat events.",
"field_incidentio_description": "Description",
"help_incidentio_description": "Template for the alert description. Markdown, rendered natively by incident.io.",
"incidentio_required_fields": "Alert source URL and token are required",
"incidentio_url_invalid": "URL must be an incident.io alert events URL (https://api.incident.io/v2/alert_events/http/<source_config_id>)",
"field_incidentio_metadata": "Additional metadata",
"help_incidentio_metadata": "Key-value pairs added to every alert's metadata, on top of the alert's labels (these win on a key clash). Values may use templates, e.g. {{ .CommonLabels.severity }}.",
"placeholder_incidentio_metadata_key": "Key",
"placeholder_incidentio_metadata_value": "Value",
"button_incidentio_add_metadata": "Add metadata",
"field_slack_recipient": "Recipient",
"field_slack_title": "Title",
"field_slack_description": "Description",

View File

@@ -21,18 +21,28 @@ import type {
CreateIngestionKey201,
CreateIngestionKeyLimit201,
CreateIngestionKeyLimitPathParameters,
CreateIngestionLimit201,
DeleteIngestionKeyLimitPathParameters,
DeleteIngestionKeyPathParameters,
DeleteIngestionLimitPathParameters,
GatewaytypesDeprecatedPostableIngestionKeyLimitDTO,
GatewaytypesPostableIngestionKeyDTO,
GatewaytypesPostableIngestionKeyLimitDTO,
GatewaytypesUpdatableIngestionKeyLimitDTO,
GetIngestionKey200,
GetIngestionKeyLimits200,
GetIngestionKeyLimitsPathParameters,
GetIngestionKeyPathParameters,
GetIngestionKeys200,
GetIngestionKeysParams,
GetIngestionLimit200,
GetIngestionLimitPathParameters,
RenderErrorResponseDTO,
SearchIngestionKeys200,
SearchIngestionKeysParams,
UpdateIngestionKeyLimitPathParameters,
UpdateIngestionKeyPathParameters,
UpdateIngestionLimitPathParameters,
} from '../sigNoz.schemas';
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
@@ -300,6 +310,108 @@ export const useDeleteIngestionKey = <
> => {
return useMutation(getDeleteIngestionKeyMutationOptions(options));
};
/**
* This endpoint returns an ingestion key for the workspace
* @summary Get ingestion key for workspace
*/
export const getIngestionKey = (
{ keyId }: GetIngestionKeyPathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<GetIngestionKey200>({
url: `/api/v2/gateway/ingestion_keys/${keyId}`,
method: 'GET',
signal,
});
};
export const getGetIngestionKeyQueryKey = ({
keyId,
}: GetIngestionKeyPathParameters) => {
return [`/api/v2/gateway/ingestion_keys/${keyId}`] as const;
};
export const getGetIngestionKeyQueryOptions = <
TData = Awaited<ReturnType<typeof getIngestionKey>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ keyId }: GetIngestionKeyPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getIngestionKey>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getGetIngestionKeyQueryKey({ keyId });
const queryFn: QueryFunction<Awaited<ReturnType<typeof getIngestionKey>>> = ({
signal,
}) => getIngestionKey({ keyId }, signal);
return {
queryKey,
queryFn,
enabled: !!keyId,
...queryOptions,
} as UseQueryOptions<
Awaited<ReturnType<typeof getIngestionKey>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type GetIngestionKeyQueryResult = NonNullable<
Awaited<ReturnType<typeof getIngestionKey>>
>;
export type GetIngestionKeyQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Get ingestion key for workspace
*/
export function useGetIngestionKey<
TData = Awaited<ReturnType<typeof getIngestionKey>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ keyId }: GetIngestionKeyPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getIngestionKey>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getGetIngestionKeyQueryOptions({ keyId }, options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return { ...query, queryKey: queryOptions.queryKey };
}
/**
* @summary Get ingestion key for workspace
*/
export const invalidateGetIngestionKey = async (
queryClient: QueryClient,
{ keyId }: GetIngestionKeyPathParameters,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getGetIngestionKeyQueryKey({ keyId }) },
options,
);
return queryClient;
};
/**
* This endpoint updates an ingestion key for the workspace
* @summary Update ingestion key for workspace
@@ -399,20 +511,123 @@ export const useUpdateIngestionKey = <
> => {
return useMutation(getUpdateIngestionKeyMutationOptions(options));
};
/**
* This endpoint returns the ingestion limits for an ingestion key
* @summary Get limits for the ingestion key
*/
export const getIngestionKeyLimits = (
{ keyId }: GetIngestionKeyLimitsPathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<GetIngestionKeyLimits200>({
url: `/api/v2/gateway/ingestion_keys/${keyId}/limits`,
method: 'GET',
signal,
});
};
export const getGetIngestionKeyLimitsQueryKey = ({
keyId,
}: GetIngestionKeyLimitsPathParameters) => {
return [`/api/v2/gateway/ingestion_keys/${keyId}/limits`] as const;
};
export const getGetIngestionKeyLimitsQueryOptions = <
TData = Awaited<ReturnType<typeof getIngestionKeyLimits>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ keyId }: GetIngestionKeyLimitsPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getIngestionKeyLimits>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getGetIngestionKeyLimitsQueryKey({ keyId });
const queryFn: QueryFunction<
Awaited<ReturnType<typeof getIngestionKeyLimits>>
> = ({ signal }) => getIngestionKeyLimits({ keyId }, signal);
return {
queryKey,
queryFn,
enabled: !!keyId,
...queryOptions,
} as UseQueryOptions<
Awaited<ReturnType<typeof getIngestionKeyLimits>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type GetIngestionKeyLimitsQueryResult = NonNullable<
Awaited<ReturnType<typeof getIngestionKeyLimits>>
>;
export type GetIngestionKeyLimitsQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Get limits for the ingestion key
*/
export function useGetIngestionKeyLimits<
TData = Awaited<ReturnType<typeof getIngestionKeyLimits>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ keyId }: GetIngestionKeyLimitsPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getIngestionKeyLimits>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getGetIngestionKeyLimitsQueryOptions({ keyId }, options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return { ...query, queryKey: queryOptions.queryKey };
}
/**
* @summary Get limits for the ingestion key
*/
export const invalidateGetIngestionKeyLimits = async (
queryClient: QueryClient,
{ keyId }: GetIngestionKeyLimitsPathParameters,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getGetIngestionKeyLimitsQueryKey({ keyId }) },
options,
);
return queryClient;
};
/**
* This endpoint creates an ingestion key limit
* @deprecated
* @summary Create limit for the ingestion key
*/
export const createIngestionKeyLimit = (
{ keyId }: CreateIngestionKeyLimitPathParameters,
gatewaytypesPostableIngestionKeyLimitDTO?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>,
gatewaytypesDeprecatedPostableIngestionKeyLimitDTO?: BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<CreateIngestionKeyLimit201>({
url: `/api/v2/gateway/ingestion_keys/${keyId}/limits`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: gatewaytypesPostableIngestionKeyLimitDTO,
data: gatewaytypesDeprecatedPostableIngestionKeyLimitDTO,
signal,
});
};
@@ -426,7 +641,7 @@ export const getCreateIngestionKeyLimitMutationOptions = <
TError,
{
pathParams: CreateIngestionKeyLimitPathParameters;
data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>;
data?: BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>;
},
TContext
>;
@@ -435,7 +650,7 @@ export const getCreateIngestionKeyLimitMutationOptions = <
TError,
{
pathParams: CreateIngestionKeyLimitPathParameters;
data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>;
data?: BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>;
},
TContext
> => {
@@ -452,7 +667,7 @@ export const getCreateIngestionKeyLimitMutationOptions = <
Awaited<ReturnType<typeof createIngestionKeyLimit>>,
{
pathParams: CreateIngestionKeyLimitPathParameters;
data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>;
data?: BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>;
}
> = (props) => {
const { pathParams, data } = props ?? {};
@@ -467,12 +682,13 @@ export type CreateIngestionKeyLimitMutationResult = NonNullable<
Awaited<ReturnType<typeof createIngestionKeyLimit>>
>;
export type CreateIngestionKeyLimitMutationBody =
| BodyType<GatewaytypesPostableIngestionKeyLimitDTO>
| BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>
| undefined;
export type CreateIngestionKeyLimitMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @deprecated
* @summary Create limit for the ingestion key
*/
export const useCreateIngestionKeyLimit = <
@@ -484,7 +700,7 @@ export const useCreateIngestionKeyLimit = <
TError,
{
pathParams: CreateIngestionKeyLimitPathParameters;
data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>;
data?: BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>;
},
TContext
>;
@@ -493,7 +709,7 @@ export const useCreateIngestionKeyLimit = <
TError,
{
pathParams: CreateIngestionKeyLimitPathParameters;
data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>;
data?: BodyType<GatewaytypesDeprecatedPostableIngestionKeyLimitDTO>;
},
TContext
> => {
@@ -501,6 +717,7 @@ export const useCreateIngestionKeyLimit = <
};
/**
* This endpoint deletes an ingestion key limit
* @deprecated
* @summary Delete limit for the ingestion key
*/
export const deleteIngestionKeyLimit = (
@@ -559,6 +776,7 @@ export type DeleteIngestionKeyLimitMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @deprecated
* @summary Delete limit for the ingestion key
*/
export const useDeleteIngestionKeyLimit = <
@@ -581,6 +799,7 @@ export const useDeleteIngestionKeyLimit = <
};
/**
* This endpoint updates an ingestion key limit
* @deprecated
* @summary Update limit for the ingestion key
*/
export const updateIngestionKeyLimit = (
@@ -653,6 +872,7 @@ export type UpdateIngestionKeyLimitMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @deprecated
* @summary Update limit for the ingestion key
*/
export const useUpdateIngestionKeyLimit = <
@@ -779,3 +999,370 @@ export const invalidateSearchIngestionKeys = async (
return queryClient;
};
/**
* This endpoint creates an ingestion limit for the ingestion key referenced by keyId
* @summary Create ingestion limit
*/
export const createIngestionLimit = (
gatewaytypesPostableIngestionKeyLimitDTO?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<CreateIngestionLimit201>({
url: `/api/v2/gateway/ingestion_limits`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: gatewaytypesPostableIngestionKeyLimitDTO,
signal,
});
};
export const getCreateIngestionLimitMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof createIngestionLimit>>,
TError,
{ data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof createIngestionLimit>>,
TError,
{ data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO> },
TContext
> => {
const mutationKey = ['createIngestionLimit'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
options.mutation.mutationKey
? options
: { ...options, mutation: { ...options.mutation, mutationKey } }
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof createIngestionLimit>>,
{ data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO> }
> = (props) => {
const { data } = props ?? {};
return createIngestionLimit(data);
};
return { mutationFn, ...mutationOptions };
};
export type CreateIngestionLimitMutationResult = NonNullable<
Awaited<ReturnType<typeof createIngestionLimit>>
>;
export type CreateIngestionLimitMutationBody =
| BodyType<GatewaytypesPostableIngestionKeyLimitDTO>
| undefined;
export type CreateIngestionLimitMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @summary Create ingestion limit
*/
export const useCreateIngestionLimit = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof createIngestionLimit>>,
TError,
{ data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof createIngestionLimit>>,
TError,
{ data?: BodyType<GatewaytypesPostableIngestionKeyLimitDTO> },
TContext
> => {
return useMutation(getCreateIngestionLimitMutationOptions(options));
};
/**
* This endpoint deletes an ingestion limit
* @summary Delete ingestion limit
*/
export const deleteIngestionLimit = (
{ limitId }: DeleteIngestionLimitPathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v2/gateway/ingestion_limits/${limitId}`,
method: 'DELETE',
signal,
});
};
export const getDeleteIngestionLimitMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof deleteIngestionLimit>>,
TError,
{ pathParams: DeleteIngestionLimitPathParameters },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof deleteIngestionLimit>>,
TError,
{ pathParams: DeleteIngestionLimitPathParameters },
TContext
> => {
const mutationKey = ['deleteIngestionLimit'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
options.mutation.mutationKey
? options
: { ...options, mutation: { ...options.mutation, mutationKey } }
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof deleteIngestionLimit>>,
{ pathParams: DeleteIngestionLimitPathParameters }
> = (props) => {
const { pathParams } = props ?? {};
return deleteIngestionLimit(pathParams);
};
return { mutationFn, ...mutationOptions };
};
export type DeleteIngestionLimitMutationResult = NonNullable<
Awaited<ReturnType<typeof deleteIngestionLimit>>
>;
export type DeleteIngestionLimitMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @summary Delete ingestion limit
*/
export const useDeleteIngestionLimit = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof deleteIngestionLimit>>,
TError,
{ pathParams: DeleteIngestionLimitPathParameters },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof deleteIngestionLimit>>,
TError,
{ pathParams: DeleteIngestionLimitPathParameters },
TContext
> => {
return useMutation(getDeleteIngestionLimitMutationOptions(options));
};
/**
* This endpoint returns an ingestion limit
* @summary Get ingestion limit
*/
export const getIngestionLimit = (
{ limitId }: GetIngestionLimitPathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<GetIngestionLimit200>({
url: `/api/v2/gateway/ingestion_limits/${limitId}`,
method: 'GET',
signal,
});
};
export const getGetIngestionLimitQueryKey = ({
limitId,
}: GetIngestionLimitPathParameters) => {
return [`/api/v2/gateway/ingestion_limits/${limitId}`] as const;
};
export const getGetIngestionLimitQueryOptions = <
TData = Awaited<ReturnType<typeof getIngestionLimit>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ limitId }: GetIngestionLimitPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getIngestionLimit>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getGetIngestionLimitQueryKey({ limitId });
const queryFn: QueryFunction<
Awaited<ReturnType<typeof getIngestionLimit>>
> = ({ signal }) => getIngestionLimit({ limitId }, signal);
return {
queryKey,
queryFn,
enabled: !!limitId,
...queryOptions,
} as UseQueryOptions<
Awaited<ReturnType<typeof getIngestionLimit>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type GetIngestionLimitQueryResult = NonNullable<
Awaited<ReturnType<typeof getIngestionLimit>>
>;
export type GetIngestionLimitQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Get ingestion limit
*/
export function useGetIngestionLimit<
TData = Awaited<ReturnType<typeof getIngestionLimit>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ limitId }: GetIngestionLimitPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getIngestionLimit>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getGetIngestionLimitQueryOptions({ limitId }, options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return { ...query, queryKey: queryOptions.queryKey };
}
/**
* @summary Get ingestion limit
*/
export const invalidateGetIngestionLimit = async (
queryClient: QueryClient,
{ limitId }: GetIngestionLimitPathParameters,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getGetIngestionLimitQueryKey({ limitId }) },
options,
);
return queryClient;
};
/**
* This endpoint updates an ingestion limit
* @summary Update ingestion limit
*/
export const updateIngestionLimit = (
{ limitId }: UpdateIngestionLimitPathParameters,
gatewaytypesUpdatableIngestionKeyLimitDTO?: BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v2/gateway/ingestion_limits/${limitId}`,
method: 'PATCH',
headers: { 'Content-Type': 'application/json' },
data: gatewaytypesUpdatableIngestionKeyLimitDTO,
signal,
});
};
export const getUpdateIngestionLimitMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof updateIngestionLimit>>,
TError,
{
pathParams: UpdateIngestionLimitPathParameters;
data?: BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>;
},
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof updateIngestionLimit>>,
TError,
{
pathParams: UpdateIngestionLimitPathParameters;
data?: BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>;
},
TContext
> => {
const mutationKey = ['updateIngestionLimit'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
options.mutation.mutationKey
? options
: { ...options, mutation: { ...options.mutation, mutationKey } }
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof updateIngestionLimit>>,
{
pathParams: UpdateIngestionLimitPathParameters;
data?: BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>;
}
> = (props) => {
const { pathParams, data } = props ?? {};
return updateIngestionLimit(pathParams, data);
};
return { mutationFn, ...mutationOptions };
};
export type UpdateIngestionLimitMutationResult = NonNullable<
Awaited<ReturnType<typeof updateIngestionLimit>>
>;
export type UpdateIngestionLimitMutationBody =
| BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>
| undefined;
export type UpdateIngestionLimitMutationError =
ErrorType<RenderErrorResponseDTO>;
/**
* @summary Update ingestion limit
*/
export const useUpdateIngestionLimit = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof updateIngestionLimit>>,
TError,
{
pathParams: UpdateIngestionLimitPathParameters;
data?: BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>;
},
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof updateIngestionLimit>>,
TError,
{
pathParams: UpdateIngestionLimitPathParameters;
data?: BodyType<GatewaytypesUpdatableIngestionKeyLimitDTO>;
},
TContext
> => {
return useMutation(getUpdateIngestionLimitMutationOptions(options));
};

View File

@@ -0,0 +1,316 @@
/**
* ! Do not edit manually
* * The file has been auto-generated using Orval for SigNoz
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
import { useMutation, useQuery } from 'react-query';
import type {
InvalidateOptions,
MutationFunction,
QueryClient,
QueryFunction,
QueryKey,
UseMutationOptions,
UseMutationResult,
UseQueryOptions,
UseQueryResult,
} from 'react-query';
import type {
GetQuickFilters200,
GetQuickFiltersPathParameters,
ListQuickFilters200,
QuickfiltertypesUpdatableQuickFiltersDTO,
RenderErrorResponseDTO,
UpdateQuickFiltersPathParameters,
} from '../sigNoz.schemas';
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
/**
* Returns the org's quick filters for every source, each filter as a telemetry field key.
* @summary List quick filters
*/
export const listQuickFilters = (signal?: AbortSignal) => {
return GeneratedAPIInstance<ListQuickFilters200>({
url: `/api/v2/quick_filters`,
method: 'GET',
signal,
});
};
export const getListQuickFiltersQueryKey = () => {
return [`/api/v2/quick_filters`] as const;
};
export const getListQuickFiltersQueryOptions = <
TData = Awaited<ReturnType<typeof listQuickFilters>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listQuickFilters>>,
TError,
TData
>;
}) => {
const { query: queryOptions } = options ?? {};
const queryKey = queryOptions?.queryKey ?? getListQuickFiltersQueryKey();
const queryFn: QueryFunction<Awaited<ReturnType<typeof listQuickFilters>>> = ({
signal,
}) => listQuickFilters(signal);
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
Awaited<ReturnType<typeof listQuickFilters>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type ListQuickFiltersQueryResult = NonNullable<
Awaited<ReturnType<typeof listQuickFilters>>
>;
export type ListQuickFiltersQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary List quick filters
*/
export function useListQuickFilters<
TData = Awaited<ReturnType<typeof listQuickFilters>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listQuickFilters>>,
TError,
TData
>;
}): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getListQuickFiltersQueryOptions(options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return { ...query, queryKey: queryOptions.queryKey };
}
/**
* @summary List quick filters
*/
export const invalidateListQuickFilters = async (
queryClient: QueryClient,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getListQuickFiltersQueryKey() },
options,
);
return queryClient;
};
/**
* Returns the org's quick filters for one source, each filter as a telemetry field key.
* @summary Get a source's quick filters
*/
export const getQuickFilters = (
{ source }: GetQuickFiltersPathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<GetQuickFilters200>({
url: `/api/v2/quick_filters/${source}`,
method: 'GET',
signal,
});
};
export const getGetQuickFiltersQueryKey = ({
source,
}: GetQuickFiltersPathParameters) => {
return [`/api/v2/quick_filters/${source}`] as const;
};
export const getGetQuickFiltersQueryOptions = <
TData = Awaited<ReturnType<typeof getQuickFilters>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ source }: GetQuickFiltersPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getQuickFilters>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getGetQuickFiltersQueryKey({ source });
const queryFn: QueryFunction<Awaited<ReturnType<typeof getQuickFilters>>> = ({
signal,
}) => getQuickFilters({ source }, signal);
return {
queryKey,
queryFn,
enabled: !!source,
...queryOptions,
} as UseQueryOptions<
Awaited<ReturnType<typeof getQuickFilters>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type GetQuickFiltersQueryResult = NonNullable<
Awaited<ReturnType<typeof getQuickFilters>>
>;
export type GetQuickFiltersQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Get a source's quick filters
*/
export function useGetQuickFilters<
TData = Awaited<ReturnType<typeof getQuickFilters>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ source }: GetQuickFiltersPathParameters,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getQuickFilters>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getGetQuickFiltersQueryOptions({ source }, options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return { ...query, queryKey: queryOptions.queryKey };
}
/**
* @summary Get a source's quick filters
*/
export const invalidateGetQuickFilters = async (
queryClient: QueryClient,
{ source }: GetQuickFiltersPathParameters,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getGetQuickFiltersQueryKey({ source }) },
options,
);
return queryClient;
};
/**
* Replaces the org's quick filters for the source named in the path.
* @summary Update quick filters
*/
export const updateQuickFilters = (
{ source }: UpdateQuickFiltersPathParameters,
quickfiltertypesUpdatableQuickFiltersDTO?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v2/quick_filters/${source}`,
method: 'PUT',
headers: { 'Content-Type': 'application/json' },
data: quickfiltertypesUpdatableQuickFiltersDTO,
signal,
});
};
export const getUpdateQuickFiltersMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof updateQuickFilters>>,
TError,
{
pathParams: UpdateQuickFiltersPathParameters;
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
},
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof updateQuickFilters>>,
TError,
{
pathParams: UpdateQuickFiltersPathParameters;
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
},
TContext
> => {
const mutationKey = ['updateQuickFilters'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
options.mutation.mutationKey
? options
: { ...options, mutation: { ...options.mutation, mutationKey } }
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof updateQuickFilters>>,
{
pathParams: UpdateQuickFiltersPathParameters;
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
}
> = (props) => {
const { pathParams, data } = props ?? {};
return updateQuickFilters(pathParams, data);
};
return { mutationFn, ...mutationOptions };
};
export type UpdateQuickFiltersMutationResult = NonNullable<
Awaited<ReturnType<typeof updateQuickFilters>>
>;
export type UpdateQuickFiltersMutationBody =
| BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>
| undefined;
export type UpdateQuickFiltersMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Update quick filters
*/
export const useUpdateQuickFilters = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof updateQuickFilters>>,
TError,
{
pathParams: UpdateQuickFiltersPathParameters;
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
},
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof updateQuickFilters>>,
TError,
{
pathParams: UpdateQuickFiltersPathParameters;
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
},
TContext
> => {
return useMutation(getUpdateQuickFiltersMutationOptions(options));
};

View File

@@ -385,6 +385,38 @@ export interface AlertmanagertypesGoogleChatReceiverConfigDTO {
webhook_url?: ConfigSecretURLDTO;
}
export type AlertmanagertypesIncidentIOReceiverConfigDTOMetadata = {
[key: string]: string;
};
export interface AlertmanagertypesIncidentIOReceiverConfigDTO {
/**
* @type string
*/
description?: string;
http_config?: ConfigHTTPClientConfigDTO;
/**
* @type object
*/
metadata?: AlertmanagertypesIncidentIOReceiverConfigDTOMetadata;
/**
* @type boolean
*/
send_resolved?: boolean;
/**
* @type string
*/
title?: string;
/**
* @type string
*/
token?: string;
/**
* @type string
*/
url?: string;
}
export interface AlertmanagertypesJSMOpsReceiverConfigDTO {
/**
* @type string
@@ -685,39 +717,6 @@ export interface ConfigEmailConfigDTO {
to?: string;
}
export type TimeDurationDTO = number;
export interface ConfigURLType2DTO {
[key: string]: unknown;
}
export interface ConfigIncidentioConfigDTO {
/**
* @type string
*/
alert_source_token?: string;
/**
* @type string
*/
alert_source_token_file?: string;
http_config?: ConfigHTTPClientConfigDTO;
/**
* @type integer
* @minimum 0
*/
max_alerts?: number;
/**
* @type boolean
*/
send_resolved?: boolean;
timeout?: TimeDurationDTO;
url?: ConfigURLType2DTO;
/**
* @type string
*/
url_file?: string;
}
export interface ConfigMattermostFieldDTO {
/**
* @type boolean,null
@@ -903,6 +902,10 @@ export interface ConfigMSTeamsV2ConfigDTO {
webhook_url_file?: string;
}
export interface ConfigURLType2DTO {
[key: string]: unknown;
}
export interface ConfigOpsGenieConfigResponderDTO {
/**
* @type string
@@ -1011,6 +1014,8 @@ export interface ConfigPagerdutyLinkDTO {
text?: string;
}
export type TimeDurationDTO = number;
export type ConfigPagerdutyConfigDTODetails = { [key: string]: unknown };
export interface ConfigPagerdutyConfigDTO {
@@ -1672,7 +1677,7 @@ export type AlertmanagertypesPostableChannelDTO = unknown & {
/**
* @type array
*/
incidentio_configs?: ConfigIncidentioConfigDTO[];
incidentio_configs?: AlertmanagertypesIncidentIOReceiverConfigDTO[];
/**
* @type array
*/
@@ -1803,7 +1808,7 @@ export interface AlertmanagertypesReceiverDTO {
/**
* @type array
*/
incidentio_configs?: ConfigIncidentioConfigDTO[];
incidentio_configs?: AlertmanagertypesIncidentIOReceiverConfigDTO[];
/**
* @type array
*/
@@ -3300,6 +3305,33 @@ export interface CommonJSONRefDTO {
$ref?: string;
}
export interface ConfigIncidentioConfigDTO {
/**
* @type string
*/
alert_source_token?: string;
/**
* @type string
*/
alert_source_token_file?: string;
http_config?: ConfigHTTPClientConfigDTO;
/**
* @type integer
* @minimum 0
*/
max_alerts?: number;
/**
* @type boolean
*/
send_resolved?: boolean;
timeout?: TimeDurationDTO;
url?: ConfigURLType2DTO;
/**
* @type string
*/
url_file?: string;
}
export type ConfigJiraConfigDTOCustomFields = { [key: string]: unknown };
export interface ConfigJiraFieldConfigDTO {
@@ -5461,6 +5493,34 @@ export interface FeaturetypesGettableFeatureDTO {
variants?: FeaturetypesGettableFeatureDTOVariants;
}
export interface GatewaytypesLimitValueDTO {
/**
* @type integer,null
*/
count?: number | null;
/**
* @type integer,null
*/
size?: number | null;
}
export interface GatewaytypesLimitConfigDTO {
day?: GatewaytypesLimitValueDTO;
second?: GatewaytypesLimitValueDTO;
}
export interface GatewaytypesDeprecatedPostableIngestionKeyLimitDTO {
config?: GatewaytypesLimitConfigDTO;
/**
* @type string
*/
signal?: string;
/**
* @type array,null
*/
tags?: string[] | null;
}
export interface GatewaytypesGettableCreatedIngestionKeyDTO {
/**
* @type string
@@ -5498,22 +5558,6 @@ export interface GatewaytypesPaginationDTO {
total?: number;
}
export interface GatewaytypesLimitValueDTO {
/**
* @type integer,null
*/
count?: number | null;
/**
* @type integer,null
*/
size?: number | null;
}
export interface GatewaytypesLimitConfigDTO {
day?: GatewaytypesLimitValueDTO;
second?: GatewaytypesLimitValueDTO;
}
export interface GatewaytypesLimitMetricValueDTO {
/**
* @type integer
@@ -5631,6 +5675,10 @@ export interface GatewaytypesPostableIngestionKeyDTO {
export interface GatewaytypesPostableIngestionKeyLimitDTO {
config?: GatewaytypesLimitConfigDTO;
/**
* @type string
*/
keyId: string;
/**
* @type string
*/
@@ -8978,6 +9026,47 @@ export enum Querybuildertypesv5QueryTypeDTO {
clickhouse_sql = 'clickhouse_sql',
promql = 'promql',
}
export enum QuickfiltertypesSourceDTO {
traces = 'traces',
logs = 'logs',
api_monitoring = 'api_monitoring',
exceptions = 'exceptions',
meter = 'meter',
ai_observability = 'ai_observability',
}
export interface QuickfiltertypesSourceFiltersDTO {
/**
* @type string
* @format date-time
*/
createdAt?: string;
/**
* @type array
*/
filters: TelemetrytypesTelemetryFieldKeyDTO[];
/**
* @type string
*/
id: string;
/**
* @type string
*/
orgId: string;
source: QuickfiltertypesSourceDTO;
/**
* @type string
* @format date-time
*/
updatedAt?: string;
}
export interface QuickfiltertypesUpdatableQuickFiltersDTO {
/**
* @type array
*/
filters: TelemetrytypesTelemetryFieldKeyDTO[];
}
export interface RenderErrorResponseDTO {
error: ErrorsJSONDTO;
/**
@@ -10790,6 +10879,11 @@ export type GetAIObservabilityFieldsValuesParams = {
* @description undefined
*/
name?: string;
/**
* @type string
* @description undefined
*/
existingQuery?: string;
};
export type GetAIObservabilityFieldsValues200 = {
@@ -11908,9 +12002,34 @@ export type CreateIngestionKey201 = {
export type DeleteIngestionKeyPathParameters = {
keyId: string;
};
export type GetIngestionKeyPathParameters = {
keyId: string;
};
export type GetIngestionKey200 = {
data: GatewaytypesIngestionKeyDTO;
/**
* @type string
*/
status: string;
};
export type UpdateIngestionKeyPathParameters = {
keyId: string;
};
export type GetIngestionKeyLimitsPathParameters = {
keyId: string;
};
export type GetIngestionKeyLimits200 = {
/**
* @type array,null
*/
data: GatewaytypesLimitDTO[] | null;
/**
* @type string
*/
status: string;
};
export type CreateIngestionKeyLimitPathParameters = {
keyId: string;
};
@@ -11954,6 +12073,31 @@ export type SearchIngestionKeys200 = {
status: string;
};
export type CreateIngestionLimit201 = {
data: TypesIdentifiableDTO;
/**
* @type string
*/
status: string;
};
export type DeleteIngestionLimitPathParameters = {
limitId: string;
};
export type GetIngestionLimitPathParameters = {
limitId: string;
};
export type GetIngestionLimit200 = {
data: GatewaytypesLimitDTO;
/**
* @type string
*/
status: string;
};
export type UpdateIngestionLimitPathParameters = {
limitId: string;
};
export type Healthz200 = {
data: FactoryResponseDTO;
/**
@@ -12379,6 +12523,31 @@ export type GetPublicDashboardPanelQueryRangeV2200 = {
status: string;
};
export type ListQuickFilters200 = {
/**
* @type array
*/
data: QuickfiltertypesSourceFiltersDTO[];
/**
* @type string
*/
status: string;
};
export type GetQuickFiltersPathParameters = {
source: string;
};
export type GetQuickFilters200 = {
data: QuickfiltertypesSourceFiltersDTO;
/**
* @type string
*/
status: string;
};
export type UpdateQuickFiltersPathParameters = {
source: string;
};
export type Readyz200 = {
data: FactoryResponseDTO;
/**

View File

@@ -2,6 +2,7 @@ import CreateAlertChannels from 'container/CreateAlertChannels';
import { ChannelType } from 'container/CreateAlertChannels/config';
import {
GoogleChatInitialConfig,
IncidentIOInitialConfig,
JiraInitialConfig,
JsmOpsInitialConfig,
} from 'container/CreateAlertChannels/defaults';
@@ -585,7 +586,7 @@ describe('Create Alert Channel', () => {
description: 'jira_site_invalid',
}),
);
});
}, 15000);
it('Should send a jira_configs payload with basic auth', async () => {
let requestBody: unknown;
@@ -737,6 +738,122 @@ describe('Create Alert Channel', () => {
});
});
});
describe('incident.io', () => {
const incidentIOURL =
'https://api.incident.io/v2/alert_events/http/01M0D1JNVBGBGVTWX053EM12XV';
beforeEach(() => {
render(<CreateAlertChannels preType={ChannelType.IncidentIO} />);
});
it('Should display the URL and token fields with the docs tip', () => {
testLabelInputAndHelpValue({
labelText: 'field_incidentio_url',
testId: 'incidentio-url-textbox',
});
testLabelInputAndHelpValue({
labelText: 'field_incidentio_token',
testId: 'incidentio-token-textbox',
});
expect(screen.getByTestId('incidentio-tip')).toBeInTheDocument();
expect(
screen.getByRole('link', { name: 'incidentio_tip_link' }),
).toHaveAttribute(
'href',
'https://signoz.io/docs/alerts-management/notification-channel/incidentio/',
);
});
it('Should block save when the URL or token is missing', async () => {
const user = userEvent.setup();
await user.type(
screen.getByTestId('channel-name-textbox'),
'incidentio-channel',
);
await user.click(screen.getByTestId('save-channel-button'));
await waitFor(() =>
expect(errorNotification).toHaveBeenCalledWith({
message: 'Error',
description: 'incidentio_required_fields',
}),
);
});
it('Should display an error when the URL is not an alert events URL', async () => {
const user = userEvent.setup();
await user.type(
screen.getByTestId('channel-name-textbox'),
'incidentio-channel',
);
await user.type(
screen.getByTestId('incidentio-url-textbox'),
'https://api.incident.io/v2/incidents',
);
await user.type(screen.getByTestId('incidentio-token-textbox'), 'tok-abc');
await user.click(screen.getByTestId('save-channel-button'));
await waitFor(() =>
expect(errorNotification).toHaveBeenCalledWith({
message: 'Error',
description: 'incidentio_url_invalid',
}),
);
}, 15000);
it('Should send an incidentio_configs payload with prefilled defaults', async () => {
let requestBody: unknown;
server.use(
rest.post('http://localhost/api/v1/channels', async (req, res, ctx) => {
requestBody = await req.json();
return res(
ctx.status(201),
ctx.json({ status: 'success', data: 'channel created' }),
);
}),
);
const user = userEvent.setup();
await user.type(
screen.getByTestId('channel-name-textbox'),
'incidentio-channel',
);
await user.type(
screen.getByTestId('incidentio-url-textbox'),
incidentIOURL,
);
await user.type(screen.getByTestId('incidentio-token-textbox'), 'tok-abc');
await user.click(screen.getByTestId('incidentio-metadata-add'));
await user.type(screen.getByTestId('incidentio-metadata-key-0'), 'team');
await user.type(screen.getByTestId('incidentio-metadata-value-0'), 'core');
await user.click(screen.getByTestId('save-channel-button'));
await waitFor(() =>
expect(successNotification).toHaveBeenCalledWith({
message: 'Success',
description: 'channel_creation_done',
}),
);
expect(requestBody).toStrictEqual({
name: 'incidentio-channel',
incidentio_configs: [
{
url: incidentIOURL,
token: 'tok-abc',
send_resolved: true,
title: IncidentIOInitialConfig.title,
description: IncidentIOInitialConfig.description,
metadata: { team: 'core' },
},
],
});
}, 15000);
});
describe('Changing the channel type', () => {
async function selectType(
user: ReturnType<typeof userEvent.setup>,

View File

@@ -90,6 +90,40 @@ describe('EditAlertChannels save', () => {
await waitFor(() => expect(edit.calls).toHaveLength(1));
});
it('sends an incidentio_configs payload when editing an incident.io channel', async () => {
const edit = mockEditChannel();
render(
<EditAlertChannels
channelId="4"
initialValue={{
type: 'incidentio',
name: 'incidentio-channel',
url: 'https://api.incident.io/v2/alert_events/http/01M0D1JNVBGBGVTWX053EM12XV',
token: 'tok-abc',
send_resolved: true,
metadata: { env: 'prod' },
}}
/>,
);
const user = userEvent.setup();
await user.click(screen.getByTestId('save-channel-button'));
await waitFor(() => expect(edit.calls).toHaveLength(1));
expect(edit.calls[0].id).toBe('4');
expect(edit.calls[0].body).toStrictEqual({
name: 'incidentio-channel',
incidentio_configs: [
{
url: 'https://api.incident.io/v2/alert_events/http/01M0D1JNVBGBGVTWX053EM12XV',
token: 'tok-abc',
send_resolved: true,
metadata: { env: 'prod' },
},
],
});
});
it('persists send_resolved toggle in the edit request', async () => {
const edit = mockEditChannel();
render(

View File

@@ -107,6 +107,7 @@ export enum ChannelType {
GoogleChat = 'googlechat',
Jira = 'jira',
JsmOps = 'jsmops',
IncidentIO = 'incidentio',
}
// LabelFilterStatement will be used for preparing filter conditions / matchers
@@ -159,6 +160,22 @@ export interface JiraChannel extends Channel {
reopen_duration?: string;
}
// IncidentIOChannel configures the incident.io alert channel, backed by an
// incident.io HTTP alert source (Alert Events V2 API).
export interface IncidentIOChannel extends Channel {
// per-source alert events URL, e.g.
// https://api.incident.io/v2/alert_events/http/<source_config_id>
url: string;
// the alert source's secret token
token: string;
// alert title template
title?: string;
// alert body template (markdown, rendered natively by incident.io)
description?: string;
// extra metadata pairs merged over the alert's labels (channel wins on clash)
metadata?: Record<string, string>;
}
// JsmOpsChannel configures the Jira Service Management Ops alert channel
// (ex-Opsgenie alert API). Auth is the JSM integration API key.
export interface JsmOpsChannel extends Channel {

View File

@@ -2,6 +2,7 @@ import {
ChannelType,
EmailChannel,
GoogleChatChannel,
IncidentIOChannel,
JiraChannel,
JsmOpsChannel,
MsTeamsChannel,
@@ -144,6 +145,29 @@ export const JsmOpsInitialConfig: Partial<JsmOpsChannel> = {
tags: ['signoz-alert'],
};
// mirrors DefaultIncidentIOTitleTemplate / DefaultIncidentIODescriptionTemplate
// in pkg/types/alertmanagertypes/incidentio.go, applied by the backend when
// title / description are left empty. send_resolved is seeded on so incident.io
// alerts resolve with the rule (the backend cannot default it).
export const IncidentIOInitialConfig: Partial<IncidentIOChannel> = {
send_resolved: true,
title: `[{{ .Status | toUpper }}{{ if eq .Status "firing" }}:{{ .Alerts.Firing | len }}{{ end }}] {{ .CommonLabels.alertname }}`,
description: `{{ range .Alerts -}}
**Alert:** {{ .Labels.alertname }}{{ if .Labels.severity }} ({{ .Labels.severity }}){{ end }}
{{ if .Annotations.summary }}**Summary:** {{ .Annotations.summary }}
{{ end }}{{ if .Annotations.description }}**Description:** {{ .Annotations.description }}
{{ end }}{{ if .GeneratorURL }}[View in SigNoz]({{ .GeneratorURL }})
{{ end }}{{ if .Annotations.related_logs }}[View related logs]({{ .Annotations.related_logs }})
{{ end }}{{ if .Annotations.related_traces }}[View related traces]({{ .Annotations.related_traces }})
{{ end }}{{ end }}`,
};
export const EmailInitialConfig: Partial<EmailChannel> = {
send_resolved: true,
html: `<!--
@@ -553,7 +577,8 @@ export const ChannelInitialConfig: Record<
EmailChannel &
GoogleChatChannel &
JiraChannel &
JsmOpsChannel
JsmOpsChannel &
IncidentIOChannel
>
> = {
[ChannelType.Slack]: SlackInitialConfig,
@@ -561,6 +586,7 @@ export const ChannelInitialConfig: Record<
[ChannelType.GoogleChat]: GoogleChatInitialConfig,
[ChannelType.Jira]: JiraInitialConfig,
[ChannelType.JsmOps]: JsmOpsInitialConfig,
[ChannelType.IncidentIO]: IncidentIOInitialConfig,
[ChannelType.Pagerduty]: PagerInitialConfig,
[ChannelType.Opsgenie]: OpsgenieInitialConfig,
[ChannelType.Email]: EmailInitialConfig,

View File

@@ -32,6 +32,7 @@ import {
ChannelType,
EmailChannel,
GoogleChatChannel,
IncidentIOChannel,
JiraChannel,
JsmOpsChannel,
MsTeamsChannel,
@@ -45,9 +46,11 @@ import { ChannelInitialConfig } from './defaults';
import {
isChannelType,
isValidGoogleChatWebhookURL,
isValidIncidentIOURL,
isValidJiraReopenDuration,
isValidJiraSiteURL,
prepareGoogleChatRequest,
prepareIncidentIORequest,
prepareJiraRequest,
prepareJsmOpsRequest,
} from './utils';
@@ -77,7 +80,8 @@ function CreateAlertChannels({
EmailChannel &
GoogleChatChannel &
JiraChannel &
JsmOpsChannel
JsmOpsChannel &
IncidentIOChannel
>
>(() => ({
send_resolved: true,
@@ -550,6 +554,56 @@ function CreateAlertChannels({
showErrorModal,
]);
const validateIncidentIOConfig = useCallback((): boolean => {
if (!selectedConfig.url || !selectedConfig.token) {
notifications.error({
message: 'Error',
description: t('incidentio_required_fields'),
});
return false;
}
if (!isValidIncidentIOURL(selectedConfig.url)) {
notifications.error({
message: 'Error',
description: t('incidentio_url_invalid'),
});
return false;
}
return true;
}, [selectedConfig.url, selectedConfig.token, notifications, t]);
const onIncidentIOHandler = useCallback(async () => {
if (!validateIncidentIOConfig()) {
return { status: 'failed', statusMessage: t('channel_creation_failed') };
}
setSavingState(true);
try {
await createChannel({ data: prepareIncidentIORequest(selectedConfig) });
notifications.success({
message: 'Success',
description: t('channel_creation_done'),
});
history.replace(ROUTES.ALL_CHANNELS);
return { status: 'success', statusMessage: t('channel_creation_done') };
} catch (error) {
showErrorModal(toAPIError(error as ErrorType<RenderErrorResponseDTO>));
return { status: 'failed', statusMessage: t('channel_creation_failed') };
} finally {
setSavingState(false);
}
}, [
validateIncidentIOConfig,
createChannel,
selectedConfig,
notifications,
t,
showErrorModal,
]);
const onSaveHandler = useCallback(
async (value: ChannelType) => {
if (!selectedConfig.name) {
@@ -570,6 +624,7 @@ function CreateAlertChannels({
[ChannelType.GoogleChat]: onGoogleChatHandler,
[ChannelType.Jira]: onJiraHandler,
[ChannelType.JsmOps]: onJsmOpsHandler,
[ChannelType.IncidentIO]: onIncidentIOHandler,
};
if (isChannelType(value)) {
@@ -604,6 +659,7 @@ function CreateAlertChannels({
onGoogleChatHandler,
onJiraHandler,
onJsmOpsHandler,
onIncidentIOHandler,
notifications,
t,
],
@@ -662,6 +718,13 @@ function CreateAlertChannels({
}
await testChannel({ data: prepareJsmOpsRequest(selectedConfig) });
break;
case ChannelType.IncidentIO:
if (!validateIncidentIOConfig()) {
setTestingState(false);
return;
}
await testChannel({ data: prepareIncidentIORequest(selectedConfig) });
break;
default:
notifications.error({
message: 'Error',
@@ -712,6 +775,7 @@ function CreateAlertChannels({
validateGoogleChatConfig,
validateJiraConfig,
validateJsmOpsConfig,
validateIncidentIOConfig,
testChannel,
notifications,
],

View File

@@ -1,4 +1,5 @@
import {
AlertmanagertypesIncidentIOReceiverConfigDTO,
AlertmanagertypesJiraReceiverConfigDTO,
AlertmanagertypesJSMOpsReceiverConfigDTO,
AlertmanagertypesPostableChannelDTO,
@@ -9,6 +10,7 @@ import {
import {
ChannelType,
GoogleChatChannel,
IncidentIOChannel,
JiraChannel,
JsmOpsChannel,
} from './config';
@@ -168,3 +170,50 @@ export const prepareJsmOpsRequest = (
jsmops_configs: [jsmops],
};
};
const INCIDENTIO_EVENTS_PATH_PREFIX = '/v2/alert_events/http/';
// the backend enforces the same rule, this is only for a nicer error experience
export const isValidIncidentIOURL = (url: string): boolean => {
try {
const { protocol, pathname } = new URL(url);
const idx = pathname.indexOf(INCIDENTIO_EVENTS_PATH_PREFIX);
return (
protocol === 'https:' &&
idx !== -1 &&
pathname.length > idx + INCIDENTIO_EVENTS_PATH_PREFIX.length
);
} catch {
return false;
}
};
// create, update and test all send the same body shape. Optional fields are
// omitted when empty so the backend applies its defaults.
export const prepareIncidentIORequest = (
config: Partial<IncidentIOChannel>,
): AlertmanagertypesPostableChannelDTO => {
const incidentio: AlertmanagertypesIncidentIOReceiverConfigDTO = {
url: config.url || '',
token: config.token || '',
send_resolved: config.send_resolved || false,
};
if (config.title) {
incidentio.title = config.title;
}
if (config.description) {
incidentio.description = config.description;
}
const metadata = Object.fromEntries(
Object.entries(config.metadata || {}).filter(([key]) => key.trim() !== ''),
);
if (Object.keys(metadata).length > 0) {
incidentio.metadata = metadata;
}
return {
name: config.name || '',
incidentio_configs: [incidentio],
};
};

View File

@@ -25,6 +25,7 @@ import {
ChannelType,
EmailChannel,
GoogleChatChannel,
IncidentIOChannel,
JiraChannel,
JsmOpsChannel,
MsTeamsChannel,
@@ -36,9 +37,11 @@ import {
} from 'container/CreateAlertChannels/config';
import {
isValidGoogleChatWebhookURL,
isValidIncidentIOURL,
isValidJiraReopenDuration,
isValidJiraSiteURL,
prepareGoogleChatRequest,
prepareIncidentIORequest,
prepareJiraRequest,
prepareJsmOpsRequest,
} from 'container/CreateAlertChannels/utils';
@@ -66,7 +69,8 @@ function EditAlertChannels({
EmailChannel &
GoogleChatChannel &
JiraChannel &
JsmOpsChannel
JsmOpsChannel &
IncidentIOChannel
>
>({
...initialValue,
@@ -578,6 +582,61 @@ function EditAlertChannels({
t,
]);
const validateIncidentIOConfig = useCallback((): string => {
if (!selectedConfig.url || !selectedConfig.token) {
return t('incidentio_required_fields');
}
if (!isValidIncidentIOURL(selectedConfig.url)) {
return t('incidentio_url_invalid');
}
return '';
}, [selectedConfig, t]);
const onIncidentIOEditHandler = useCallback(async () => {
const validationError = validateIncidentIOConfig();
if (validationError !== '') {
notifications.error({
message: 'Error',
description: validationError,
});
return { status: 'failed', statusMessage: validationError };
}
setSavingState(true);
try {
await updateChannel({
pathParams: { id },
data: prepareIncidentIORequest(selectedConfig),
});
notifications.success({
message: 'Success',
description: t('channel_edit_done'),
});
history.replace(ROUTES.ALL_CHANNELS);
return { status: 'success', statusMessage: t('channel_edit_done') };
} catch (error) {
const apiError = notifyError(error);
return {
status: 'failed',
statusMessage: apiError.getErrorMessage() || t('channel_edit_failed'),
};
} finally {
setSavingState(false);
}
}, [
validateIncidentIOConfig,
updateChannel,
id,
selectedConfig,
notifications,
notifyError,
t,
]);
const onSaveHandler = useCallback(
async (value: ChannelType) => {
let result;
@@ -599,6 +658,8 @@ function EditAlertChannels({
result = await onJiraEditHandler();
} else if (value === ChannelType.JsmOps) {
result = await onJsmOpsEditHandler();
} else if (value === ChannelType.IncidentIO) {
result = await onIncidentIOEditHandler();
}
logEvent('Alert Channel: Save channel', {
type: value,
@@ -620,6 +681,7 @@ function EditAlertChannels({
onGoogleChatEditHandler,
onJiraEditHandler,
onJsmOpsEditHandler,
onIncidentIOEditHandler,
],
);
@@ -701,6 +763,19 @@ function EditAlertChannels({
await testChannel({ data: prepareJsmOpsRequest(selectedConfig) });
break;
}
case ChannelType.IncidentIO: {
const validationError = validateIncidentIOConfig();
if (validationError !== '') {
notifications.error({
message: 'Error',
description: validationError,
});
setTestingState(false);
return;
}
await testChannel({ data: prepareIncidentIORequest(selectedConfig) });
break;
}
default:
notifications.error({
message: 'Error',
@@ -740,6 +815,7 @@ function EditAlertChannels({
validateGoogleChatConfig,
validateJiraConfig,
validateJsmOpsConfig,
validateIncidentIOConfig,
testChannel,
prepareWebhookRequest,
preparePagerRequest,

View File

@@ -0,0 +1,172 @@
import { Dispatch, SetStateAction, useState } from 'react';
import { useTranslation } from 'react-i18next';
import { Minus, Plus } from '@signozhq/icons';
import { Button, Form, Input } from 'antd';
import { Typography } from '@signozhq/ui/typography';
import { IncidentIOChannel } from '../../CreateAlertChannels/config';
interface MetadataRow {
key: string;
value: string;
}
function IncidentIOSettings({
setSelectedConfig,
initialMetadata,
}: IncidentIOProps): JSX.Element {
const { t } = useTranslation('channels');
const [metadataRows, setMetadataRows] = useState<MetadataRow[]>(() =>
Object.entries(initialMetadata || {}).map(([key, value]) => ({
key,
value,
})),
);
const update = (patch: Partial<IncidentIOChannel>): void =>
setSelectedConfig((value) => ({ ...value, ...patch }));
const syncMetadata = (rows: MetadataRow[]): void => {
setMetadataRows(rows);
update({
metadata: Object.fromEntries(
rows
.filter((row) => row.key.trim() !== '')
.map((row) => [row.key, row.value]),
),
});
};
return (
<>
<Typography.Text
color="muted"
size="sm"
testId="incidentio-tip"
style={{ display: 'block', marginBottom: 16 }}
>
{t('incidentio_tip')}{' '}
<Typography.Link
href="https://signoz.io/docs/alerts-management/notification-channel/incidentio/"
target="_blank"
rel="noopener noreferrer"
>
{t('incidentio_tip_link')}
</Typography.Link>
</Typography.Text>
<Form.Item
name="url"
label={t('field_incidentio_url')}
help={t('help_incidentio_url')}
required
>
<Input
onChange={(event): void => update({ url: event.target.value })}
data-testid="incidentio-url-textbox"
/>
</Form.Item>
<Form.Item
name="token"
label={t('field_incidentio_token')}
help={t('help_incidentio_token')}
required
>
<Input
type="password"
onChange={(event): void => update({ token: event.target.value })}
data-testid="incidentio-token-textbox"
/>
</Form.Item>
<Form.Item
name="title"
label={t('field_incidentio_title')}
help={t('help_incidentio_title')}
>
<Input.TextArea
rows={2}
onChange={(event): void => update({ title: event.target.value })}
data-testid="incidentio-title-textarea"
/>
</Form.Item>
<Form.Item
name="description"
label={t('field_incidentio_description')}
help={t('help_incidentio_description')}
>
<Input.TextArea
rows={6}
onChange={(event): void => update({ description: event.target.value })}
data-testid="incidentio-description-textarea"
/>
</Form.Item>
<Form.Item
label={t('field_incidentio_metadata')}
help={t('help_incidentio_metadata')}
>
{metadataRows.map((row, index) => (
// rows have no stable identity beyond their position
// eslint-disable-next-line react/no-array-index-key
<div key={index} style={{ display: 'flex', gap: 8, marginBottom: 8 }}>
<Input
placeholder={t('placeholder_incidentio_metadata_key')}
value={row.key}
onChange={(event): void =>
syncMetadata(
metadataRows.map((r, i) =>
i === index ? { ...r, key: event.target.value } : r,
),
)
}
data-testid={`incidentio-metadata-key-${index}`}
/>
<Input
placeholder={t('placeholder_incidentio_metadata_value')}
value={row.value}
onChange={(event): void =>
syncMetadata(
metadataRows.map((r, i) =>
i === index ? { ...r, value: event.target.value } : r,
),
)
}
data-testid={`incidentio-metadata-value-${index}`}
/>
<Button
icon={<Minus size={14} />}
onClick={(): void =>
syncMetadata(metadataRows.filter((_, i) => i !== index))
}
data-testid={`incidentio-metadata-remove-${index}`}
/>
</div>
))}
<Button
type="dashed"
icon={<Plus size={14} />}
onClick={(): void =>
syncMetadata([...metadataRows, { key: '', value: '' }])
}
data-testid="incidentio-metadata-add"
>
{t('button_incidentio_add_metadata')}
</Button>
</Form.Item>
</>
);
}
interface IncidentIOProps {
setSelectedConfig: Dispatch<SetStateAction<Partial<IncidentIOChannel>>>;
initialMetadata?: Record<string, string>;
}
IncidentIOSettings.defaultProps = {
initialMetadata: undefined,
};
export default IncidentIOSettings;

View File

@@ -10,6 +10,7 @@ import {
ChannelType,
EmailChannel,
GoogleChatChannel,
IncidentIOChannel,
JiraChannel,
JsmOpsChannel,
OpsgenieChannel,
@@ -21,6 +22,7 @@ import history from 'lib/history';
import EmailSettings from './Settings/Email';
import GoogleChatSettings from './Settings/GoogleChat';
import IncidentIOSettings from './Settings/IncidentIo';
import JiraSettings from './Settings/Jira';
import JsmOpsSettings from './Settings/JsmOps';
import MsTeamsSettings from './Settings/MsTeams';
@@ -61,6 +63,13 @@ function FormAlertChannels({
return <JiraSettings setSelectedConfig={setSelectedConfig} />;
case ChannelType.JsmOps:
return <JsmOpsSettings setSelectedConfig={setSelectedConfig} />;
case ChannelType.IncidentIO:
return (
<IncidentIOSettings
setSelectedConfig={setSelectedConfig}
initialMetadata={initialValue?.metadata as Record<string, string>}
/>
);
case ChannelType.Opsgenie:
return <OpsgenieSettings setSelectedConfig={setSelectedConfig} />;
case ChannelType.Email:
@@ -157,6 +166,14 @@ function FormAlertChannels({
<Select.Option value="jsmops" key="jsmops" data-testid="select-option">
Jira Service Management Ops
</Select.Option>
<Select.Option
value="incidentio"
key="incidentio"
data-testid="select-option"
>
incident.io
</Select.Option>
</Select>
</Form.Item>
@@ -207,7 +224,8 @@ interface FormAlertChannelsProps {
EmailChannel &
GoogleChatChannel &
JiraChannel &
JsmOpsChannel
JsmOpsChannel &
IncidentIOChannel
>
>
>;

View File

@@ -11,6 +11,7 @@ import ROUTES from 'constants/routes';
import {
ChannelType,
GoogleChatChannel,
IncidentIOChannel,
JiraChannel,
JsmOpsChannel,
MsTeamsChannel,
@@ -69,7 +70,8 @@ function ChannelsEdit(): JSX.Element {
MsTeamsChannel &
GoogleChatChannel &
JiraChannel &
JsmOpsChannel
JsmOpsChannel &
IncidentIOChannel
>;
} => {
let channel: Partial<
@@ -79,7 +81,8 @@ function ChannelsEdit(): JSX.Element {
MsTeamsChannel &
GoogleChatChannel &
JiraChannel &
JsmOpsChannel
JsmOpsChannel &
IncidentIOChannel
> = {
name: '',
};
@@ -135,6 +138,15 @@ function ChannelsEdit(): JSX.Element {
};
}
if (value && 'incidentio_configs' in value) {
const [incidentIOConfig] = value.incidentio_configs;
channel = incidentIOConfig;
return {
type: ChannelType.IncidentIO,
channel,
};
}
if (value && 'jsmops_configs' in value) {
const [jsmopsConfig] = value.jsmops_configs;
channel = jsmopsConfig;

View File

@@ -0,0 +1,188 @@
package incidentio
import (
"bytes"
"context"
"encoding/json"
"log/slog"
"maps"
"net/http"
"strings"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagertemplate"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
"github.com/SigNoz/signoz/pkg/types/ruletypes"
"github.com/prometheus/alertmanager/notify"
"github.com/prometheus/alertmanager/template"
"github.com/prometheus/alertmanager/types"
)
const (
Integration = "incidentio"
// incident.io rejects payloads over 512 KB with a 413. Runes cap the
// description at 4 bytes each worst case (~400 KB), leaving headroom for
// the other fields.
maxDescriptionLenRunes = 100000
statusFiring = "firing"
statusResolved = "resolved"
)
// alertEvent is the body of incident.io's HTTP alert source endpoint
// (Alert Events V2 API). Title, status and deduplication_key are required;
// metadata values must be flat scalars. Repeat events for a firing key are
// dropped server-side and resolves for unknown keys are safe no-ops, so
// events are sent unconditionally.
type alertEvent struct {
Title string `json:"title"`
Description string `json:"description,omitempty"`
Status string `json:"status"`
DeduplicationKey string `json:"deduplication_key"`
SourceURL string `json:"source_url,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
}
type Notifier struct {
conf *alertmanagertypes.IncidentIOReceiverConfig
tmpl *template.Template
logger *slog.Logger
client *http.Client
retrier *notify.Retrier
templater alertmanagertypes.Templater
}
func New(conf *alertmanagertypes.IncidentIOReceiverConfig, t *template.Template, l *slog.Logger, templater alertmanagertypes.Templater) (*Notifier, error) {
if conf.HTTPConfig == nil {
return nil, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "incidentio http_config is nil")
}
client, err := notify.NewClientWithTracing(*conf.HTTPConfig, Integration)
if err != nil {
return nil, err
}
return &Notifier{
conf: conf,
tmpl: t,
logger: l,
client: client,
retrier: &notify.Retrier{RetryCodes: []int{http.StatusTooManyRequests}},
templater: templater,
}, nil
}
func (n *Notifier) Notify(ctx context.Context, as ...*types.Alert) (bool, error) {
key, err := notify.ExtractGroupKey(ctx)
if err != nil {
return false, err
}
firing := types.Alerts(as...).HasFiring()
n.logger.DebugContext(ctx, "sending incidentio notification", slog.String("group_key", key.String()), slog.Bool("firing", firing))
customTitle, customBody := alertmanagertemplate.ExtractTemplatesFromAnnotations(as)
result, err := n.templater.Expand(ctx, alertmanagertypes.ExpandRequest{
TitleTemplate: customTitle,
BodyTemplate: customBody,
DefaultTitleTemplate: n.conf.Title,
DefaultBodyTemplate: n.conf.Description,
}, as)
if err != nil {
return false, err
}
// title is required by the API; a channel title template can render empty,
// so fall back to the rule name, then to a static last resort.
title := result.Title
if strings.TrimSpace(title) == "" && len(as) > 0 {
title = string(as[0].Labels[ruletypes.LabelAlertName])
}
if strings.TrimSpace(title) == "" {
title = "SigNoz alert"
}
var parts []string
for _, body := range result.Body {
if body != "" {
parts = append(parts, body)
}
}
description := truncateRunes(strings.Join(parts, "\n\n---\n\n"), maxDescriptionLenRunes)
status := statusFiring
if !firing {
status = statusResolved
}
event := alertEvent{
Title: title,
Description: description,
Status: status,
DeduplicationKey: key.Hash(),
SourceURL: sourceURL(as),
Metadata: n.metadata(ctx, as),
}
var buf bytes.Buffer
if err := json.NewEncoder(&buf).Encode(event); err != nil {
return false, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, n.conf.URL, &buf)
if err != nil {
return false, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+string(n.conf.Token))
resp, err := n.client.Do(req) //nolint:bodyclose // notify.Drain closes the body
if err != nil {
return true, notify.RedactURL(err)
}
defer notify.Drain(resp)
shouldRetry, err := n.retrier.Check(resp.StatusCode, resp.Body)
if err != nil {
return shouldRetry, notify.NewErrorWithReason(notify.GetFailureReasonFromStatusCode(resp.StatusCode), err)
}
return shouldRetry, nil
}
// metadata copies the group's common labels wholesale (the Opsgenie details
// precedent), so severity, ruleId and any user-defined rule labels arrive as
// flat strings ready for incident.io attribute mapping. Channel-configured
// pairs are template-expanded and laid on top (channel wins on key clash);
// a value that fails to expand is sent raw so delivery never breaks on it.
func (n *Notifier) metadata(ctx context.Context, as []*types.Alert) map[string]string {
data := notify.GetTemplateData(ctx, n.tmpl, as, n.logger)
out := make(map[string]string, len(data.CommonLabels)+len(n.conf.Metadata))
maps.Copy(out, data.CommonLabels)
for k, v := range n.conf.Metadata {
expanded, err := n.tmpl.ExecuteTextString(v, data)
if err != nil {
n.logger.WarnContext(ctx, "failed to expand incidentio metadata value, sending it raw", slog.String("metadata_key", k))
expanded = v
}
out[k] = expanded
}
if len(out) == 0 {
return nil
}
return out
}
// sourceURL returns the per-rule SigNoz link from the ruleSource label, which
// is identical for every alert in the group.
func sourceURL(as []*types.Alert) string {
if len(as) == 0 {
return ""
}
return string(as[0].Labels[ruletypes.LabelRuleSource])
}
func truncateRunes(s string, max int) string {
runes := []rune(s)
if len(runes) <= max {
return s
}
return string(runes[:max-1]) + "…"
}

View File

@@ -0,0 +1,216 @@
package incidentio
import (
"context"
"encoding/json"
"log/slog"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagertemplate"
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
"github.com/prometheus/alertmanager/notify"
"github.com/prometheus/alertmanager/notify/test"
"github.com/prometheus/alertmanager/types"
commoncfg "github.com/prometheus/common/config"
"github.com/prometheus/common/model"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
type mockIncidentIO struct {
srv *httptest.Server
mu sync.Mutex
events []alertEvent
auths []string
status int
}
func newMockIncidentIO(t *testing.T) *mockIncidentIO {
t.Helper()
m := &mockIncidentIO{status: http.StatusAccepted}
m.srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var ev alertEvent
_ = json.NewDecoder(r.Body).Decode(&ev)
m.mu.Lock()
m.events = append(m.events, ev)
m.auths = append(m.auths, r.Header.Get("Authorization"))
status := m.status
m.mu.Unlock()
w.WriteHeader(status)
if status == http.StatusAccepted {
_, _ = w.Write([]byte(`{"status":"accepted","message":"Event accepted for processing","deduplication_key":"` + ev.DeduplicationKey + `"}`))
} else {
_, _ = w.Write([]byte(`{"type":"validation_error","status":422,"errors":[{"code":"is_required","message":"Deduplication key is required"}]}`))
}
}))
t.Cleanup(m.srv.Close)
return m
}
func (m *mockIncidentIO) lastEvent(t *testing.T) alertEvent {
t.Helper()
m.mu.Lock()
defer m.mu.Unlock()
require.NotEmpty(t, m.events)
return m.events[len(m.events)-1]
}
func newNotifier(t *testing.T, m *mockIncidentIO) *Notifier {
t.Helper()
tmpl := test.CreateTmpl(t)
n, err := New(&alertmanagertypes.IncidentIOReceiverConfig{
URL: m.srv.URL + "/v2/alert_events/http/src-1",
Token: "tok-1",
Title: alertmanagertypes.DefaultIncidentIOTitleTemplate,
Description: alertmanagertypes.DefaultIncidentIODescriptionTemplate,
HTTPConfig: &commoncfg.HTTPClientConfig{},
}, tmpl, slog.New(slog.DiscardHandler), alertmanagertemplate.New(tmpl, slog.New(slog.DiscardHandler)))
require.NoError(t, err)
return n
}
func alert(firing bool) *types.Alert {
a := &types.Alert{Alert: model.Alert{
Labels: model.LabelSet{
"alertname": "HighCPU",
"severity": "critical",
"ruleSource": "https://signoz.example/alerts/edit?ruleId=1",
},
Annotations: model.LabelSet{"summary": "cpu high", "related_logs": "https://signoz.example/logs?q=1"},
GeneratorURL: "https://signoz.example/alerts/edit?ruleId=1",
StartsAt: time.Now().Add(-time.Minute),
}}
if firing {
a.EndsAt = time.Now().Add(time.Hour)
} else {
a.EndsAt = time.Now().Add(-time.Minute)
}
return a
}
func ctx() context.Context {
return notify.WithGroupKey(context.Background(), "test-incidentio")
}
func TestNotifyFiringEvent(t *testing.T) {
m := newMockIncidentIO(t)
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
require.NoError(t, err)
assert.False(t, retry)
ev := m.lastEvent(t)
assert.Equal(t, "[FIRING:1] HighCPU", ev.Title)
assert.Equal(t, "firing", ev.Status)
assert.NotEmpty(t, ev.DeduplicationKey)
assert.Equal(t, "https://signoz.example/alerts/edit?ruleId=1", ev.SourceURL)
assert.Contains(t, ev.Description, "**Alert:** HighCPU (critical)")
assert.Contains(t, ev.Description, "**Summary:** cpu high")
assert.Contains(t, ev.Description, "[View in SigNoz](https://signoz.example/alerts/edit?ruleId=1)")
assert.Contains(t, ev.Description, "[View related logs](https://signoz.example/logs?q=1)")
assert.Equal(t, map[string]string{
"alertname": "HighCPU",
"severity": "critical",
"ruleSource": "https://signoz.example/alerts/edit?ruleId=1",
}, ev.Metadata)
assert.Equal(t, "Bearer tok-1", m.auths[0])
}
func TestNotifyResolvedEventReusesDedupKey(t *testing.T) {
m := newMockIncidentIO(t)
n := newNotifier(t, m)
_, err := n.Notify(ctx(), alert(true))
require.NoError(t, err)
firingKey := m.lastEvent(t).DeduplicationKey
_, err = n.Notify(ctx(), alert(false))
require.NoError(t, err)
ev := m.lastEvent(t)
assert.Equal(t, "resolved", ev.Status)
assert.Equal(t, firingKey, ev.DeduplicationKey)
}
func TestNotifyPermanentFailureDoesNotRetry(t *testing.T) {
m := newMockIncidentIO(t)
m.status = http.StatusUnprocessableEntity
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
require.Error(t, err)
assert.False(t, retry)
assert.Contains(t, err.Error(), "Deduplication key is required") // response body surfaces to the user
}
func TestNotifyRateLimitRetries(t *testing.T) {
m := newMockIncidentIO(t)
m.status = http.StatusTooManyRequests
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
require.Error(t, err)
assert.True(t, retry)
}
func TestNotifyMergesChannelMetadata(t *testing.T) {
m := newMockIncidentIO(t)
tmpl := test.CreateTmpl(t)
n, err := New(&alertmanagertypes.IncidentIOReceiverConfig{
URL: m.srv.URL + "/v2/alert_events/http/src-1",
Token: "tok-1",
Title: alertmanagertypes.DefaultIncidentIOTitleTemplate,
Description: alertmanagertypes.DefaultIncidentIODescriptionTemplate,
HTTPConfig: &commoncfg.HTTPClientConfig{},
Metadata: map[string]string{
"env": "prod",
"sev": "{{ .CommonLabels.severity }}",
"alertname": "channel-wins",
"broken": "{{ .Nope",
},
}, tmpl, slog.New(slog.DiscardHandler), alertmanagertemplate.New(tmpl, slog.New(slog.DiscardHandler)))
require.NoError(t, err)
_, err = n.Notify(ctx(), alert(true))
require.NoError(t, err) // a broken metadata template must not fail delivery
md := m.lastEvent(t).Metadata
assert.Equal(t, "prod", md["env"])
assert.Equal(t, "critical", md["sev"]) // values are template-expanded
assert.Equal(t, "channel-wins", md["alertname"]) // channel overrides the rule label
assert.Equal(t, "{{ .Nope", md["broken"]) // unexpandable value sent raw
assert.Equal(t, "critical", md["severity"]) // rule labels still present
}
func TestNotifyEmptyTitleFallsBackToRuleName(t *testing.T) {
m := newMockIncidentIO(t)
tmpl := test.CreateTmpl(t)
n, err := New(&alertmanagertypes.IncidentIOReceiverConfig{
URL: m.srv.URL + "/v2/alert_events/http/src-1",
Token: "tok-1",
Title: `{{ .CommonLabels.nonexistent }}`,
Description: alertmanagertypes.DefaultIncidentIODescriptionTemplate,
HTTPConfig: &commoncfg.HTTPClientConfig{},
}, tmpl, slog.New(slog.DiscardHandler), alertmanagertemplate.New(tmpl, slog.New(slog.DiscardHandler)))
require.NoError(t, err)
_, err = n.Notify(ctx(), alert(true))
require.NoError(t, err)
assert.Equal(t, "HighCPU", m.lastEvent(t).Title)
}
func TestNotifyTruncatesLongDescription(t *testing.T) {
m := newMockIncidentIO(t)
a := alert(true)
a.Annotations["description"] = model.LabelValue(strings.Repeat("x", maxDescriptionLenRunes+1000))
_, err := newNotifier(t, m).Notify(ctx(), a)
require.NoError(t, err)
desc := []rune(m.lastEvent(t).Description)
assert.LessOrEqual(t, len(desc), maxDescriptionLenRunes)
assert.Equal(t, '…', desc[len(desc)-1])
}

View File

@@ -6,6 +6,7 @@ import (
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/email"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/googlechat"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/incidentio"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/jira"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/jsmops"
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/msteamsv2"
@@ -30,6 +31,7 @@ var customNotifierIntegrations = []string{
googlechat.Integration,
jira.Integration,
jsmops.Integration,
incidentio.Integration,
}
func NewReceiverIntegrations(nc *alertmanagertypes.Receiver, tmpl *template.Template, logger *slog.Logger, templater alertmanagertypes.Templater) ([]notify.Integration, error) {
@@ -95,6 +97,11 @@ func NewReceiverIntegrations(nc *alertmanagertypes.Receiver, tmpl *template.Temp
return jsmops.New(c, tmpl, l, templater, true)
})
}
for i, c := range nc.IncidentIOConfigs {
add(incidentio.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) {
return incidentio.New(c, tmpl, l, templater)
})
}
if errs.Len() > 0 {
return nil, &errs

View File

@@ -75,7 +75,7 @@ func (store *config) CreateChannel(ctx context.Context, channel *alertmanagertyp
NewInsert().
Model(channel).
Exec(ctx); err != nil {
return err
return store.sqlstore.WrapAlreadyExistsErrf(err, alertmanagertypes.ErrCodeAlertmanagerChannelAlreadyExists, "channel with name %q already exists", channel.Name)
}
return nil

View File

@@ -0,0 +1,71 @@
package sqlalertmanagerstore
import (
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
func TestCreateChannelRejectsDuplicateNameInSameOrg(t *testing.T) {
sqlstore := newTestStore(t)
_, err := sqlstore.BunDB().NewCreateTable().
Model((*alertmanagertypes.Channel)(nil)).
IfNotExists().
Exec(t.Context())
require.NoError(t, err)
_, err = sqlstore.BunDB().NewCreateIndex().
Model((*alertmanagertypes.Channel)(nil)).
Index("notification_channel_org_id_name_idx").
Column("org_id", "name").
Unique().
Exec(t.Context())
require.NoError(t, err)
store := NewConfigStore(sqlstore)
orgID := valuer.GenerateUUID().StringValue()
now := time.Now().UTC()
firstChannel := &alertmanagertypes.Channel{
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
TimeAuditable: types.TimeAuditable{CreatedAt: now, UpdatedAt: now},
Name: "shared-name",
DisplayName: "First Channel",
Type: "slack",
Data: `{"name":"First Channel","slack_configs":[{"api_url":"https://hooks.slack.com/services/first"}]}`,
OrgID: orgID,
}
require.NoError(t, store.CreateChannel(t.Context(), firstChannel))
duplicateChannel := &alertmanagertypes.Channel{
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
TimeAuditable: types.TimeAuditable{CreatedAt: now, UpdatedAt: now},
Name: "shared-name",
DisplayName: "Second Channel",
Type: "slack",
Data: `{"name":"Second Channel","slack_configs":[{"api_url":"https://hooks.slack.com/services/second"}]}`,
OrgID: orgID,
}
err = store.CreateChannel(t.Context(), duplicateChannel)
require.Error(t, err)
assert.True(t, errors.Ast(err, errors.TypeAlreadyExists))
otherOrgChannel := &alertmanagertypes.Channel{
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
TimeAuditable: types.TimeAuditable{CreatedAt: now, UpdatedAt: now},
Name: "shared-name",
DisplayName: "Second Channel",
Type: "slack",
Data: `{"name":"Second Channel","slack_configs":[{"api_url":"https://hooks.slack.com/services/second"}]}`,
OrgID: valuer.GenerateUUID().StringValue(),
}
assert.NoError(t, store.CreateChannel(t.Context(), otherOrgChannel))
}

View File

@@ -187,7 +187,7 @@ func (provider *provider) DeleteChannelByID(ctx context.Context, orgID string, c
}
// Check if channel is referenced by any route policy (rule-based or policy-based)
policies, err := provider.notificationManager.GetRoutePoliciesByChannel(ctx, orgID, channel.Name)
policies, err := provider.notificationManager.GetRoutePoliciesByChannel(ctx, orgID, channel.DisplayName)
if err != nil {
return err
}
@@ -198,7 +198,7 @@ func (provider *provider) DeleteChannelByID(ctx context.Context, orgID string, c
}
return errors.NewInvalidInputf(errors.CodeInvalidInput,
"channel %q cannot be deleted because it is used by the following routing policies: %v",
channel.Name, names)
channel.DisplayName, names)
}
config, err := provider.configStore.Get(ctx, orgID)
@@ -206,7 +206,7 @@ func (provider *provider) DeleteChannelByID(ctx context.Context, orgID string, c
return err
}
if err := config.DeleteReceiver(channel.Name); err != nil {
if err := config.DeleteReceiver(channel.DisplayName); err != nil {
return err
}

View File

@@ -97,23 +97,57 @@ func (provider *provider) addGatewayRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_keys/{keyId}/limits", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.CreateIngestionKeyLimit), handler.OpenAPIDef{
if err := router.Handle("/api/v2/gateway/ingestion_keys/{keyId}", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.GetIngestionKey), handler.OpenAPIDef{
ID: "GetIngestionKey",
Tags: []string{"gateway"},
Summary: "Get ingestion key for workspace",
Description: "This endpoint returns an ingestion key for the workspace",
Request: nil,
RequestContentType: "",
Response: new(gatewaytypes.IngestionKey),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_keys/{keyId}/limits", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.DeprecatedCreateIngestionKeyLimit), handler.OpenAPIDef{
ID: "CreateIngestionKeyLimit",
Tags: []string{"gateway"},
Summary: "Create limit for the ingestion key",
Description: "This endpoint creates an ingestion key limit",
Request: new(gatewaytypes.PostableIngestionKeyLimit),
Request: new(gatewaytypes.DeprecatedPostableIngestionKeyLimit),
RequestContentType: "application/json",
Response: new(gatewaytypes.GettableCreatedIngestionKeyLimit),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{},
Deprecated: false,
Deprecated: true,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_keys/{keyId}/limits", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.GetIngestionKeyLimits), handler.OpenAPIDef{
ID: "GetIngestionKeyLimits",
Tags: []string{"gateway"},
Summary: "Get limits for the ingestion key",
Description: "This endpoint returns the ingestion limits for an ingestion key",
Request: nil,
RequestContentType: "",
Response: new([]gatewaytypes.Limit),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_keys/limits/{limitId}", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.UpdateIngestionKeyLimit), handler.OpenAPIDef{
ID: "UpdateIngestionKeyLimit",
Tags: []string{"gateway"},
@@ -125,7 +159,7 @@ func (provider *provider) addGatewayRoutes(router *mux.Router) error {
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{},
Deprecated: false,
Deprecated: true,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodPatch).GetError(); err != nil {
return err
@@ -142,6 +176,74 @@ func (provider *provider) addGatewayRoutes(router *mux.Router) error {
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{},
Deprecated: true,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodDelete).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_limits", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.CreateIngestionKeyLimit), handler.OpenAPIDef{
ID: "CreateIngestionLimit",
Tags: []string{"gateway"},
Summary: "Create ingestion limit",
Description: "This endpoint creates an ingestion limit for the ingestion key referenced by keyId",
Request: new(gatewaytypes.PostableIngestionKeyLimit),
RequestContentType: "application/json",
Response: new(types.Identifiable),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_limits/{limitId}", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.GetIngestionKeyLimit), handler.OpenAPIDef{
ID: "GetIngestionLimit",
Tags: []string{"gateway"},
Summary: "Get ingestion limit",
Description: "This endpoint returns an ingestion limit",
Request: nil,
RequestContentType: "",
Response: new(gatewaytypes.Limit),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_limits/{limitId}", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.UpdateIngestionKeyLimit), handler.OpenAPIDef{
ID: "UpdateIngestionLimit",
Tags: []string{"gateway"},
Summary: "Update ingestion limit",
Description: "This endpoint updates an ingestion limit",
Request: new(gatewaytypes.UpdatableIngestionKeyLimit),
RequestContentType: "application/json",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodPatch).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/gateway/ingestion_limits/{limitId}", handler.New(provider.authzMiddleware.EditAccess(provider.gatewayHandler.DeleteIngestionKeyLimit), handler.OpenAPIDef{
ID: "DeleteIngestionLimit",
Tags: []string{"gateway"},
Summary: "Delete ingestion limit",
Description: "This endpoint deletes an ingestion limit",
Request: nil,
RequestContentType: "",
Response: nil,
ResponseContentType: "",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodDelete).GetError(); err != nil {

View File

@@ -25,6 +25,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/organization"
"github.com/SigNoz/signoz/pkg/modules/preference"
"github.com/SigNoz/signoz/pkg/modules/promote"
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
"github.com/SigNoz/signoz/pkg/modules/rawdataexport"
"github.com/SigNoz/signoz/pkg/modules/rulestatehistory"
"github.com/SigNoz/signoz/pkg/modules/savedview"
@@ -84,6 +85,8 @@ type provider struct {
llmPricingRuleHandler llmpricingrule.Handler
statsHandler statsreporter.Handler
savedViewHandler savedview.Handler
quickFilterModule quickfilter.Module
quickFilterHandler quickfilter.Handler
}
func NewFactory(
@@ -124,6 +127,8 @@ func NewFactory(
rulerHandler ruler.Handler,
statsHandler statsreporter.Handler,
savedViewHandler savedview.Handler,
quickFilterModule quickfilter.Module,
quickFilterHandler quickfilter.Handler,
) factory.ProviderFactory[apiserver.APIServer, apiserver.Config] {
return factory.NewProviderFactory(factory.MustNewName("signoz"), func(ctx context.Context, providerSettings factory.ProviderSettings, config apiserver.Config) (apiserver.APIServer, error) {
return newProvider(
@@ -167,6 +172,8 @@ func NewFactory(
rulerHandler,
statsHandler,
savedViewHandler,
quickFilterModule,
quickFilterHandler,
)
})
}
@@ -212,6 +219,8 @@ func newProvider(
rulerHandler ruler.Handler,
statsHandler statsreporter.Handler,
savedViewHandler savedview.Handler,
quickFilterModule quickfilter.Module,
quickFilterHandler quickfilter.Handler,
) (apiserver.APIServer, error) {
settings := factory.NewScopedProviderSettings(providerSettings, "github.com/SigNoz/signoz/pkg/apiserver/signozapiserver")
router := mux.NewRouter().UseEncodedPath()
@@ -256,6 +265,8 @@ func newProvider(
llmPricingRuleHandler: llmPricingRuleHandler,
statsHandler: statsHandler,
savedViewHandler: savedViewHandler,
quickFilterModule: quickFilterModule,
quickFilterHandler: quickFilterHandler,
}
provider.authzMiddleware = middleware.NewAuthZ(settings.Logger(), orgGetter, authzService)
@@ -404,6 +415,10 @@ func (provider *provider) AddToRouter(router *mux.Router) error {
return err
}
if err := provider.addQuickFilterRoutes(router); err != nil {
return err
}
return nil
}

View File

@@ -0,0 +1,120 @@
package signozapiserver
import (
"context"
"net/http"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/http/handler"
"github.com/SigNoz/signoz/pkg/types/authtypes"
"github.com/SigNoz/signoz/pkg/types/coretypes"
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/gorilla/mux"
)
func (provider *provider) addQuickFilterRoutes(router *mux.Router) error {
if err := router.Handle("/api/v2/quick_filters", handler.New(
provider.authzMiddleware.CheckResources(provider.quickFilterHandler.ListQuickFiltersV2, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName),
handler.OpenAPIDef{
ID: "ListQuickFilters",
Tags: []string{"quick_filter"},
Summary: "List quick filters",
Description: "Returns the org's quick filters for every source, each filter as a telemetry field key.",
Request: nil,
RequestContentType: "",
Response: make([]*quickfiltertypes.SourceFilters, 0),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceQuickFilter.Scope(coretypes.VerbList)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceMetaResourceQuickFilter,
Verb: coretypes.VerbList,
Category: coretypes.ActionCategoryDataAccess,
Selector: coretypes.WildcardSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/quick_filters/{source}", handler.New(
provider.authzMiddleware.CheckResources(provider.quickFilterHandler.GetQuickFiltersV2, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName),
handler.OpenAPIDef{
ID: "GetQuickFilters",
Tags: []string{"quick_filter"},
Summary: "Get a source's quick filters",
Description: "Returns the org's quick filters for one source, each filter as a telemetry field key.",
Request: nil,
RequestContentType: "",
Response: new(quickfiltertypes.SourceFilters),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceQuickFilter.Scope(coretypes.VerbRead)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceMetaResourceQuickFilter,
Verb: coretypes.VerbRead,
Category: coretypes.ActionCategoryDataAccess,
ID: coretypes.PathParam("source"),
Selector: provider.quickFilterSelector,
}),
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/quick_filters/{source}", handler.New(
provider.authzMiddleware.CheckResources(provider.quickFilterHandler.UpdateQuickFiltersV2, authtypes.SigNozAdminRoleName),
handler.OpenAPIDef{
ID: "UpdateQuickFilters",
Tags: []string{"quick_filter"},
Summary: "Update quick filters",
Description: "Replaces the org's quick filters for the source named in the path.",
Request: new(quickfiltertypes.UpdatableQuickFilters),
RequestContentType: "application/json",
Response: nil,
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest},
Deprecated: false,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceQuickFilter.Scope(coretypes.VerbUpdate)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
Resource: coretypes.ResourceMetaResourceQuickFilter,
Verb: coretypes.VerbUpdate,
Category: coretypes.ActionCategoryConfigurationChange,
ID: coretypes.PathParam("source"),
Selector: provider.quickFilterSelector,
}),
)).Methods(http.MethodPut).GetError(); err != nil {
return err
}
return nil
}
func (provider *provider) quickFilterSelector(ctx context.Context, resource coretypes.Resource, source string, orgID valuer.UUID) ([]coretypes.Selector, error) {
validatedSource, err := quickfiltertypes.NewSource(source)
if err != nil {
return nil, err
}
// A source can have no stored row yet: GET serves it as empty and PUT
// creates it, so only the wildcard grant applies until the row exists.
quickFilter, err := provider.quickFilterModule.Get(ctx, orgID, validatedSource)
if err != nil {
if errors.Ast(err, errors.TypeNotFound) {
return []coretypes.Selector{resource.Type().MustSelector(coretypes.WildCardSelectorString)}, nil
}
return nil, err
}
return []coretypes.Selector{
resource.Type().MustSelector(quickFilter.ID.StringValue()),
resource.Type().MustSelector(coretypes.WildCardSelectorString),
}, nil
}

View File

@@ -22,6 +22,9 @@ type Gateway interface {
// Search Ingestion Keys by Name (this is supposed to be for the current user but for now in gateway code this is ignoring the consumer user)
SearchIngestionKeysByName(ctx context.Context, orgID valuer.UUID, name string, page, perPage int) (*gatewaytypes.GettableIngestionKeys, error)
// Get Ingestion Key
GetIngestionKey(ctx context.Context, orgID valuer.UUID, keyID string) (*gatewaytypes.IngestionKey, error)
// Create Ingestion Key
CreateIngestionKey(ctx context.Context, orgID valuer.UUID, name string, tags []string, expiresAt time.Time) (*gatewaytypes.GettableCreatedIngestionKey, error)
@@ -34,6 +37,12 @@ type Gateway interface {
// Create Ingestion Key Limit
CreateIngestionKeyLimit(ctx context.Context, orgID valuer.UUID, keyID string, signal string, limitConfig gatewaytypes.LimitConfig, tags []string) (*gatewaytypes.GettableCreatedIngestionKeyLimit, error)
// Get Ingestion Key Limit
GetIngestionKeyLimit(ctx context.Context, orgID valuer.UUID, limitID string) (*gatewaytypes.Limit, error)
// Get Ingestion Key Limits
GetIngestionKeyLimits(ctx context.Context, orgID valuer.UUID, keyID string) ([]gatewaytypes.Limit, error)
// Update Ingestion Key Limit
UpdateIngestionKeyLimit(ctx context.Context, orgID valuer.UUID, limitID string, limitConfig gatewaytypes.LimitConfig, tags []string) error
@@ -46,6 +55,8 @@ type Handler interface {
SearchIngestionKeys(http.ResponseWriter, *http.Request)
GetIngestionKey(http.ResponseWriter, *http.Request)
CreateIngestionKey(http.ResponseWriter, *http.Request)
UpdateIngestionKey(http.ResponseWriter, *http.Request)
@@ -54,7 +65,13 @@ type Handler interface {
CreateIngestionKeyLimit(http.ResponseWriter, *http.Request)
GetIngestionKeyLimit(http.ResponseWriter, *http.Request)
GetIngestionKeyLimits(http.ResponseWriter, *http.Request)
UpdateIngestionKeyLimit(http.ResponseWriter, *http.Request)
DeleteIngestionKeyLimit(http.ResponseWriter, *http.Request)
DeprecatedCreateIngestionKeyLimit(http.ResponseWriter, *http.Request)
}

View File

@@ -1,11 +1,11 @@
package gateway
import (
"encoding/json"
"net/http"
"strconv"
"github.com/SigNoz/signoz/pkg/errors"
"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/gatewaytypes"
@@ -111,8 +111,8 @@ func (handler *handler) CreateIngestionKey(rw http.ResponseWriter, r *http.Reque
orgID := valuer.MustNewUUID(claims.OrgID)
var req gatewaytypes.PostableIngestionKey
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid request body"))
if err := binding.JSON.BindBody(r.Body, &req); err != nil {
render.Error(rw, err)
return
}
@@ -143,8 +143,8 @@ func (handler *handler) UpdateIngestionKey(rw http.ResponseWriter, r *http.Reque
}
var req gatewaytypes.PostableIngestionKey
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid request body"))
if err := binding.JSON.BindBody(r.Body, &req); err != nil {
render.Error(rw, err)
return
}
@@ -183,7 +183,7 @@ func (handler *handler) DeleteIngestionKey(rw http.ResponseWriter, r *http.Reque
render.Success(rw, http.StatusNoContent, nil)
}
func (handler *handler) CreateIngestionKeyLimit(rw http.ResponseWriter, r *http.Request) {
func (handler *handler) DeprecatedCreateIngestionKeyLimit(rw http.ResponseWriter, r *http.Request) {
ctx := r.Context()
claims, err := authtypes.ClaimsFromContext(ctx)
@@ -200,9 +200,9 @@ func (handler *handler) CreateIngestionKeyLimit(rw http.ResponseWriter, r *http.
return
}
var req gatewaytypes.PostableIngestionKeyLimit
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid request body"))
var req gatewaytypes.DeprecatedPostableIngestionKeyLimit
if err := binding.JSON.BindBody(r.Body, &req); err != nil {
render.Error(rw, err)
return
}
@@ -233,8 +233,8 @@ func (handler *handler) UpdateIngestionKeyLimit(rw http.ResponseWriter, r *http.
}
var req gatewaytypes.UpdatableIngestionKeyLimit
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid request body"))
if err := binding.JSON.BindBody(r.Body, &req); err != nil {
render.Error(rw, err)
return
}
@@ -273,6 +273,115 @@ func (handler *handler) DeleteIngestionKeyLimit(rw http.ResponseWriter, r *http.
render.Success(rw, http.StatusNoContent, nil)
}
func (handler *handler) CreateIngestionKeyLimit(rw http.ResponseWriter, r *http.Request) {
ctx := r.Context()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID := valuer.MustNewUUID(claims.OrgID)
var req gatewaytypes.PostableIngestionKeyLimit
if err := binding.JSON.BindBody(r.Body, &req); err != nil {
render.Error(rw, err)
return
}
if req.KeyID == "" {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "keyId is required"))
return
}
response, err := handler.gateway.CreateIngestionKeyLimit(ctx, orgID, req.KeyID, req.Signal, req.Config, req.Tags)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusCreated, response)
}
func (handler *handler) GetIngestionKeyLimit(rw http.ResponseWriter, r *http.Request) {
ctx := r.Context()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID := valuer.MustNewUUID(claims.OrgID)
limitID := mux.Vars(r)["limitId"]
if limitID == "" {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "limitId is required"))
return
}
response, err := handler.gateway.GetIngestionKeyLimit(ctx, orgID, limitID)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, response)
}
func (handler *handler) GetIngestionKey(rw http.ResponseWriter, r *http.Request) {
ctx := r.Context()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID := valuer.MustNewUUID(claims.OrgID)
keyID := mux.Vars(r)["keyId"]
if keyID == "" {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "keyId is required"))
return
}
response, err := handler.gateway.GetIngestionKey(ctx, orgID, keyID)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, response)
}
func (handler *handler) GetIngestionKeyLimits(rw http.ResponseWriter, r *http.Request) {
ctx := r.Context()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID := valuer.MustNewUUID(claims.OrgID)
keyID := mux.Vars(r)["keyId"]
if keyID == "" {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "keyId is required"))
return
}
response, err := handler.gateway.GetIngestionKeyLimits(ctx, orgID, keyID)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, response)
}
func parseIntWithDefaultValue(value string, defaultValue int) (int, error) {
if value == "" {
return defaultValue, nil

View File

@@ -31,6 +31,10 @@ func (p *provider) SearchIngestionKeysByName(_ context.Context, _ valuer.UUID, _
return nil, errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}
func (p *provider) GetIngestionKey(_ context.Context, _ valuer.UUID, _ string) (*gatewaytypes.IngestionKey, error) {
return nil, errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}
func (p *provider) CreateIngestionKey(_ context.Context, _ valuer.UUID, _ string, _ []string, _ time.Time) (*gatewaytypes.GettableCreatedIngestionKey, error) {
return nil, errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}
@@ -47,6 +51,14 @@ func (p *provider) CreateIngestionKeyLimit(_ context.Context, _ valuer.UUID, _ s
return nil, errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}
func (p *provider) GetIngestionKeyLimit(_ context.Context, _ valuer.UUID, _ string) (*gatewaytypes.Limit, error) {
return nil, errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}
func (p *provider) GetIngestionKeyLimits(_ context.Context, _ valuer.UUID, _ string) ([]gatewaytypes.Limit, error) {
return nil, errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}
func (p *provider) UpdateIngestionKeyLimit(_ context.Context, _ valuer.UUID, _ string, _ gatewaytypes.LimitConfig, _ []string) error {
return errors.New(errors.TypeUnsupported, gateway.ErrCodeGatewayUnsupported, "unsupported call")
}

View File

@@ -2,10 +2,12 @@ package implaiobservability
import (
"context"
"log/slog"
"net/http"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/http/binding"
"github.com/SigNoz/signoz/pkg/http/render"
"github.com/SigNoz/signoz/pkg/modules/aiobservability"
@@ -17,11 +19,13 @@ import (
)
type handler struct {
settings factory.ScopedProviderSettings
telemetryMetadataStore telemetrytypes.MetadataStore
}
func NewHandler(telemetryMetadataStore telemetrytypes.MetadataStore) aiobservability.Handler {
func NewHandler(providerSettings factory.ProviderSettings, telemetryMetadataStore telemetrytypes.MetadataStore) aiobservability.Handler {
return &handler{
settings: factory.NewScopedProviderSettings(providerSettings, "github.com/SigNoz/signoz/pkg/modules/aiobservability/implaiobservability"),
telemetryMetadataStore: telemetryMetadataStore,
}
}
@@ -66,13 +70,6 @@ func (handler *handler) GetFieldsValues(rw http.ResponseWriter, req *http.Reques
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
defer cancel()
// binding ignores query params the struct does not declare, so an unsupported
// existingQuery would silently return values it did not narrow
if req.URL.Query().Has("existingQuery") {
render.Error(rw, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "existingQuery is not supported"))
return
}
var params aiobservabilitytypes.PostableFieldValueParams
if err := binding.Query.BindQuery(req.URL.Query(), &params); err != nil {
render.Error(rw, err)
@@ -84,18 +81,32 @@ func (handler *handler) GetFieldsValues(rw http.ResponseWriter, req *http.Reques
render.Error(rw, err)
return
}
orgID := valuer.MustNewUUID(claims.OrgID)
scopedQuery, err := aitelemetryschema.ScopedExistingQuery(params.ExistingQuery)
if err != nil {
handler.settings.Logger().WarnContext(ctx, "dropping unparseable existing query", slog.String("query", params.ExistingQuery), errors.Attr(err))
}
params.ExistingQuery = scopedQuery
fieldValueSelector := aiobservabilitytypes.NewFieldValueSelectorFromPostableFieldValueParams(params)
values := &telemetrytypes.TelemetryFieldValues{}
complete := true
// the trace context names the computed per-trace aggregates, which are never ingested
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextTrace {
values, complete, err = handler.telemetryMetadataStore.GetAllValues(ctx, valuer.MustNewUUID(claims.OrgID), fieldValueSelector)
values, complete, err = handler.telemetryMetadataStore.GetAllValues(ctx, orgID, fieldValueSelector)
if err != nil {
render.Error(rw, err)
return
}
// related values are best-effort: on failure the plain values still serve the filter bar
relatedValues, relatedComplete, err := handler.telemetryMetadataStore.GetRelatedValues(ctx, orgID, fieldValueSelector)
if err != nil {
relatedValues = []string{}
}
values.RelatedValues = relatedValues
complete = complete && relatedComplete
}
render.Success(rw, http.StatusOK, &telemetrytypes.GettableFieldValues{

View File

@@ -4,10 +4,13 @@ import (
"encoding/json"
"net/http"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/http/render"
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/types/authtypes"
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/gorilla/mux"
)
@@ -20,6 +23,13 @@ func NewHandler(module quickfilter.Module) quickfilter.Handler {
return &handler{module: module}
}
// legacySourceFilters is the v1 API shape: filters as v3 attribute keys,
// with the source still spelled "signal" on the wire.
type legacySourceFilters struct {
Source quickfiltertypes.Source `json:"signal"`
Filters []v3.AttributeKey `json:"filters"`
}
func (handler *handler) GetQuickFilters(rw http.ResponseWriter, r *http.Request) {
claims, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
@@ -27,13 +37,41 @@ func (handler *handler) GetQuickFilters(rw http.ResponseWriter, r *http.Request)
return
}
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID))
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), quickfiltertypes.Source{})
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, filters)
legacyFilters := make([]*legacySourceFilters, 0, len(filters))
for _, sourceFilters := range filters {
legacyFilters = append(legacyFilters, newLegacySourceFilters(sourceFilters))
}
render.Success(rw, http.StatusOK, legacyFilters)
}
func (handler *handler) GetSourceFilters(rw http.ResponseWriter, r *http.Request) {
claims, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
render.Error(rw, err)
return
}
source := mux.Vars(r)["signal"]
validatedSource, err := quickfiltertypes.NewSource(source)
if err != nil {
render.Error(rw, err)
return
}
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSource)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, newLegacySourceFilters(handler.sourceFiltersOrEmpty(filters, validatedSource)))
}
func (handler *handler) UpdateQuickFilters(rw http.ResponseWriter, r *http.Request) {
@@ -43,14 +81,19 @@ func (handler *handler) UpdateQuickFilters(rw http.ResponseWriter, r *http.Reque
return
}
var req quickfiltertypes.UpdatableQuickFilters
decodeErr := json.NewDecoder(r.Body).Decode(&req)
if decodeErr != nil {
render.Error(rw, decodeErr)
var req legacySourceFilters
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
render.Error(rw, err)
return
}
err = handler.module.UpdateQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), req.Signal, req.Filters)
fieldKeys, err := newTelemetryFieldKeysFromLegacy(req.Source, req.Filters)
if err != nil {
render.Error(rw, err)
return
}
err = handler.module.UpsertQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), req.Source, fieldKeys)
if err != nil {
render.Error(rw, err)
return
@@ -59,21 +102,14 @@ func (handler *handler) UpdateQuickFilters(rw http.ResponseWriter, r *http.Reque
render.Success(rw, http.StatusNoContent, nil)
}
func (handler *handler) GetSignalFilters(rw http.ResponseWriter, r *http.Request) {
func (handler *handler) ListQuickFiltersV2(rw http.ResponseWriter, r *http.Request) {
claims, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
render.Error(rw, err)
return
}
signal := mux.Vars(r)["signal"]
validatedSignal, err := quickfiltertypes.NewSignal(signal)
if err != nil {
render.Error(rw, err)
return
}
filters, err := handler.module.GetSignalFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSignal)
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), quickfiltertypes.Source{})
if err != nil {
render.Error(rw, err)
return
@@ -81,3 +117,141 @@ func (handler *handler) GetSignalFilters(rw http.ResponseWriter, r *http.Request
render.Success(rw, http.StatusOK, filters)
}
func (handler *handler) UpdateQuickFiltersV2(rw http.ResponseWriter, r *http.Request) {
claims, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
render.Error(rw, err)
return
}
source := mux.Vars(r)["source"]
validatedSource, err := quickfiltertypes.NewSource(source)
if err != nil {
render.Error(rw, err)
return
}
var req quickfiltertypes.UpdatableQuickFilters
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
render.Error(rw, err)
return
}
err = handler.module.UpsertQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSource, req.Filters)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusNoContent, nil)
}
func (handler *handler) GetQuickFiltersV2(rw http.ResponseWriter, r *http.Request) {
claims, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
render.Error(rw, err)
return
}
source := mux.Vars(r)["source"]
validatedSource, err := quickfiltertypes.NewSource(source)
if err != nil {
render.Error(rw, err)
return
}
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSource)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, handler.sourceFiltersOrEmpty(filters, validatedSource))
}
// sourceFiltersOrEmpty keeps the single-source response contract: a source
// with no stored filters is served as an empty filter list, not an error.
func (handler *handler) sourceFiltersOrEmpty(filters []*quickfiltertypes.SourceFilters, source quickfiltertypes.Source) *quickfiltertypes.SourceFilters {
if len(filters) == 0 {
return quickfiltertypes.NewSourceFiltersFromSource(source)
}
return filters[0]
}
// newTelemetryFieldKeysFromLegacy converts a v1 write payload with the same
// normalizations as the storage migration: alias contexts, numerics to number.
// The v1 shape carries no per filter signal, so meter keys get it restored.
func newTelemetryFieldKeysFromLegacy(source quickfiltertypes.Source, filters []v3.AttributeKey) ([]telemetrytypes.TelemetryFieldKey, error) {
var fieldSignal telemetrytypes.Signal
if source == quickfiltertypes.SourceMeter {
fieldSignal = telemetrytypes.SignalMetrics
}
fieldKeys := make([]telemetrytypes.TelemetryFieldKey, 0, len(filters))
for _, filter := range filters {
if err := filter.Validate(); err != nil {
return nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter: %v", err)
}
fieldContext, ok := telemetrytypes.FieldContextFromText(string(filter.Type))
if !ok {
fieldContext = telemetrytypes.FieldContextUnspecified
}
var fieldDataType telemetrytypes.FieldDataType
if err := fieldDataType.Scan(string(filter.DataType)); err != nil {
fieldDataType = telemetrytypes.FieldDataTypeUnspecified
}
if fieldDataType == telemetrytypes.FieldDataTypeInt64 {
fieldDataType = telemetrytypes.FieldDataTypeNumber
}
fieldKeys = append(fieldKeys, telemetrytypes.TelemetryFieldKey{
Name: filter.Key,
Signal: fieldSignal,
FieldContext: fieldContext,
FieldDataType: fieldDataType,
})
}
return fieldKeys, nil
}
// newLegacySourceFilters renders stored telemetry field keys
// back into the v1 shape, restoring the legacy spellings v1 clients expect.
func newLegacySourceFilters(sourceFilters *quickfiltertypes.SourceFilters) *legacySourceFilters {
filters := make([]v3.AttributeKey, 0, len(sourceFilters.Filters))
for _, fieldKey := range sourceFilters.Filters {
// Only tag and resource exist in the v3 enum; other contexts render as
// unspecified so v1 clients never see spellings their queries can't use.
var attributeType v3.AttributeKeyType
switch fieldKey.FieldContext {
case telemetrytypes.FieldContextAttribute:
attributeType = v3.AttributeKeyTypeTag
case telemetrytypes.FieldContextResource:
attributeType = v3.AttributeKeyTypeResource
default:
attributeType = v3.AttributeKeyTypeUnspecified
}
var dataType v3.AttributeKeyDataType
switch fieldKey.FieldDataType {
case telemetrytypes.FieldDataTypeNumber:
dataType = v3.AttributeKeyDataTypeFloat64
default:
dataType = v3.AttributeKeyDataType(fieldKey.FieldDataType.StringValue())
}
filters = append(filters, v3.AttributeKey{
Key: fieldKey.Name,
Type: attributeType,
DataType: dataType,
})
}
return &legacySourceFilters{
Source: sourceFilters.Source,
Filters: filters,
}
}

View File

@@ -0,0 +1,62 @@
package implquickfilter
import (
"testing"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestNewTelemetryFieldKeysFromLegacy(t *testing.T) {
fieldKeys, err := newTelemetryFieldKeysFromLegacy(quickfiltertypes.SourceTraces, []v3.AttributeKey{
{Key: "service.name", Type: v3.AttributeKeyTypeResource, DataType: v3.AttributeKeyDataTypeString},
{Key: "http.method", Type: v3.AttributeKeyTypeTag, DataType: v3.AttributeKeyDataTypeString},
{Key: "duration_nano", Type: v3.AttributeKeyTypeTag, DataType: v3.AttributeKeyDataTypeFloat64},
{Key: "code_line", Type: v3.AttributeKeyTypeTag, DataType: v3.AttributeKeyDataTypeInt64},
})
require.NoError(t, err)
require.Len(t, fieldKeys, 4)
assert.Equal(t, telemetrytypes.TelemetryFieldKey{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString}, fieldKeys[0])
assert.Equal(t, telemetrytypes.FieldContextAttribute, fieldKeys[1].FieldContext)
assert.Equal(t, telemetrytypes.FieldDataTypeNumber, fieldKeys[2].FieldDataType)
assert.Equal(t, telemetrytypes.FieldDataTypeNumber, fieldKeys[3].FieldDataType)
t.Run("meter writes restore the per-filter telemetry signal", func(t *testing.T) {
fieldKeys, err := newTelemetryFieldKeysFromLegacy(quickfiltertypes.SourceMeter, []v3.AttributeKey{
{Key: "host.name", DataType: v3.AttributeKeyDataTypeString},
})
require.NoError(t, err)
require.Len(t, fieldKeys, 1)
assert.Equal(t, telemetrytypes.SignalMetrics, fieldKeys[0].Signal)
})
t.Run("rejects a filter without a key", func(t *testing.T) {
_, err := newTelemetryFieldKeysFromLegacy(quickfiltertypes.SourceTraces, []v3.AttributeKey{{DataType: v3.AttributeKeyDataTypeString}})
require.Error(t, err)
})
}
func TestNewLegacySourceFilters(t *testing.T) {
legacy := newLegacySourceFilters(&quickfiltertypes.SourceFilters{
Source: quickfiltertypes.SourceLogs,
Filters: []telemetrytypes.TelemetryFieldKey{
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "http.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "duration_nano", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeNumber},
{Name: "severity_text", FieldContext: telemetrytypes.FieldContextLog, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "host.name", Signal: telemetrytypes.SignalMetrics},
},
})
assert.Equal(t, quickfiltertypes.SourceLogs, legacy.Source)
require.Len(t, legacy.Filters, 5)
assert.Equal(t, v3.AttributeKey{Key: "service.name", Type: v3.AttributeKeyTypeResource, DataType: v3.AttributeKeyDataTypeString}, legacy.Filters[0])
assert.Equal(t, v3.AttributeKeyTypeTag, legacy.Filters[1].Type)
assert.Equal(t, v3.AttributeKeyDataTypeFloat64, legacy.Filters[2].DataType)
assert.Equal(t, v3.AttributeKeyTypeUnspecified, legacy.Filters[3].Type, "contexts outside the v3 enum must render as unspecified")
assert.Equal(t, v3.AttributeKey{Key: "host.name"}, legacy.Filters[4])
}

View File

@@ -2,12 +2,11 @@ package implquickfilter
import (
"context"
"encoding/json"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
@@ -19,91 +18,54 @@ func NewModule(store quickfiltertypes.QuickFilterStore) quickfilter.Module {
return &module{store: store}
}
// GetQuickFilters returns all quick filters for an organization.
func (module *module) GetQuickFilters(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes.SignalFilters, error) {
storedFilters, err := module.store.Get(ctx, orgID)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error fetching organization filters")
}
result := make([]*quickfiltertypes.SignalFilters, 0, len(storedFilters))
for _, storedFilter := range storedFilters {
signalFilter, err := quickfiltertypes.NewSignalFilterFromStorableQuickFilter(storedFilter)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for signal: %s", storedFilter.Signal)
}
result = append(result, signalFilter)
}
return result, nil
func (module *module) Get(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) (*quickfiltertypes.StorableQuickFilter, error) {
return module.store.GetBySource(ctx, orgID, source.StringValue())
}
// GetSignalFilters returns quick filters for a specific signal in an organization.
func (m *module) GetSignalFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal) (*quickfiltertypes.SignalFilters, error) {
storedFilter, err := m.store.GetBySignal(ctx, orgID, signal.StringValue())
// GetQuickFilters returns quick filters for a source, or for every source when source is zero.
func (module *module) GetQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) ([]*quickfiltertypes.SourceFilters, error) {
if source.IsZero() {
storedFilters, err := module.store.Get(ctx, orgID)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error fetching organization filters")
}
result := make([]*quickfiltertypes.SourceFilters, 0, len(storedFilters))
for _, storedFilter := range storedFilters {
sourceFilter, err := quickfiltertypes.NewSourceFilterFromStorableQuickFilter(storedFilter)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for source: %s", storedFilter.Source)
}
result = append(result, sourceFilter)
}
return result, nil
}
storedFilter, err := module.store.GetBySource(ctx, orgID, source.StringValue())
if err != nil {
if errors.Ast(err, errors.TypeNotFound) {
return []*quickfiltertypes.SourceFilters{}, nil
}
return nil, err
}
// If no filter exists for this signal, return empty filters with the requested signal
if storedFilter == nil {
return &quickfiltertypes.SignalFilters{
Signal: signal,
Filters: []v3.AttributeKey{},
}, nil
}
// Convert stored filter to signal filter
signalFilter, err := quickfiltertypes.NewSignalFilterFromStorableQuickFilter(storedFilter)
sourceFilter, err := quickfiltertypes.NewSourceFilterFromStorableQuickFilter(storedFilter)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for signal: %s", storedFilter.Signal)
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for source: %s", storedFilter.Source)
}
return signalFilter, nil
return []*quickfiltertypes.SourceFilters{sourceFilter}, nil
}
// UpdateQuickFilters updates quick filters for a specific signal in an organization.
func (module *module) UpdateQuickFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal, filters []v3.AttributeKey) error {
// Validate each filter
for _, filter := range filters {
if err := filter.Validate(); err != nil {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter: %v", err)
}
}
// Marshal filters to JSON
filterJSON, err := json.Marshal(filters)
// UpsertQuickFilters replaces quick filters for a specific source in an organization, creating them if absent.
func (module *module) UpsertQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source, filters []telemetrytypes.TelemetryFieldKey) error {
filter, err := quickfiltertypes.NewStorableQuickFilter(orgID, source, filters)
if err != nil {
return errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error marshalling filters")
}
// Check if filter exists
existingFilter, err := module.store.GetBySignal(ctx, orgID, signal.StringValue())
if err != nil {
return errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error checking existing filters")
}
var filter *quickfiltertypes.StorableQuickFilter
if existingFilter != nil {
// Update in place
if err := existingFilter.Update(filterJSON); err != nil {
return errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "error updating existing filter")
}
filter = existingFilter
} else {
// Create new
filter, err = quickfiltertypes.NewStorableQuickFilter(orgID, signal, filterJSON)
if err != nil {
return errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "error creating new filter")
}
}
// Persist filter
if err := module.store.Upsert(ctx, filter); err != nil {
return err
}
return nil
return module.store.Upsert(ctx, filter)
}
func (module *module) SetDefaultConfig(ctx context.Context, orgID valuer.UUID) error {

View File

@@ -26,7 +26,7 @@ func (s *store) Get(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes
NewSelect().
Model(&filters).
Where("org_id = ?", orgID).
Order("signal ASC").
Order("source ASC").
Scan(ctx)
if err != nil {
@@ -36,7 +36,7 @@ func (s *store) Get(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes
return filters, nil
}
func (s *store) GetBySignal(ctx context.Context, orgID valuer.UUID, signal string) (*quickfiltertypes.StorableQuickFilter, error) {
func (s *store) GetBySource(ctx context.Context, orgID valuer.UUID, source string) (*quickfiltertypes.StorableQuickFilter, error) {
filter := new(quickfiltertypes.StorableQuickFilter)
err := s.store.
@@ -44,12 +44,12 @@ func (s *store) GetBySignal(ctx context.Context, orgID valuer.UUID, signal strin
NewSelect().
Model(filter).
Where("org_id = ?", orgID).
Where("signal = ?", signal).
Where("source = ?", source).
Scan(ctx)
if err != nil {
if err == sql.ErrNoRows {
return nil, s.store.WrapNotFoundErrf(err, errors.CodeNotFound, "No rows found for org_id: "+orgID.StringValue()+" signal: "+signal)
return nil, s.store.WrapNotFoundErrf(err, errors.CodeNotFound, "No rows found for org_id: "+orgID.StringValue()+" source: "+source)
}
return nil, err
}
@@ -62,7 +62,7 @@ func (s *store) Upsert(ctx context.Context, filter *quickfiltertypes.StorableQui
BunDB().
NewInsert().
Model(filter).
On("CONFLICT (id) DO UPDATE").
On("CONFLICT (org_id, source) DO UPDATE").
Set("filter = EXCLUDED.filter").
Set("updated_at = EXCLUDED.updated_at").
Exec(ctx)
@@ -78,7 +78,7 @@ func (s *store) Create(ctx context.Context, filters []*quickfiltertypes.Storable
BunDBCtx(ctx).
NewInsert().
Model(&filters).
On("CONFLICT (org_id, signal) DO NOTHING").
On("CONFLICT (org_id, source) DO NOTHING").
Exec(ctx)
if err != nil {

View File

@@ -4,20 +4,27 @@ import (
"context"
"net/http"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
type Module interface {
GetQuickFilters(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes.SignalFilters, error)
UpdateQuickFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal, filters []v3.AttributeKey) error
GetSignalFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal) (*quickfiltertypes.SignalFilters, error)
// Get returns the stored quick filter row for a source.
Get(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) (*quickfiltertypes.StorableQuickFilter, error)
// GetQuickFilters returns quick filters for a source, or for every source when source is zero.
GetQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) ([]*quickfiltertypes.SourceFilters, error)
UpsertQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source, filters []telemetrytypes.TelemetryFieldKey) error
SetDefaultConfig(ctx context.Context, orgID valuer.UUID) error
}
type Handler interface {
// Legacy v1 endpoints, served by converting to and from the v3 attribute key shape.
GetQuickFilters(http.ResponseWriter, *http.Request)
UpdateQuickFilters(http.ResponseWriter, *http.Request)
GetSignalFilters(http.ResponseWriter, *http.Request)
GetSourceFilters(http.ResponseWriter, *http.Request)
ListQuickFiltersV2(http.ResponseWriter, *http.Request)
GetQuickFiltersV2(http.ResponseWriter, *http.Request)
UpdateQuickFiltersV2(http.ResponseWriter, *http.Request)
}

View File

@@ -451,9 +451,9 @@ func (aH *APIHandler) RegisterRoutes(router *mux.Router, am *middleware.AuthZ) {
router.HandleFunc("/api/v1/disks", am.ViewAccess(aH.getDisks)).Methods(http.MethodGet)
// Quick Filters
// Quick Filters (v1 routes serve the legacy v3 shape; v2 lives in signozapiserver)
router.HandleFunc("/api/v1/orgs/me/filters", am.ViewAccess(aH.Signoz.Handlers.QuickFilter.GetQuickFilters)).Methods(http.MethodGet)
router.HandleFunc("/api/v1/orgs/me/filters/{signal}", am.ViewAccess(aH.Signoz.Handlers.QuickFilter.GetSignalFilters)).Methods(http.MethodGet)
router.HandleFunc("/api/v1/orgs/me/filters/{signal}", am.ViewAccess(aH.Signoz.Handlers.QuickFilter.GetSourceFilters)).Methods(http.MethodGet)
router.HandleFunc("/api/v1/orgs/me/filters", am.AdminAccess(aH.Signoz.Handlers.QuickFilter.UpdateQuickFilters)).Methods(http.MethodPut)
router.HandleFunc("/api/v1/register", am.OpenAccess(aH.registerUser)).Methods(http.MethodPost)

View File

@@ -333,7 +333,7 @@ func (m *Manager) validateChannels(ctx context.Context, orgID string, rule *rule
known := make(map[string]struct{}, len(orgChannels))
for _, ch := range orgChannels {
known[ch.Name] = struct{}{}
known[ch.DisplayName] = struct{}{}
}
var unknown []string

View File

@@ -127,7 +127,7 @@ func NewHandlers(
FlaggerHandler: flagger.NewHandler(flaggerService),
GatewayHandler: gateway.NewHandler(gatewayService),
Fields: implfields.NewHandler(providerSettings, telemetryMetadataStore),
AIObservability: implaiobservability.NewHandler(telemetryMetadataStore),
AIObservability: implaiobservability.NewHandler(providerSettings, telemetryMetadataStore),
AuthzHandler: signozauthzapi.NewHandler(authz),
ZeusHandler: zeus.NewHandler(zeusService, licensingService),
LicensingHandler: licensing.NewHandler(licensingService),

View File

@@ -30,6 +30,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/organization"
"github.com/SigNoz/signoz/pkg/modules/preference"
"github.com/SigNoz/signoz/pkg/modules/promote"
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
"github.com/SigNoz/signoz/pkg/modules/rawdataexport"
"github.com/SigNoz/signoz/pkg/modules/rulestatehistory"
"github.com/SigNoz/signoz/pkg/modules/savedview"
@@ -97,6 +98,8 @@ func NewOpenAPI(ctx context.Context, instrumentation instrumentation.Instrumenta
struct{ ruler.Handler }{},
struct{ statsreporter.Handler }{},
struct{ savedview.Handler }{},
struct{ quickfilter.Module }{},
struct{ quickfilter.Handler }{},
).New(ctx, instrumentation.ToProviderSettings(), apiserver.Config{})
if err != nil {
return nil, err

View File

@@ -247,6 +247,9 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddDeploymentHostTuplesFactory(sqlstore),
sqlmigration.NewAddSystemDashboardFactory(sqlstore, sqlschema),
sqlmigration.NewAddLicenseTuplesFactory(sqlstore),
sqlmigration.NewAddChannelDisplayNameFactory(sqlstore, sqlschema),
sqlmigration.NewMigrateQuickFiltersFactory(sqlstore),
sqlmigration.NewAddQuickFilterTuplesFactory(sqlstore),
)
}
@@ -352,6 +355,8 @@ func NewAPIServerProviderFactories(orgGetter organization.Getter, authz authz.Au
handlers.RulerHandler,
handlers.StatsHandler,
handlers.SavedView,
modules.QuickFilter,
handlers.QuickFilter,
),
)
}

View File

@@ -153,8 +153,23 @@ func (migration *addAlertmanager) populateOrgIDInChannels(ctx context.Context, t
return nil
}
// storableChannel is the channel table as it stood at this migration.
// alertmanagertypes.Channel has grown columns since, and reading through it here
// would put those columns into statements that run before they exist.
type storableChannel struct {
bun.BaseModel `bun:"table:notification_channel"`
ID string `bun:"id,pk,type:text"`
CreatedAt time.Time `bun:"created_at"`
UpdatedAt time.Time `bun:"updated_at"`
Name string `bun:"name"`
Type string `bun:"type"`
Data string `bun:"data"`
OrgID string `bun:"org_id"`
}
func (migration *addAlertmanager) populateAlertmanagerConfig(ctx context.Context, tx bun.Tx, orgID string) error {
var channels []*alertmanagertypes.Channel
var channels []*storableChannel
err := tx.
NewSelect().
@@ -205,7 +220,13 @@ func (migration *addAlertmanager) populateAlertmanagerConfig(ctx context.Context
}
}
config, err := alertmanagertypes.NewConfigFromChannels(alertmanagerserver.NewConfig().Global, alertmanagerserver.NewConfig().Route, channels, orgID)
// NewConfigFromChannels reads nothing off a channel but Data.
receiverChannels := make(alertmanagertypes.Channels, 0, len(channels))
for _, channel := range channels {
receiverChannels = append(receiverChannels, &alertmanagertypes.Channel{Data: channel.Data})
}
config, err := alertmanagertypes.NewConfigFromChannels(alertmanagerserver.NewConfig().Global, alertmanagerserver.NewConfig().Route, receiverChannels, orgID)
if err != nil {
return err
}
@@ -273,7 +294,7 @@ func newReceiver(input string) (config.Receiver, error) {
return receiverWithDefaults, nil
}
func (migration *addAlertmanager) msTeamsChannelToMSTeamsV2Channel(c *alertmanagertypes.Channel) error {
func (migration *addAlertmanager) msTeamsChannelToMSTeamsV2Channel(c *storableChannel) error {
if c.Type != "msteams" {
return nil
}

View File

@@ -0,0 +1,172 @@
package sqlmigration
import (
"context"
"crypto/rand"
"strings"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlschema"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
)
type addChannelDisplayName struct {
sqlstore sqlstore.SQLStore
sqlschema sqlschema.SQLSchema
}
func NewAddChannelDisplayNameFactory(sqlstore sqlstore.SQLStore, sqlschema sqlschema.SQLSchema) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(
factory.MustNewName("channel_display_name"),
func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addChannelDisplayName{sqlstore: sqlstore, sqlschema: sqlschema}, nil
},
)
}
func (migration *addChannelDisplayName) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
// Up moves the free-text name onto display_name and gives name the DNS1123
// identity, matching organizations and dashboards. Nothing outside the channel
// reads the new name yet: routing policies and rules still reference
// display_name, and the alertmanager config still keys receivers by that same
// string.
func (migration *addChannelDisplayName) Up(ctx context.Context, db *bun.DB) error {
// Adding a NOT NULL column rebuilds the whole table on SQLite (create temp,
// copy, drop, rename), and notification_channel has a foreign key on org_id
// that the copy would re-validate. Enforcement stays off for the rebuild.
if err := migration.sqlschema.ToggleFKEnforcement(ctx, db, false); err != nil {
return err
}
table, uniqueConstraints, err := migration.sqlschema.GetTable(ctx, sqlschema.TableName("notification_channel"))
if err != nil {
return err
}
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() {
_ = tx.Rollback()
}()
if _, err := migration.sqlstore.Dialect().RenameColumn(ctx, tx, "notification_channel", "name", "display_name"); err != nil {
return err
}
// The table was inspected before the rename, and the recreate-table fallback
// below rebuilds the table from this description, so it has to follow.
for _, column := range table.Columns {
if column.Name == sqlschema.ColumnName("name") {
column.Name = sqlschema.ColumnName("display_name")
}
}
nameColumn := &sqlschema.Column{
Name: sqlschema.ColumnName("name"),
DataType: sqlschema.DataTypeText,
Nullable: false,
}
sqls := migration.sqlschema.Operator().AddColumn(table, uniqueConstraints, nameColumn, "")
for _, sql := range sqls {
if _, err := tx.ExecContext(ctx, string(sql)); err != nil {
return err
}
}
type channel struct {
bun.BaseModel `bun:"table:notification_channel"`
ID valuer.UUID `bun:"id,pk"`
DisplayName string `bun:"display_name"`
}
var channels []channel
if err := tx.
NewSelect().
Model(&channels).
Column("id", "display_name").
Scan(ctx); err != nil {
return err
}
for _, existing := range channels {
if _, err := tx.
NewUpdate().
Model((*channel)(nil)).
Set("name = ?", slugifyChannelName(existing.DisplayName)).
Where("id = ?", existing.ID).
Exec(ctx); err != nil {
return err
}
}
indexSQLs := migration.sqlschema.Operator().CreateIndex(&sqlschema.UniqueIndex{
TableName: "notification_channel",
ColumnNames: []sqlschema.ColumnName{"org_id", "name"},
})
for _, sql := range indexSQLs {
if _, err := tx.ExecContext(ctx, string(sql)); err != nil {
return err
}
}
if err := tx.Commit(); err != nil {
return err
}
return migration.sqlschema.ToggleFKEnforcement(ctx, db, true)
}
func (migration *addChannelDisplayName) Down(context.Context, *bun.DB) error {
return nil
}
const migrationChannelNameSuffixLen = 8
// slugifyChannelName is a copy of dashboardtypes.generateDashboardName. The
// random suffix is what makes the unique index safe to add without a collision
// loop over the existing free-text names.
func slugifyChannelName(displayName string) string {
const dns1123LabelMaxLen = 63
suffixAlphabet := []byte("abcdefghijklmnopqrstuvwxyz0123456789")
var b strings.Builder
b.Grow(len(displayName))
prevHyphen := false
for _, r := range strings.ToLower(displayName) {
switch {
case (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9'):
b.WriteRune(r)
prevHyphen = false
case b.Len() > 0 && !prevHyphen:
b.WriteByte('-')
prevHyphen = true
}
}
prefix := strings.TrimRight(b.String(), "-")
suffix := make([]byte, migrationChannelNameSuffixLen)
if _, err := rand.Read(suffix); err != nil {
panic(err)
}
for i := range suffix {
suffix[i] = suffixAlphabet[int(suffix[i])%len(suffixAlphabet)]
}
maxPrefix := dns1123LabelMaxLen - 1 - migrationChannelNameSuffixLen
if len(prefix) > maxPrefix {
prefix = strings.TrimRight(prefix[:maxPrefix], "-")
}
if prefix == "" {
return string(suffix)
}
return prefix + "-" + string(suffix)
}

View File

@@ -0,0 +1,180 @@
package sqlmigration
import (
"context"
"encoding/json"
"log/slog"
"strings"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlstore"
)
type storableQuickFilterRow struct {
bun.BaseModel `bun:"table:quick_filter"`
ID string `bun:"id,pk"`
Filter string `bun:"filter"`
}
// legacyQuickFilterEntry carries both shapes a stored entry can be in: the
// legacy key/type/dataType shape and the current name-carrying shape.
type legacyQuickFilterEntry struct {
Name string `json:"name"`
Key string `json:"key"`
Type string `json:"type"`
DataType string `json:"dataType"`
Signal string `json:"signal"`
}
// quickFilterLegacyTypeToFieldContext maps the v3 attribute key types the v1
// write path could store. Materialized top-level fields carried no type, and
// anything unknown (e.g. "Sum" in the old meter defaults) normalizes to
// unspecified, matching what the v1 write path does at runtime.
var quickFilterLegacyTypeToFieldContext = map[string]string{
"tag": "attribute",
"resource": "resource",
"scope": "scope",
}
// quickFilterLegacyDataTypeToFieldDataType maps the v3 attribute key data
// types the v1 write path could store, with numerics collapsed to number,
// matching the fields API and the v1 write path.
var quickFilterLegacyDataTypeToFieldDataType = map[string]string{
"string": "string",
"bool": "bool",
"int64": "number",
"float64": "number",
}
type migrateQuickFilters struct {
sqlstore sqlstore.SQLStore
settings factory.ProviderSettings
}
func NewMigrateQuickFiltersFactory(sqlstore sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(factory.MustNewName("migrate_quick_filters"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &migrateQuickFilters{sqlstore: sqlstore, settings: ps}, nil
})
}
func (migration *migrateQuickFilters) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func (migration *migrateQuickFilters) Up(ctx context.Context, db *bun.DB) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
var rows []*storableQuickFilterRow
if err := tx.NewSelect().Model(&rows).Scan(ctx); err != nil {
return err
}
var migrated, skipped int
for _, row := range rows {
migratedFilter, changed, ok := migrateQuickFilterEntries(row.Filter)
if !ok {
migration.settings.Logger.WarnContext(ctx, "quick filter could not be parsed, leaving it untouched", slog.String("quick_filter_id", row.ID), slog.String("raw_filter", row.Filter))
skipped++
continue
}
if !changed {
continue
}
migrated++
if _, err := tx.NewUpdate().Model((*storableQuickFilterRow)(nil)).Set("filter = ?", migratedFilter).Where("id = ?", row.ID).Exec(ctx); err != nil {
return err
}
}
migration.settings.Logger.InfoContext(ctx, "migrated quick filters to telemetry field keys", slog.Int("total", len(rows)), slog.Int("migrated", migrated), slog.Int("skipped", skipped))
if _, err := migration.sqlstore.Dialect().RenameColumn(ctx, tx, "quick_filter", "signal", "source"); err != nil {
return err
}
for _, column := range []string{"created_by", "updated_by"} {
if err := migration.sqlstore.Dialect().DropColumn(ctx, tx, "quick_filter", column); err != nil {
return err
}
}
return tx.Commit()
}
func (migration *migrateQuickFilters) Down(context.Context, *bun.DB) error {
return nil
}
// migrateQuickFilterEntries rewrites a stored filter list from the legacy
// key/dataType/type shape to telemetry field keys; ok=false means unparseable.
func migrateQuickFilterEntries(filter string) (migrated string, changed bool, ok bool) {
var entriesRaw []json.RawMessage
if err := json.Unmarshal([]byte(filter), &entriesRaw); err != nil {
return "", false, false
}
migratedEntries := make([]json.RawMessage, 0, len(entriesRaw))
for _, rawEntry := range entriesRaw {
var entry legacyQuickFilterEntry
if err := json.Unmarshal(rawEntry, &entry); err != nil {
// Some stored entries are plain strings rather than objects; treat
// the string as the filter key name, dropping empty ones.
var name string
if err := json.Unmarshal(rawEntry, &name); err != nil {
return "", false, false
}
entry = legacyQuickFilterEntry{Key: name}
}
switch {
case entry.Name != "":
migratedEntries = append(migratedEntries, rawEntry)
case entry.Key != "":
migratedJSON, err := marshalUnescaped(telemetryFieldKeyOutput{
Name: entry.Key,
Signal: entry.Signal,
FieldContext: quickFilterFieldContext(entry.Type),
FieldDataType: quickFilterFieldDataType(entry.DataType),
})
if err != nil {
return "", false, false
}
migratedEntries = append(migratedEntries, migratedJSON)
changed = true
default:
changed = true
}
}
if !changed {
return "", false, true
}
migratedJSON, err := marshalUnescaped(migratedEntries)
if err != nil {
return "", false, false
}
return string(migratedJSON), true, true
}
// quickFilterFieldDataType resolves legacy datatype spellings, with unknowns
// normalized to unspecified.
func quickFilterFieldDataType(legacyDataType string) string {
return quickFilterLegacyDataTypeToFieldDataType[strings.ToLower(strings.TrimSpace(legacyDataType))]
}
// quickFilterFieldContext resolves legacy type spellings, with unknowns
// normalized to unspecified.
func quickFilterFieldContext(legacyType string) string {
return quickFilterLegacyTypeToFieldContext[strings.ToLower(strings.TrimSpace(legacyType))]
}

View File

@@ -0,0 +1,139 @@
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 addQuickFilterTuples struct {
sqlstore sqlstore.SQLStore
}
func NewAddQuickFilterTuplesFactory(sqlstore sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(factory.MustNewName("add_quick_filter_tuples"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addQuickFilterTuples{sqlstore: sqlstore}, nil
})
}
func (migration *addQuickFilterTuples) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func (migration *addQuickFilterTuples) 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
// quick-filter moved from the legacy ViewAccess/AdminAccess role gate to
// CheckResources, which on enterprise requires real tuples -- existing orgs
// never had these written, only new orgs get them from the registry at bootstrap.
tuples := []migrationTuple{
{authtypes.SigNozAdminRoleName, "metaresource", "quick-filter", "read"},
{authtypes.SigNozAdminRoleName, "metaresource", "quick-filter", "update"},
{authtypes.SigNozAdminRoleName, "metaresource", "quick-filter", "list"},
{authtypes.SigNozEditorRoleName, "metaresource", "quick-filter", "read"},
{authtypes.SigNozEditorRoleName, "metaresource", "quick-filter", "list"},
{authtypes.SigNozViewerRoleName, "metaresource", "quick-filter", "read"},
{authtypes.SigNozViewerRoleName, "metaresource", "quick-filter", "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 *addQuickFilterTuples) Down(context.Context, *bun.DB) error {
return nil
}

View File

@@ -0,0 +1,133 @@
package aistatementbuilder
import (
"context"
"testing"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// Mixed filter: the span-level part gates the scan, the trace-level part becomes the
// __trace_scope qualification.
func TestBuild_FullSQL_SpanList_TraceScoped(t *testing.T) {
b := newTestBuilder(t)
stmt, err := b.Build(context.Background(), valuer.UUID{}, testStartMs, testEndMs, qbtypes.RequestTypeRaw,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
Filter: &qbtypes.Filter{Expression: "gen_ai.request.model = 'gpt-4o-mini' AND trace.output_tokens > 1000"},
Limit: 10,
}, nil)
require.NoError(t, err)
assertSQLEqual(t, `
WITH __trace_scope AS (
SELECT trace_id,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens'), toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens
FROM signoz_traces.distributed_signoz_index_v3
WHERE timestamp >= '1747947419000000000'
AND timestamp < '1747983448000000000'
AND ts_bucket_start >= 1747945619
AND ts_bucket_start <= 1747983448
AND (mapContains(attributes_string, 'gen_ai.request.model') OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name'))
GROUP BY trace_id
HAVING output_tokens > 1000
)
SELECT timestamp AS __SELECT_KEY_0_timestamp, trace_id AS __SELECT_KEY_1_trace_id, span_id AS __SELECT_KEY_2_span_id,
trace_state AS __SELECT_KEY_3_trace_state, parent_span_id AS __SELECT_KEY_4_parent_span_id, flags AS __SELECT_KEY_5_flags,
name AS __SELECT_KEY_6_name, kind AS __SELECT_KEY_7_kind, kind_string AS __SELECT_KEY_8_kind_string, duration_nano AS __SELECT_KEY_9_duration_nano,
status_code AS __SELECT_KEY_10_status_code, status_message AS __SELECT_KEY_11_status_message,
status_code_string AS __SELECT_KEY_12_status_code_string, events AS __SELECT_KEY_13_events, links AS __SELECT_KEY_14_links,
response_status_code AS __SELECT_KEY_15_response_status_code, external_http_url AS __SELECT_KEY_16_external_http_url,
http_url AS __SELECT_KEY_17_http_url, external_http_method AS __SELECT_KEY_18_external_http_method,
http_method AS __SELECT_KEY_19_http_method, http_host AS __SELECT_KEY_20_http_host, db_name AS __SELECT_KEY_21_db_name,
db_operation AS __SELECT_KEY_22_db_operation, has_error AS __SELECT_KEY_23_has_error, is_remote AS __SELECT_KEY_24_is_remote,
attributes_string, attributes_number, attributes_bool, resources_string
FROM signoz_traces.distributed_signoz_index_v3
WHERE trace_id GLOBAL IN (SELECT trace_id FROM __trace_scope)
AND (((mapContains(attributes_string, 'gen_ai.request.model')
OR mapContains(attributes_string, 'gen_ai.tool.name')
OR mapContains(attributes_string, 'gen_ai.agent.name')))
AND ((attributes_string['gen_ai.request.model'] = 'gpt-4o-mini'
AND mapContains(attributes_string, 'gen_ai.request.model'))))
AND timestamp >= '1747947419000000000'
AND timestamp < '1747983448000000000'
AND ts_bucket_start >= 1747945619
AND ts_bucket_start <= 1747983448
LIMIT 10
`, stmt)
}
// A resource attribute mixed with a trace-level condition: the resource part flows
// through the fingerprint machinery, the trace-level part through __trace_scope.
func TestBuild_SpanList_ResourcePlusTraceFilter(t *testing.T) {
b := newTestBuilder(t)
stmt, err := b.Build(context.Background(), valuer.UUID{}, testStartMs, testEndMs, qbtypes.RequestTypeRaw,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
Filter: &qbtypes.Filter{Expression: "resource.service.name = 'checkout' AND trace.output_tokens > 1000"},
Limit: 10,
}, nil)
require.NoError(t, err)
got := renderSQL(t, stmt)
assert.Contains(t, got, "__resource_filter AS (")
assert.Contains(t, got, "resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter)")
assert.Contains(t, got, "__trace_scope AS (")
assert.Contains(t, got, "trace_id GLOBAL IN (SELECT trace_id FROM __trace_scope)")
assert.Contains(t, got, "HAVING output_tokens > 1000")
}
// Trace-level order keys are rejected — known aggregate alias or not — while a bare
// span column sharing an alias (duration_nano) stays orderable.
func TestBuild_SpanList_OrderKeyValidation(t *testing.T) {
b := newTestBuilder(t)
build := func(q qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]) error {
q.Signal = telemetrytypes.SignalTraces
_, err := b.Build(context.Background(), valuer.UUID{}, testStartMs, testEndMs, qbtypes.RequestTypeRaw, q, nil)
return err
}
err := build(qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Order: []qbtypes.OrderBy{{Key: qbtypes.OrderByKey{TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: "trace.output_tokens"}}}},
})
require.ErrorContains(t, err, `ordering the span list by trace-level key "trace.output_tokens" is not supported`)
err = build(qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Order: []qbtypes.OrderBy{{Key: qbtypes.OrderByKey{TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: "trace.foo"}}}},
})
require.ErrorContains(t, err, `trace-level key "trace.foo"`)
err = build(qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Order: []qbtypes.OrderBy{{Key: qbtypes.OrderByKey{TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: "duration_nano"}}, Direction: qbtypes.OrderDirectionDesc}},
Limit: 10,
})
require.NoError(t, err, "bare duration_nano is a span column, not a trace-level key")
}
// Variables in a trace-level condition resolve through the standard pipeline; a
// dynamic __all__ drops the condition (no scope CTE).
func TestBuild_SpanList_TraceFilter_Variables(t *testing.T) {
b := newTestBuilder(t)
build := func(expr string, vars map[string]qbtypes.VariableItem) (*qbtypes.Statement, error) {
return b.Build(context.Background(), valuer.UUID{}, testStartMs, testEndMs, qbtypes.RequestTypeRaw,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
Filter: &qbtypes.Filter{Expression: expr},
Limit: 10,
}, vars)
}
stmt, err := build("trace.output_tokens > $threshold",
map[string]qbtypes.VariableItem{"threshold": {Value: 700}})
require.NoError(t, err)
assert.Contains(t, renderSQL(t, stmt), "HAVING output_tokens > 700")
stmt, err = build("trace.output_tokens > $threshold",
map[string]qbtypes.VariableItem{"threshold": {Type: qbtypes.DynamicVariableType, Value: "__all__"}})
require.NoError(t, err)
assert.NotContains(t, stmt.Query, "__trace_scope")
}

View File

@@ -1,8 +1,6 @@
package aistatementbuilder
import (
"strings"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/statementbuilder"
@@ -27,11 +25,8 @@ func NewFactory(
// Scope describes gen_ai for the scoped trace builder: an AI trace has >=1 gen_ai
// LLM, tool, or agent span, and its list adds AI/LLM per-trace metrics.
func Scope() scopedtraces.TraceScope {
gateKeyNames := []string{aiobservabilitytypes.GenAIRequestModel, aiobservabilitytypes.GenAIToolName, aiobservabilitytypes.GenAIAgentName}
gateExprs := make([]string, 0, len(gateKeyNames))
gateKeys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(gateKeyNames))
for _, name := range gateKeyNames {
gateExprs = append(gateExprs, name+" EXISTS")
gateKeys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(aiobservabilitytypes.GenAISpanGateKeys))
for _, name := range aiobservabilitytypes.GenAISpanGateKeys {
gateKeys = append(gateKeys, &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalTraces,
@@ -79,7 +74,7 @@ func Scope() scopedtraces.TraceScope {
}
return scopedtraces.TraceScope{
FilterExpression: strings.Join(gateExprs, " OR "),
FilterExpression: aiobservabilitytypes.GenAISpanFilterExpression(),
FieldKeys: gateKeys,
Columns: columns,
DefaultOrderAlias: "last_activity_time",

View File

@@ -114,6 +114,9 @@ func (b *scopedTraceStatementBuilder) Build(
case qbtypes.RequestTypeTrace:
return b.buildTraceListQuery(ctx, orgID, querybuilder.ToNanoSecs(start), querybuilder.ToNanoSecs(end), query, variables)
case qbtypes.RequestTypeRaw:
if err := b.validateRawOrderKeys(query); err != nil {
return nil, err
}
return b.buildDelegated(ctx, orgID, start, end, requestType, query, variables)
case qbtypes.RequestTypeScalar, qbtypes.RequestTypeTimeSeries:
return b.buildAggregation(ctx, orgID, start, end, requestType, query, variables)
@@ -122,27 +125,19 @@ func (b *scopedTraceStatementBuilder) Build(
}
}
// buildDelegated ANDs the base gate into the user filter and delegates to the
// standard trace builder (the span-list / raw path).
func (b *scopedTraceStatementBuilder) buildDelegated(
ctx context.Context,
orgID valuer.UUID,
start, end uint64,
requestType qbtypes.RequestType,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
gate := b.scope.FilterExpression
expr := gate
if query.Filter != nil && strings.TrimSpace(query.Filter.Expression) != "" {
expr = fmt.Sprintf("(%s) AND (%s)", gate, query.Filter.Expression)
// validateRawOrderKeys rejects trace-level order keys — no per-trace value exists on
// span rows. A bare name may be a span column sharing an alias (duration_nano), so it passes.
// TODO: move this into the request validation layer (querybuildertypesv5/validation.go).
func (b *scopedTraceStatementBuilder) validateRawOrderKeys(query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]) error {
for _, o := range query.Order {
key := o.Key.TelemetryFieldKey
key.Normalize()
if key.FieldContext == telemetrytypes.FieldContextTrace {
return errors.NewInvalidInputf(errors.CodeInvalidInput,
"ordering the span list by trace-level key %q is not supported; order by span columns instead (e.g. timestamp, duration_nano)", o.Key.Name)
}
}
// shallow copy; only Filter is replaced, caller's query untouched
gated := query
gated.Filter = &qbtypes.Filter{Expression: expr}
return b.traceStmtBuilder.Build(ctx, orgID, start, end, requestType, gated, variables)
return nil
}
// traceScopedStatementBuilder is the delegate's optional capability of constraining a
@@ -153,10 +148,10 @@ type traceScopedStatementBuilder interface {
BuildTraceScoped(ctx context.Context, orgID valuer.UUID, start, end uint64, requestType qbtypes.RequestType, query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation], variables map[string]qbtypes.VariableItem, traceScope, traceScopeResource *qbtypes.Statement) (*qbtypes.Statement, error)
}
// buildDelegatedAggregation serves span-level scalar/time-series through the standard
// trace builder, with the gate ANDed into the span-level filter part; a trace-level
// part becomes a qualification the delegate constrains trace_id by.
func (b *scopedTraceStatementBuilder) buildDelegatedAggregation(
// buildDelegated serves the raw span list and span-level scalar/time-series through
// the standard trace builder, with the gate ANDed into the span-level filter part; a
// trace-level part becomes a qualification the delegate constrains trace_id by.
func (b *scopedTraceStatementBuilder) buildDelegated(
ctx context.Context,
orgID valuer.UUID,
start, end uint64,

View File

@@ -45,7 +45,7 @@ func (b *scopedTraceStatementBuilder) buildAggregation(
return nil, err
}
if len(traceAggs) == 0 {
return b.buildDelegatedAggregation(ctx, orgID, start, end, requestType, query, variables)
return b.buildDelegated(ctx, orgID, start, end, requestType, query, variables)
}
return b.buildTraceAggregationQuery(ctx, orgID, querybuilder.ToNanoSecs(start), querybuilder.ToNanoSecs(end), requestType, query, variables, traceAggs)
}

View File

@@ -377,6 +377,11 @@ func (b *traceQueryStatementBuilder) buildListQuery(
cteArgs = append(cteArgs, args)
}
if scopeFrags, scopeArgs := b.attachTraceScope(sb, frag != ""); len(scopeFrags) > 0 {
cteFragments = append(cteFragments, scopeFrags...)
cteArgs = append(cteArgs, scopeArgs...)
}
for i, field := range query.SelectFields {
expr, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &field, telemetrytypes.FieldDataTypeUnspecified, keys)
if err != nil {

View File

@@ -0,0 +1,31 @@
package aitelemetryschema
import (
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
)
var (
traceAggregateNames = func() map[string]struct{} {
names := make(map[string]struct{}, len(TraceAggregateFields))
for name := range TraceAggregateFields {
names[name] = struct{}{}
}
return names
}()
genAISpanGate = "(" + aiobservabilitytypes.GenAISpanFilterExpression() + ")"
)
// ScopedExistingQuery narrows value suggestions to gen_ai spans: the caller's
// filter minus its per-trace aggregate atoms (never ingested, so nothing can
// narrow on them), ANDed with the gen_ai span gate. An unparseable filter is
// dropped and reported through the returned error; the gate alone is still
// usable, matching how the metadata store treats a bad filter downstream.
func ScopedExistingQuery(existingQuery string) (string, error) {
spanExpr, _, err := querybuilder.SplitFilterForAggregates(existingQuery, traceAggregateNames)
if err != nil || spanExpr == "" {
return genAISpanGate, err
}
return genAISpanGate + " AND (" + spanExpr + ")", nil
}

View File

@@ -0,0 +1,103 @@
package aitelemetryschema
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestScopedExistingQuery(t *testing.T) {
gate := "(gen_ai.request.model EXISTS OR gen_ai.tool.name EXISTS OR gen_ai.agent.name EXISTS)"
testCases := []struct {
name string
existingQuery string
expected string
expectedErr string
}{
{
name: "empty query returns the gate alone",
existingQuery: "",
expected: gate,
},
{
name: "span filter is preserved under the gate",
existingQuery: "service.name = 'checkout'",
expected: gate + " AND (service.name = 'checkout')",
},
{
name: "trace aggregate filter is stripped",
existingQuery: "llm_call_count > 5",
expected: gate,
},
{
name: "mixed filter keeps only the span part",
existingQuery: "llm_call_count > 5 AND gen_ai.request.model = 'gpt-4'",
expected: gate + " AND (gen_ai.request.model = 'gpt-4')",
},
{
name: "trace context filter is stripped",
existingQuery: "trace.total_tokens > 100 AND service.name = 'checkout'",
expected: gate + " AND (service.name = 'checkout')",
},
{
name: "unparseable filter is dropped",
existingQuery: "service.name = ",
expected: gate,
expectedErr: "syntax errors while parsing the filter expression",
},
{
name: "multiple span conditions survive as one AND chain",
existingQuery: "service.name = 'checkout' AND gen_ai.request.model = 'gpt-4' AND llm_call_count > 5",
expected: gate + " AND (service.name = 'checkout' AND gen_ai.request.model = 'gpt-4')",
},
{
name: "span OR group is kept whole and parenthesized against the AND join",
existingQuery: "service.name = 'a' OR service.name = 'b'",
expected: gate + " AND ((service.name = 'a' OR service.name = 'b'))",
},
{
name: "parenthesized span OR group ANDed with an aggregate keeps only the group",
existingQuery: "(service.name = 'a' OR service.name = 'b') AND llm_call_count > 5",
expected: gate + " AND ((service.name = 'a' OR service.name = 'b'))",
},
{
name: "OR group of trace aggregates is stripped whole",
existingQuery: "llm_call_count > 5 OR total_tokens > 100",
expected: gate,
},
{
name: "OR mixing aggregate and span atoms drops the whole filter",
existingQuery: "llm_call_count > 5 OR service.name = 'checkout'",
expected: gate,
expectedErr: "trace-level and span-level filters cannot be combined within an OR/NOT group",
},
{
name: "parenthesized AND group is split, not routed whole",
existingQuery: "(llm_call_count > 5 AND service.name = 'checkout') AND gen_ai.request.model = 'gpt-4'",
expected: gate + " AND (service.name = 'checkout' AND gen_ai.request.model = 'gpt-4')",
},
{
name: "NOT over an aggregate group is stripped",
existingQuery: "NOT (llm_call_count > 5) AND service.name = 'checkout'",
expected: gate + " AND (service.name = 'checkout')",
},
{
name: "NOT over a span group is kept",
existingQuery: "NOT (service.name = 'checkout')",
expected: gate + " AND (NOT (service.name = 'checkout'))",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
scoped, err := ScopedExistingQuery(testCase.existingQuery)
if testCase.expectedErr != "" {
assert.ErrorContains(t, err, testCase.expectedErr)
} else {
assert.NoError(t, err)
}
assert.Equal(t, testCase.expected, scoped)
})
}
}

View File

@@ -15,11 +15,12 @@ type PostableFieldKeysParams struct {
Limit int `query:"limit"`
}
// existingQuery is unsupported until the computed per-trace aggregates it may
// reference can be narrowed on.
// existingQuery may reference the computed per-trace aggregates, which are
// never ingested; those filters are stripped before narrowing values.
type PostableFieldValueParams struct {
PostableFieldKeysParams
Name string `query:"name"`
Name string `query:"name"`
ExistingQuery string `query:"existingQuery"`
}
func NewFieldKeySelectorFromPostableFieldKeysParams(params PostableFieldKeysParams) *telemetrytypes.FieldKeySelector {
@@ -30,6 +31,7 @@ func NewFieldValueSelectorFromPostableFieldValueParams(params PostableFieldValue
return telemetrytypes.NewFieldValueSelectorFromPostableFieldValueParams(telemetrytypes.PostableFieldValueParams{
PostableFieldKeysParams: params.telemetryParams(),
Name: params.Name,
ExistingQuery: params.ExistingQuery,
})
}

View File

@@ -1,5 +1,7 @@
package aiobservabilitytypes
import "strings"
// OpenTelemetry gen_ai semantic-convention attribute keys. Single source of truth
// shared by the AI query builder and the LLM pricing pipeline.
const (
@@ -26,3 +28,17 @@ const (
SignozGenAICostCacheWrite = "_signoz.gen_ai.cost_cache_write"
SignozGenAITotalCost = "_signoz.gen_ai.total_cost"
)
// GenAISpanGateKeys mark a span as gen_ai: an LLM call, a tool call, or an
// agent span. A trace belongs to the AI explorer when any span carries one.
var GenAISpanGateKeys = []string{GenAIRequestModel, GenAIToolName, GenAIAgentName}
// GenAISpanFilterExpression renders the gate as a query-builder filter
// expression: each gate key ORed on EXISTS.
func GenAISpanFilterExpression() string {
exprs := make([]string, 0, len(GenAISpanGateKeys))
for _, key := range GenAISpanGateKeys {
exprs = append(exprs, key+" EXISTS")
}
return strings.Join(exprs, " OR ")
}

View File

@@ -1,6 +1,7 @@
package alertmanagertypes
import (
"crypto/rand"
"encoding/json"
"reflect"
"regexp"
@@ -16,9 +17,10 @@ import (
)
var (
ErrCodeAlertmanagerChannelNotFound = errors.MustNewCode("alertmanager_channel_not_found")
ErrCodeAlertmanagerChannelNameMismatch = errors.MustNewCode("alertmanager_channel_name_mismatch")
ErrCodeAlertmanagerChannelInvalid = errors.MustNewCode("alertmanager_channel_invalid")
ErrCodeAlertmanagerChannelNotFound = errors.MustNewCode("alertmanager_channel_not_found")
ErrCodeAlertmanagerChannelNameMismatch = errors.MustNewCode("alertmanager_channel_name_mismatch")
ErrCodeAlertmanagerChannelInvalid = errors.MustNewCode("alertmanager_channel_invalid")
ErrCodeAlertmanagerChannelAlreadyExists = errors.MustNewCode("alertmanager_channel_already_exists")
)
var (
@@ -48,14 +50,19 @@ type Channel struct {
types.Identifiable
types.TimeAuditable
Name string `json:"name" required:"true" bun:"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"`
// Name is the DNS1123 identity references will migrate onto. Until then
// DisplayName is the receiver name inside Data and what policies and rules
// 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"`
}
// 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)
@@ -70,8 +77,9 @@ func NewChannelFromReceiver(receiver *Receiver, orgID string) (*Channel, error)
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
},
Name: receiver.Name,
OrgID: orgID,
Name: generateChannelName(receiver.Name),
DisplayName: receiver.Name,
OrgID: orgID,
}
data, err := json.Marshal(receiver)
@@ -88,6 +96,60 @@ func NewChannelFromReceiver(receiver *Receiver, orgID string) (*Channel, error)
return &channel, nil
}
const channelNameSuffixLen = 8
// generateChannelName is a copy of dashboardtypes.generateDashboardName: slugify
// the display name, then append a random suffix rather than looping on collisions.
func generateChannelName(displayName string) string {
const dns1123LabelMaxLen = 63
suffixAlphabet := []byte("abcdefghijklmnopqrstuvwxyz0123456789")
var b strings.Builder
b.Grow(len(displayName))
prevHyphen := false
for _, r := range strings.ToLower(displayName) {
switch {
case (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9'):
b.WriteRune(r)
prevHyphen = false
case b.Len() > 0 && !prevHyphen:
b.WriteByte('-')
prevHyphen = true
}
}
prefix := strings.TrimRight(b.String(), "-")
suffix := make([]byte, channelNameSuffixLen)
if _, err := rand.Read(suffix); err != nil {
panic(errors.WrapInternalf(err, errors.CodeInternal, "read random for channel name suffix"))
}
for i := range suffix {
suffix[i] = suffixAlphabet[int(suffix[i])%len(suffixAlphabet)]
}
maxPrefix := dns1123LabelMaxLen - 1 - channelNameSuffixLen
if len(prefix) > maxPrefix {
prefix = strings.TrimRight(prefix[:maxPrefix], "-")
}
if prefix == "" {
return string(suffix)
}
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.
@@ -151,26 +213,6 @@ func NewConfigFromChannels(globalConfig GlobalConfig, routeConfig RouteConfig, c
return cfg, nil
}
func GetChannelByID(channels Channels, id valuer.UUID) (int, *Channel, error) {
for i, channel := range channels {
if channel.ID == id {
return i, channel, nil
}
}
return 0, nil, errors.Newf(errors.TypeNotFound, ErrCodeAlertmanagerChannelNotFound, "cannot find channel with id %s", id.StringValue())
}
func GetChannelByName(channels Channels, name string) (int, *Channel, error) {
for i, channel := range channels {
if channel.Name == name {
return i, channel, nil
}
}
return 0, nil, errors.Newf(errors.TypeNotFound, ErrCodeAlertmanagerChannelNotFound, "cannot find channel with name %s", name)
}
func NewStatsFromChannels(channels Channels) map[string]any {
stats := make(map[string]any)
for _, channel := range channels {
@@ -188,15 +230,21 @@ func NewStatsFromChannels(channels Channels) map[string]any {
}
func (c *Channel) Update(receiver *Receiver) error {
channel, err := NewChannelFromReceiver(receiver, c.OrgID)
channel, err := NewChannelFromReceiverWithName(receiver, c.Name, c.OrgID)
if err != nil {
return err
}
if c.Name != channel.Name {
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()

View File

@@ -22,9 +22,9 @@ func TestNewConfigFromChannels(t *testing.T) {
name: "OneEmailChannel",
channels: Channels{
{
Name: "email-receiver",
Type: "email",
Data: `{"name":"email-receiver","email_configs":[{"to":"test@example.com"}]}`,
DisplayName: "email-receiver",
Type: "email",
Data: `{"name":"email-receiver","email_configs":[{"to":"test@example.com"}]}`,
},
},
expectedRoutes: []map[string]any{{"receiver": "email-receiver", "continue": true, "matchers": []any{"ruleId=~\"-1\""}}},
@@ -46,9 +46,9 @@ func TestNewConfigFromChannels(t *testing.T) {
name: "OneSlackChannel",
channels: Channels{
{
Name: "slack-receiver",
Type: "slack",
Data: `{"name":"slack-receiver","slack_configs":[{"channel":"#alerts","api_url":"https://slack.com/api/test","send_resolved":true}]}`,
DisplayName: "slack-receiver",
Type: "slack",
Data: `{"name":"slack-receiver","slack_configs":[{"channel":"#alerts","api_url":"https://slack.com/api/test","send_resolved":true}]}`,
},
},
expectedRoutes: []map[string]any{{"receiver": "slack-receiver", "continue": true, "matchers": []any{"ruleId=~\"-1\""}}},
@@ -80,9 +80,9 @@ func TestNewConfigFromChannels(t *testing.T) {
name: "OnePagerdutyChannel",
channels: Channels{
{
Name: "pagerduty-receiver",
Type: "pagerduty",
Data: `{"name":"pagerduty-receiver","pagerduty_configs":[{"service_key":"test"}]}`,
DisplayName: "pagerduty-receiver",
Type: "pagerduty",
Data: `{"name":"pagerduty-receiver","pagerduty_configs":[{"service_key":"test"}]}`,
},
},
expectedRoutes: []map[string]any{{"receiver": "pagerduty-receiver", "continue": true, "matchers": []any{"ruleId=~\"-1\""}}},
@@ -112,14 +112,14 @@ func TestNewConfigFromChannels(t *testing.T) {
name: "OnePagerdutyAndOneSlackChannel",
channels: Channels{
{
Name: "pagerduty-receiver",
Type: "pagerduty",
Data: `{"name":"pagerduty-receiver","pagerduty_configs":[{"service_key":"test"}]}`,
DisplayName: "pagerduty-receiver",
Type: "pagerduty",
Data: `{"name":"pagerduty-receiver","pagerduty_configs":[{"service_key":"test"}]}`,
},
{
Name: "slack-receiver",
Type: "slack",
Data: `{"name":"slack-receiver","slack_configs":[{"channel":"#alerts","api_url":"https://slack.com/api/test","send_resolved":true}]}`,
DisplayName: "slack-receiver",
Type: "slack",
Data: `{"name":"slack-receiver","slack_configs":[{"channel":"#alerts","api_url":"https://slack.com/api/test","send_resolved":true}]}`,
},
},
expectedRoutes: []map[string]any{{"receiver": "pagerduty-receiver", "continue": true, "matchers": []any{"ruleId=~\"-1\""}}, {"receiver": "slack-receiver", "continue": true, "matchers": []any{"ruleId=~\"-1\""}}},
@@ -243,9 +243,9 @@ func TestNewChannelFromReceiver(t *testing.T) {
},
},
expected: &Channel{
Name: "test-receiver",
Type: "slack",
Data: `{"name":"test-receiver","slack_configs":[{"send_resolved":true,"api_url":"https://slack.com/api/test","channel":"#alerts","timeout":0}]}`,
DisplayName: "test-receiver",
Type: "slack",
Data: `{"name":"test-receiver","slack_configs":[{"send_resolved":true,"api_url":"https://slack.com/api/test","channel":"#alerts","timeout":0}]}`,
},
pass: true,
},
@@ -261,7 +261,7 @@ func TestNewChannelFromReceiver(t *testing.T) {
}
assert.NoError(t, err)
assert.Equal(t, testCase.expected.Name, channel.Name)
assert.Equal(t, testCase.expected.DisplayName, channel.DisplayName)
assert.Equal(t, testCase.expected.Type, channel.Type)
assert.Equal(t, testCase.expected.Data, channel.Data)
})
@@ -289,7 +289,7 @@ func TestNewChannelFromReceiverGoogleChat(t *testing.T) {
channel, err := NewChannelFromReceiver(receiver, "1")
assert.NoError(t, err)
assert.Equal(t, "googlechat-receiver", channel.Name)
assert.Equal(t, "googlechat-receiver", channel.DisplayName)
assert.Equal(t, "googlechat", channel.Type)
assert.JSONEq(t,
`{"name":"googlechat-receiver","googlechat_configs":[{"send_resolved":false,"webhook_url":"https://chat.googleapis.com/v1/spaces/test/messages","title":"Alert","text":"Body"}]}`,

View File

@@ -72,10 +72,11 @@ type customReceiverConfigs struct {
GoogleChat []*GoogleChatReceiverConfig
Jira []*JiraReceiverConfig
JSMOps []*JSMOpsReceiverConfig
IncidentIO []*IncidentIOReceiverConfig
}
func (c customReceiverConfigs) isEmpty() bool {
return len(c.GoogleChat) == 0 && len(c.Jira) == 0 && len(c.JSMOps) == 0
return len(c.GoogleChat) == 0 && len(c.Jira) == 0 && len(c.JSMOps) == 0 && len(c.IncidentIO) == 0
}
func customConfigsOf(receiver *Receiver) customReceiverConfigs {
@@ -83,6 +84,7 @@ func customConfigsOf(receiver *Receiver) customReceiverConfigs {
GoogleChat: receiver.GoogleChatConfigs,
Jira: receiver.JiraConfigs,
JSMOps: receiver.JSMOpsConfigs,
IncidentIO: receiver.IncidentIOConfigs,
}
}
@@ -193,6 +195,7 @@ func extendedReceivers(c *config.Config, customConfigs map[string]customReceiver
GoogleChatConfigs: custom.GoogleChat,
JiraConfigs: custom.Jira,
JSMOpsConfigs: custom.JSMOps,
IncidentIOConfigs: custom.IncidentIO,
}
}
@@ -370,6 +373,7 @@ func (c *Config) GetReceiver(name string) (*Receiver, error) {
GoogleChatConfigs: custom.GoogleChat,
JiraConfigs: custom.Jira,
JSMOpsConfigs: custom.JSMOps,
IncidentIOConfigs: custom.IncidentIO,
}, nil
}
}
@@ -459,6 +463,11 @@ func (c *Config) applyNativeDefaults() {
jc.HTTPConfig = httpDefault
}
}
for _, ic := range custom.IncidentIO {
if ic.HTTPConfig == nil {
ic.HTTPConfig = httpDefault
}
}
}
}

View File

@@ -0,0 +1,109 @@
package alertmanagertypes
import (
"fmt"
"net/url"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
// incidentIOEventsPathPrefix is the path of incident.io's HTTP alert source
// endpoint (Alert Events V2 API). The full URL is per-source:
// https://api.incident.io/v2/alert_events/http/<source_config_id>.
const incidentIOEventsPathPrefix = "/v2/alert_events/http/"
// The description is markdown; incident.io renders it natively. The templates
// mirror Google Chat / Jira / JSM for a consistent default across channels.
const (
DefaultIncidentIOTitleTemplate = `[{{ .Status | toUpper }}{{ if eq .Status "firing" }}:{{ .Alerts.Firing | len }}{{ end }}] {{ .CommonLabels.alertname }}`
DefaultIncidentIODescriptionTemplate = `{{ range .Alerts -}}
**Alert:** {{ .Labels.alertname }}{{ if .Labels.severity }} ({{ .Labels.severity }}){{ end }}
{{ if .Annotations.summary }}**Summary:** {{ .Annotations.summary }}
{{ end }}{{ if .Annotations.description }}**Description:** {{ .Annotations.description }}
{{ end }}{{ if .GeneratorURL }}[View in SigNoz]({{ .GeneratorURL }})
{{ end }}{{ if .Annotations.related_logs }}[View related logs]({{ .Annotations.related_logs }})
{{ end }}{{ if .Annotations.related_traces }}[View related traces]({{ .Annotations.related_traces }})
{{ end }}{{ end }}`
)
// IncidentIOReceiverConfig is the SigNoz incident.io receiver, backed by an
// incident.io HTTP alert source. URL is the per-source alert events endpoint
// and Token its secret, both copied from the source's setup page.
type IncidentIOReceiverConfig struct {
config.NotifierConfig `yaml:",inline" json:",inline"`
HTTPConfig *commoncfg.HTTPClientConfig `yaml:"http_config,omitempty" json:"http_config,omitempty"`
URL string `yaml:"url,omitempty" json:"url,omitempty"`
Token config.Secret `yaml:"token,omitempty" json:"token,omitempty"`
Title string `yaml:"title,omitempty" json:"title,omitempty"`
Description string `yaml:"description,omitempty" json:"description,omitempty"`
// Metadata is merged into the event's metadata on top of the group's common
// labels (channel wins on key clash). Values are template-expanded.
Metadata map[string]string `yaml:"metadata,omitempty" json:"metadata,omitempty"`
}
// send_resolved has no omitempty upstream, so a var default here is overwritten
// by the yaml round-trip to the request value (false when omitted); the UI sends
// it explicitly, defaulted on, so incident.io alerts resolve with the rule.
var DefaultIncidentIOReceiverConfig = IncidentIOReceiverConfig{
NotifierConfig: config.NotifierConfig{
VSendResolved: false,
},
Title: DefaultIncidentIOTitleTemplate,
Description: DefaultIncidentIODescriptionTemplate,
}
func (c *IncidentIOReceiverConfig) UnmarshalYAML(unmarshal func(any) error) error {
*c = DefaultIncidentIOReceiverConfig
type plain IncidentIOReceiverConfig
if err := unmarshal((*plain)(c)); err != nil {
return err
}
if c.Title == "" {
c.Title = DefaultIncidentIOTitleTemplate
}
if c.Description == "" {
c.Description = DefaultIncidentIODescriptionTemplate
}
// Values are stored and sent exactly as configured, so anything that is
// not already canonical is rejected rather than rewritten.
if c.URL != strings.TrimSpace(c.URL) {
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "incidentio url must not have leading or trailing whitespace")
}
u, err := url.Parse(c.URL)
if c.URL == "" || err != nil || u.Scheme != "https" || u.Host == "" ||
!strings.Contains(u.Path, incidentIOEventsPathPrefix) ||
strings.HasSuffix(u.Path, incidentIOEventsPathPrefix) {
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, fmt.Sprintf("incidentio url must be an alert events URL (https://api.incident.io%s<source_config_id>)", incidentIOEventsPathPrefix))
}
if strings.HasSuffix(c.URL, "/") {
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "incidentio url must not end with a trailing slash")
}
token := string(c.Token)
if token == "" {
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "incidentio token is required")
}
if token != strings.TrimSpace(token) {
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "incidentio token must not have leading or trailing whitespace")
}
// incident.io's setup page shows the header value as "Bearer <token>"; a
// pasted prefix would be sent doubled, so reject it instead.
if strings.EqualFold(token, "bearer") || (len(token) >= 7 && strings.EqualFold(token[:7], "bearer ")) {
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "incidentio token must be the source's secret token only, without the Bearer prefix")
}
return nil
}

View File

@@ -0,0 +1,65 @@
package alertmanagertypes
import (
"fmt"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
const testIncidentIOURL = "https://api.incident.io/v2/alert_events/http/01M0D1JNVBGBGVTWX053EM12XV"
func TestIncidentIOReceiverConfigDefaults(t *testing.T) {
r, err := NewReceiver(fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s","token":"tok-123"}]}`, testIncidentIOURL))
require.NoError(t, err)
require.Len(t, r.IncidentIOConfigs, 1)
c := r.IncidentIOConfigs[0]
assert.Equal(t, testIncidentIOURL, c.URL)
assert.Equal(t, "tok-123", string(c.Token))
assert.Equal(t, DefaultIncidentIOTitleTemplate, c.Title)
assert.Equal(t, DefaultIncidentIODescriptionTemplate, c.Description)
assert.False(t, c.SendResolved()) // default off when omitted, like other channels
ch, err := NewChannelFromReceiver(r, "org-1")
require.NoError(t, err)
assert.Equal(t, "incidentio", ch.Type)
}
func TestIncidentIOReceiverConfigOverrides(t *testing.T) {
r, err := NewReceiver(fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s","token":"k","title":"t","description":"d","send_resolved":true}]}`, testIncidentIOURL))
require.NoError(t, err)
require.Len(t, r.IncidentIOConfigs, 1)
c := r.IncidentIOConfigs[0]
assert.Equal(t, testIncidentIOURL, c.URL)
assert.Equal(t, "t", c.Title)
assert.Equal(t, "d", c.Description)
assert.True(t, c.SendResolved())
}
func TestIncidentIOReceiverConfigValidation(t *testing.T) {
cases := []struct {
name string
json string
}{
{"missing url", `{"name":"incio","incidentio_configs":[{"token":"k"}]}`},
{"http url", `{"name":"incio","incidentio_configs":[{"url":"http://api.incident.io/v2/alert_events/http/abc","token":"k"}]}`},
{"not an alert events url", `{"name":"incio","incidentio_configs":[{"url":"https://api.incident.io/v2/incidents","token":"k"}]}`},
{"missing source config id", `{"name":"incio","incidentio_configs":[{"url":"https://api.incident.io/v2/alert_events/http/","token":"k"}]}`},
{"trailing slash", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s/","token":"k"}]}`, testIncidentIOURL)},
{"whitespace around url", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":" %s ","token":"k"}]}`, testIncidentIOURL)},
{"missing token", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s"}]}`, testIncidentIOURL)},
{"bearer prefixed token", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s","token":"Bearer tok-123"}]}`, testIncidentIOURL)},
{"lowercase bearer prefixed token", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s","token":"bearer tok-123"}]}`, testIncidentIOURL)},
{"bearer only token", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s","token":"Bearer"}]}`, testIncidentIOURL)},
{"whitespace around token", fmt.Sprintf(`{"name":"incio","incidentio_configs":[{"url":"%s","token":" tok-123 "}]}`, testIncidentIOURL)},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
_, err := NewReceiver(c.json)
assert.Error(t, err)
})
}
}

View File

@@ -28,6 +28,9 @@ type Receiver struct {
JiraConfigs []*JiraReceiverConfig `json:"jira_configs,omitempty" yaml:"jira_configs,omitempty"`
// JSM Ops (ex-Opsgenie alert API); delivered by reusing the Opsgenie notifier.
JSMOpsConfigs []*JSMOpsReceiverConfig `json:"jsmops_configs,omitempty" yaml:"jsmops_configs,omitempty"`
// Shadows upstream's incidentio_configs so our custom notifier (templater,
// group-key dedup, label metadata) handles it instead of upstream's.
IncidentIOConfigs []*IncidentIOReceiverConfig `json:"incidentio_configs,omitempty" yaml:"incidentio_configs,omitempty"`
}
// NewReceiver builds a Receiver from its JSON input, applying each notifier
@@ -72,6 +75,14 @@ func NewReceiver(input string) (*Receiver, error) {
receiver.JSMOpsConfigs[i] = defaulted
}
for i, ic := range receiver.IncidentIOConfigs {
defaulted, err := defaultedNotifierConfig(ic)
if err != nil {
return nil, err
}
receiver.IncidentIOConfigs[i] = defaulted
}
return receiver, nil
}

View File

@@ -62,7 +62,7 @@ var (
ResourceMetaResourcePipeline = NewResourceMetaResource(KindPipeline)
ResourceMetaResourceUserPreference = NewResourceMetaResource(KindUserPreference)
ResourceMetaResourceOrgPreference = NewResourceMetaResource(KindOrgPreference)
ResourceMetaResourceQuickFilter = NewResourceMetaResource(KindQuickFilter)
ResourceMetaResourceQuickFilter = NewResourceMetaResource(KindQuickFilter, VerbList, VerbRead, VerbUpdate)
ResourceMetaResourceTTLSetting = NewResourceMetaResource(KindTTLSetting)
ResourceMetaResourceRule = NewResourceMetaResource(KindRule)
ResourceMetaResourcePlannedMaintenance = NewResourceMetaResource(KindPlannedMaintenance)

View File

@@ -81,7 +81,7 @@ type GettableCreatedIngestionKey struct {
Value string `json:"value" required:"true"`
}
type PostableIngestionKeyLimit struct {
type DeprecatedPostableIngestionKeyLimit struct {
Signal string `json:"signal"`
Config LimitConfig `json:"config"`
Tags []string `json:"tags"`
@@ -91,6 +91,13 @@ type GettableCreatedIngestionKeyLimit struct {
ID string `json:"id" required:"true"`
}
type PostableIngestionKeyLimit struct {
KeyID string `json:"keyId" required:"true"`
Signal string `json:"signal"`
Config LimitConfig `json:"config"`
Tags []string `json:"tags"`
}
type UpdatableIngestionKeyLimit struct {
Config LimitConfig `json:"config" required:"true"`
Tags []string `json:"tags"`

View File

@@ -5,58 +5,69 @@ import (
"time"
"github.com/SigNoz/signoz/pkg/errors"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/uptrace/bun"
)
type Signal struct {
type Source struct {
valuer.String
}
func (enum *Signal) UnmarshalJSON(data []byte) error {
func (enum *Source) UnmarshalJSON(data []byte) error {
var str string
if err := json.Unmarshal(data, &str); err != nil {
return err
}
signal, err := NewSignal(str)
source, err := NewSource(str)
if err != nil {
return err
}
*enum = signal
*enum = source
return nil
}
var (
SignalTraces = Signal{valuer.NewString("traces")}
SignalLogs = Signal{valuer.NewString("logs")}
SignalApiMonitoring = Signal{valuer.NewString("api_monitoring")}
SignalExceptions = Signal{valuer.NewString("exceptions")}
SignalMeter = Signal{valuer.NewString("meter")}
SignalAiObservability = Signal{valuer.NewString("ai_observability")}
SourceTraces = Source{valuer.NewString("traces")}
SourceLogs = Source{valuer.NewString("logs")}
SourceApiMonitoring = Source{valuer.NewString("api_monitoring")}
SourceExceptions = Source{valuer.NewString("exceptions")}
SourceMeter = Source{valuer.NewString("meter")}
SourceAiObservability = Source{valuer.NewString("ai_observability")}
)
// NewSignal creates a Signal from a string.
func NewSignal(s string) (Signal, error) {
func (Source) Enum() []any {
return []any{
SourceTraces,
SourceLogs,
SourceApiMonitoring,
SourceExceptions,
SourceMeter,
SourceAiObservability,
}
}
// NewSource creates a Source from a string.
func NewSource(s string) (Source, error) {
switch s {
case "traces":
return SignalTraces, nil
return SourceTraces, nil
case "logs":
return SignalLogs, nil
return SourceLogs, nil
case "api_monitoring":
return SignalApiMonitoring, nil
return SourceApiMonitoring, nil
case "exceptions":
return SignalExceptions, nil
return SourceExceptions, nil
case "meter":
return SignalMeter, nil
return SourceMeter, nil
case "ai_observability":
return SignalAiObservability, nil
return SourceAiObservability, nil
default:
return Signal{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid signal: %s", s)
return Source{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid source: %s", s)
}
}
@@ -65,33 +76,46 @@ type StorableQuickFilter struct {
types.Identifiable
OrgID valuer.UUID `bun:"org_id,type:text,notnull"`
Filter string `bun:"filter,type:text,notnull"`
Signal Signal `bun:"signal,type:text,notnull"`
Source Source `bun:"source,type:text,notnull"`
types.TimeAuditable
}
type SignalFilters struct {
Signal Signal `json:"signal"`
Filters []v3.AttributeKey `json:"filters"`
type SourceFilters struct {
types.Identifiable
types.TimeAuditable
OrgID valuer.UUID `json:"orgId" required:"true"`
Source Source `json:"source" required:"true"`
Filters []telemetrytypes.TelemetryFieldKey `json:"filters" required:"true" nullable:"false"`
}
type UpdatableQuickFilters struct {
Signal Signal `json:"signal"`
Filters []v3.AttributeKey `json:"filters"`
Filters []telemetrytypes.TelemetryFieldKey `json:"filters" required:"true" nullable:"false"`
}
// NewStorableQuickFilter creates a new StorableQuickFilter after validation.
func NewStorableQuickFilter(orgID valuer.UUID, signal Signal, filterJSON []byte) (*StorableQuickFilter, error) {
if orgID.StringValue() == "" {
func NewStorableQuickFilter(orgID valuer.UUID, source Source, filters []telemetrytypes.TelemetryFieldKey) (*StorableQuickFilter, error) {
if orgID.IsZero() {
return nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "orgID is required")
}
if _, err := NewSignal(signal.StringValue()); err != nil {
if _, err := NewSource(source.StringValue()); err != nil {
return nil, err
}
var filters []v3.AttributeKey
if err := json.Unmarshal(filterJSON, &filters); err != nil {
return nil, errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter JSON")
if err := validateFilters(filters); err != nil {
return nil, err
}
// A nil slice marshals to the JSON literal "null"; store an empty array so
// reads never have to render a null filter list.
if filters == nil {
filters = []telemetrytypes.TelemetryFieldKey{}
}
filterJSON, err := json.Marshal(filters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error marshalling filters")
}
now := time.Now()
@@ -100,7 +124,7 @@ func NewStorableQuickFilter(orgID valuer.UUID, signal Signal, filterJSON []byte)
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Signal: signal,
Source: source,
Filter: string(filterJSON),
TimeAuditable: types.TimeAuditable{
CreatedAt: now,
@@ -109,25 +133,21 @@ func NewStorableQuickFilter(orgID valuer.UUID, signal Signal, filterJSON []byte)
}, nil
}
// Update updates an existing StorableQuickFilter with new filter data after validation.
func (quickfilter *StorableQuickFilter) Update(filterJSON []byte) error {
var filters []v3.AttributeKey
if err := json.Unmarshal(filterJSON, &filters); err != nil {
return errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter JSON")
// NewSourceFiltersFromSource creates a SourceFilters with no filters for a source.
func NewSourceFiltersFromSource(source Source) *SourceFilters {
return &SourceFilters{
Source: source,
Filters: []telemetrytypes.TelemetryFieldKey{},
}
quickfilter.Filter = string(filterJSON)
quickfilter.UpdatedAt = time.Now()
return nil
}
// NewSignalFilterFromStorableQuickFilter converts a StorableQuickFilter to a SignalFilters object.
func NewSignalFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFilter) (*SignalFilters, error) {
// NewSourceFilterFromStorableQuickFilter converts a StorableQuickFilter to a SourceFilters object.
func NewSourceFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFilter) (*SourceFilters, error) {
if storableQuickFilter == nil {
return nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "storableQuickFilter cannot be nil")
}
var filters []v3.AttributeKey
filters := []telemetrytypes.TelemetryFieldKey{}
if storableQuickFilter.Filter != "" {
err := json.Unmarshal([]byte(storableQuickFilter.Filter), &filters)
if err != nil {
@@ -135,178 +155,114 @@ func NewSignalFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFi
}
}
return &SignalFilters{
Signal: storableQuickFilter.Signal,
Filters: filters,
// Stored filter JSON can be the literal "null" (a nil slice was upserted),
// which unmarshals to nil; the API contract requires a non-null array.
if filters == nil {
filters = []telemetrytypes.TelemetryFieldKey{}
}
return &SourceFilters{
Identifiable: storableQuickFilter.Identifiable,
OrgID: storableQuickFilter.OrgID,
Source: storableQuickFilter.Source,
Filters: filters,
TimeAuditable: storableQuickFilter.TimeAuditable,
}, nil
}
// NewDefaultQuickFilter generates default filters for all supported signals.
// NewDefaultQuickFilter generates default filters for all supported sources.
func NewDefaultQuickFilter(orgID valuer.UUID) ([]*StorableQuickFilter, error) {
tracesFilters := []map[string]interface{}{
{"key": "duration_nano", "dataType": "float64", "type": "tag"},
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "hasError", "dataType": "bool", "type": "tag"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "name", "dataType": "string", "type": "tag"},
{"key": "rpc.method", "dataType": "string", "type": "tag"},
{"key": "response_status_code", "dataType": "string", "type": "tag"},
{"key": "http_host", "dataType": "string", "type": "tag"},
{"key": "http.method", "dataType": "string", "type": "tag"},
{"key": "http.route", "dataType": "string", "type": "tag"},
{"key": "http_url", "dataType": "string", "type": "tag"},
{"key": "trace_id", "dataType": "string", "type": "tag"},
tracesFilters := []telemetrytypes.TelemetryFieldKey{
{Name: "duration_nano", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeNumber},
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "hasError", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeBool},
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "name", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "rpc.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "response_status_code", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "http_host", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "http.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "http.route", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "http_url", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "trace_id", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
}
logsFilters := []map[string]interface{}{
{"key": "severity_text", "dataType": "string", "type": "resource"},
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "host.name", "dataType": "string", "type": "resource"},
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
logsFilters := []telemetrytypes.TelemetryFieldKey{
{Name: "severity_text", FieldContext: telemetrytypes.FieldContextLog, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "host.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.cluster.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.deployment.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.namespace.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.pod.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
}
apiMonitoringFilters := []map[string]interface{}{
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "rpc.method", "dataType": "string", "type": "tag"},
apiMonitoringFilters := []telemetrytypes.TelemetryFieldKey{
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "rpc.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
}
exceptionsFilters := []map[string]interface{}{
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "host.name", "dataType": "string", "type": "resource"},
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
exceptionsFilters := []telemetrytypes.TelemetryFieldKey{
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "host.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.cluster.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.deployment.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.namespace.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "k8s.pod.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
}
meterFilters := []map[string]interface{}{
{"key": "deployment.environment", "dataType": "float64", "type": "Sum"},
{"key": "service.name", "dataType": "float64", "type": "Sum"},
{"key": "host.name", "dataType": "float64", "type": "Sum"},
// Meter keys are label names with no context or datatype: the meter fields
// API returns them as name+signal only, so the defaults mirror that shape.
meterFilters := []telemetrytypes.TelemetryFieldKey{
{Name: "deployment.environment", Signal: telemetrytypes.SignalMetrics},
{Name: "service.name", Signal: telemetrytypes.SignalMetrics},
{Name: "host.name", Signal: telemetrytypes.SignalMetrics},
}
// AI observability (builder_ai_query trace explorer), ordered by expected
// usage: env scoping, the LLM identity keys, then service and the rest.
aiObservabilityFilters := []map[string]interface{}{
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
{"key": aiobservabilitytypes.GenAIOperationName, "dataType": "string", "type": "tag"},
{"key": aiobservabilitytypes.GenAIProviderName, "dataType": "string", "type": "tag"},
{"key": aiobservabilitytypes.GenAIRequestModel, "dataType": "string", "type": "tag"},
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": aiobservabilitytypes.GenAIToolName, "dataType": "string", "type": "tag"},
{"key": aiobservabilitytypes.GenAIAgentName, "dataType": "string", "type": "tag"},
aiObservabilityFilters := []telemetrytypes.TelemetryFieldKey{
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: aiobservabilitytypes.GenAIOperationName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: aiobservabilitytypes.GenAIProviderName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: aiobservabilitytypes.GenAIRequestModel, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: aiobservabilitytypes.GenAIToolName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
{Name: aiobservabilitytypes.GenAIAgentName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
}
tracesJSON, err := json.Marshal(tracesFilters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal traces filters")
defaults := []struct {
source Source
filters []telemetrytypes.TelemetryFieldKey
}{
{SourceTraces, tracesFilters},
{SourceLogs, logsFilters},
{SourceApiMonitoring, apiMonitoringFilters},
{SourceExceptions, exceptionsFilters},
{SourceMeter, meterFilters},
{SourceAiObservability, aiObservabilityFilters},
}
logsJSON, err := json.Marshal(logsFilters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal logs filters")
storableQuickFilters := make([]*StorableQuickFilter, 0, len(defaults))
for _, def := range defaults {
storableQuickFilter, err := NewStorableQuickFilter(orgID, def.source, def.filters)
if err != nil {
return nil, err
}
storableQuickFilters = append(storableQuickFilters, storableQuickFilter)
}
apiMonitoringJSON, err := json.Marshal(apiMonitoringFilters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal api monitoring filters")
}
exceptionsJSON, err := json.Marshal(exceptionsFilters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal exceptions filters")
}
meterJSON, err := json.Marshal(meterFilters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal meter filters")
}
aiObservabilityJSON, err := json.Marshal(aiObservabilityFilters)
if err != nil {
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal ai observability filters")
}
timeRightNow := time.Now()
return []*StorableQuickFilter{
{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Filter: string(tracesJSON),
Signal: SignalTraces,
TimeAuditable: types.TimeAuditable{
CreatedAt: timeRightNow,
UpdatedAt: timeRightNow,
},
},
{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Filter: string(logsJSON),
Signal: SignalLogs,
TimeAuditable: types.TimeAuditable{
CreatedAt: timeRightNow,
UpdatedAt: timeRightNow,
},
},
{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Filter: string(apiMonitoringJSON),
Signal: SignalApiMonitoring,
TimeAuditable: types.TimeAuditable{
CreatedAt: timeRightNow,
UpdatedAt: timeRightNow,
},
},
{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Filter: string(exceptionsJSON),
Signal: SignalExceptions,
TimeAuditable: types.TimeAuditable{
CreatedAt: timeRightNow,
UpdatedAt: timeRightNow,
},
},
{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Filter: string(meterJSON),
Signal: SignalMeter,
TimeAuditable: types.TimeAuditable{
CreatedAt: timeRightNow,
UpdatedAt: timeRightNow,
},
},
{
Identifiable: types.Identifiable{
ID: valuer.GenerateUUID(),
},
OrgID: orgID,
Filter: string(aiObservabilityJSON),
Signal: SignalAiObservability,
TimeAuditable: types.TimeAuditable{
CreatedAt: timeRightNow,
UpdatedAt: timeRightNow,
},
},
}, nil
return storableQuickFilters, nil
}
func validateFilters(filters []telemetrytypes.TelemetryFieldKey) error {
for _, filter := range filters {
if filter.Name == "" {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "filter name is required")
}
}
return nil
}

View File

@@ -10,10 +10,10 @@ type QuickFilterStore interface {
// Get retrieves all filters for an organization
Get(ctx context.Context, orgID valuer.UUID) ([]*StorableQuickFilter, error)
// GetBySignal retrieves filters for a specific signal in an organization
GetBySignal(ctx context.Context, orgID valuer.UUID, signal string) (*StorableQuickFilter, error)
// GetBySource retrieves filters for a specific source in an organization
GetBySource(ctx context.Context, orgID valuer.UUID, source string) (*StorableQuickFilter, error)
// Upsert inserts or updates filters for an organization and signal
// Upsert inserts or updates filters for an organization and source
Upsert(ctx context.Context, filter *StorableQuickFilter) error
Create(ctx context.Context, filter []*StorableQuickFilter) error
}

View File

@@ -228,7 +228,7 @@ def test_email_channel_never_stores_or_serves_smtp_settings(
with signoz.sqlstore.conn.connect() as conn:
stored = conn.execute(
text("SELECT data FROM notification_channel WHERE name = :name"),
text("SELECT data FROM notification_channel WHERE display_name = :name"),
{"name": hostile_name},
).fetchone()
assert stored is not None

View File

@@ -177,6 +177,61 @@ def test_get_ingestion_keys(
assert data["_pagination"]["total"] == 1
def test_get_ingestion_key_by_id(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list], None],
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
gateway_url = f"/v1/workspaces/me/keys/{TEST_KEY_ID}"
make_http_mocks(
signoz.gateway,
[
Mapping(
request=MappingRequest(
method=HttpMethods.GET,
url=gateway_url,
headers=common_gateway_headers(),
),
response=MappingResponse(
status=200,
json_body={
"status": "success",
"data": {
"id": TEST_KEY_ID,
"name": "my-test-key",
"value": "secret",
"expires_at": "2030-01-01T00:00:00Z",
"tags": ["env:test"],
"created_at": "2024-01-01T00:00:00Z",
"updated_at": "2024-01-01T00:00:00Z",
"workspace_id": "ws-1",
},
},
),
persistent=False,
),
],
)
response = requests.get(
signoz.self.host_configs["8080"].get(f"/api/v2/gateway/ingestion_keys/{TEST_KEY_ID}"),
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.OK, f"Expected 200, got {response.status_code}: {response.text}"
data = response.json()["data"]
assert data["id"] == TEST_KEY_ID
assert data["name"] == "my-test-key"
assert data["workspace_id"] == "ws-1"
assert data["tags"] == ["env:test"]
def test_get_ingestion_keys_custom_pagination(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument

View File

@@ -65,7 +65,7 @@ def test_create_ingestion_key_limit_only_size(
status=201,
json_body={
"status": "success",
"data": {"id": "limit-created-1"},
"data": {"id": "2c5b7d3e-4f6a-4b8c-a09d-3f4a5b6c7d8e"},
},
),
persistent=False,
@@ -86,7 +86,7 @@ def test_create_ingestion_key_limit_only_size(
assert response.status_code == HTTPStatus.CREATED, f"Expected 201, got {response.status_code}: {response.text}"
assert response.json()["data"]["id"] == "limit-created-1"
assert response.json()["data"]["id"] == "2c5b7d3e-4f6a-4b8c-a09d-3f4a5b6c7d8e"
body = get_latest_gateway_request_body(signoz, "POST", gateway_url)
assert body is not None, "Expected a POST request to reach the gateway"
@@ -121,7 +121,7 @@ def test_create_ingestion_key_limit_only_count(
status=201,
json_body={
"status": "success",
"data": {"id": "limit-created-2"},
"data": {"id": "3d6c8e4f-5a7b-4c9d-b1ae-4a5b6c7d8e9f"},
},
),
persistent=False,
@@ -174,7 +174,7 @@ def test_create_ingestion_key_limit_both_size_and_count(
status=201,
json_body={
"status": "success",
"data": {"id": "limit-created-3"},
"data": {"id": "4e7d9f5a-6b8c-4dae-c2bf-5b6c7d8e9fa0"},
},
),
persistent=False,
@@ -391,3 +391,265 @@ def test_delete_ingestion_key_limit(
# Verify at least one DELETE reached the gateway
matched = get_gateway_requests(signoz, "DELETE", gateway_url)
assert len(matched) >= 1, "Expected a DELETE request to reach the gateway"
def test_create_ingestion_limit(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list], None],
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
gateway_url = f"/v1/workspaces/me/keys/{TEST_KEY_ID}/limits"
make_http_mocks(
signoz.gateway,
[
Mapping(
request=MappingRequest(
method=HttpMethods.POST,
url=gateway_url,
headers=common_gateway_headers(),
),
response=MappingResponse(
status=201,
json_body={
"status": "success",
"data": {"id": "5f8ea06b-7c9d-4ebf-a3c0-6c7d8e9fa0b1"},
},
),
persistent=False,
),
],
)
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/gateway/ingestion_limits"),
json={
"keyId": TEST_KEY_ID,
"signal": "logs",
"config": {"day": {"size": 3000}},
"tags": ["test"],
},
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.CREATED, f"Expected 201, got {response.status_code}: {response.text}"
assert response.json()["data"]["id"] == "5f8ea06b-7c9d-4ebf-a3c0-6c7d8e9fa0b1"
body = get_latest_gateway_request_body(signoz, "POST", gateway_url)
assert body is not None, "Expected a POST request to reach the gateway"
assert body["signal"] == "logs"
assert body["config"]["day"]["size"] == 3000
assert body["tags"] == ["test"]
def test_create_ingestion_limit_without_key_id(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/gateway/ingestion_limits"),
json={
"signal": "logs",
"config": {"day": {"size": 3000}},
},
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, f"Expected 400, got {response.status_code}: {response.text}"
def test_get_ingestion_limit(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list], None],
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
gateway_url = f"/v1/workspaces/me/limits/{TEST_LIMIT_ID}"
make_http_mocks(
signoz.gateway,
[
Mapping(
request=MappingRequest(
method=HttpMethods.GET,
url=gateway_url,
headers=common_gateway_headers(),
),
response=MappingResponse(
status=200,
json_body={
"status": "success",
"data": {
"id": TEST_LIMIT_ID,
"key_id": TEST_KEY_ID,
"signal": "logs",
"config": {"day": {"size": 1000}},
"tags": ["test"],
},
},
),
persistent=False,
),
],
)
response = requests.get(
signoz.self.host_configs["8080"].get(f"/api/v2/gateway/ingestion_limits/{TEST_LIMIT_ID}"),
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.OK, f"Expected 200, got {response.status_code}: {response.text}"
data = response.json()["data"]
assert data["id"] == TEST_LIMIT_ID
assert data["key_id"] == TEST_KEY_ID
assert data["signal"] == "logs"
assert data["config"]["day"]["size"] == 1000
def test_get_ingestion_key_limits(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list], None],
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
gateway_url = f"/v1/workspaces/me/keys/{TEST_KEY_ID}/limits"
make_http_mocks(
signoz.gateway,
[
Mapping(
request=MappingRequest(
method=HttpMethods.GET,
url=gateway_url,
headers=common_gateway_headers(),
),
response=MappingResponse(
status=200,
json_body={
"status": "success",
"data": [
{
"id": TEST_LIMIT_ID,
"key_id": TEST_KEY_ID,
"signal": "logs",
"config": {"day": {"size": 1000}},
"tags": ["test"],
}
],
},
),
persistent=False,
),
],
)
response = requests.get(
signoz.self.host_configs["8080"].get(f"/api/v2/gateway/ingestion_keys/{TEST_KEY_ID}/limits"),
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.OK, f"Expected 200, got {response.status_code}: {response.text}"
data = response.json()["data"]
assert len(data) == 1
assert data[0]["id"] == TEST_LIMIT_ID
assert data[0]["key_id"] == TEST_KEY_ID
assert data[0]["signal"] == "logs"
assert data[0]["config"]["day"]["size"] == 1000
def test_update_ingestion_limit(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list], None],
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
gateway_url = f"/v1/workspaces/me/limits/{TEST_LIMIT_ID}"
make_http_mocks(
signoz.gateway,
[
Mapping(
request=MappingRequest(
method=HttpMethods.PATCH,
url=gateway_url,
headers=common_gateway_headers(),
),
response=MappingResponse(status=204),
persistent=False,
),
],
)
response = requests.patch(
signoz.self.host_configs["8080"].get(f"/api/v2/gateway/ingestion_limits/{TEST_LIMIT_ID}"),
json={
"config": {"day": {"size": 4000, "count": 250}},
"tags": ["test"],
},
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, f"Expected 204, got {response.status_code}: {response.text}"
body = get_latest_gateway_request_body(signoz, "PATCH", gateway_url)
assert body is not None, "Expected a PATCH request to reach the gateway"
assert body["config"]["day"]["size"] == 4000
assert body["config"]["day"]["count"] == 250
assert body["tags"] == ["test"]
def test_delete_ingestion_limit(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
make_http_mocks: Callable[[types.TestContainerDocker, list], None],
get_token: Callable[[str, str], str],
) -> None:
editor_token = get_token(GATEWAY_APIS_EDITOR_EMAIL, GATEWAY_APIS_EDITOR_PASSWORD)
gateway_url = f"/v1/workspaces/me/limits/{TEST_LIMIT_ID}"
make_http_mocks(
signoz.gateway,
[
Mapping(
request=MappingRequest(
method=HttpMethods.DELETE,
url=gateway_url,
headers=common_gateway_headers(),
),
response=MappingResponse(status=204),
persistent=False,
),
],
)
response = requests.delete(
signoz.self.host_configs["8080"].get(f"/api/v2/gateway/ingestion_limits/{TEST_LIMIT_ID}"),
headers={"Authorization": f"Bearer {editor_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, f"Expected 204, got {response.status_code}: {response.text}"
matched = get_gateway_requests(signoz, "DELETE", gateway_url)
assert len(matched) >= 1, "Expected a DELETE request to reach the gateway"

View File

@@ -209,6 +209,52 @@ def test_ai_span_list_excludes_non_gen_ai_spans(
assert "POST /api/chat" not in names # root span excluded
def test_ai_span_list_trace_level_filter(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
"""Span list (raw) with a trace-level condition returns only the gen_ai spans
of traces whose window-clipped aggregates qualify."""
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
service = "ai-it-spanlist-tracefilter"
small = ai_trace(now=now, service=service, user="a", in_tokens=10, out_tokens=100)
large = ai_trace(now=now, service=service, user="b", in_tokens=30, out_tokens=300)
insert_traces(small + large)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms, end_ms = query_window(now)
query = BuilderQuery(
signal="traces",
query_type="builder_ai_query",
name="A",
filter_expression=f"service.name = '{service}' AND trace.output_tokens > 100",
limit=10,
)
response = make_query_request(signoz, token, start_ms, end_ms, [query.to_dict()], request_type=RequestType.RAW)
assert response.status_code == HTTPStatus.OK, response.text
rows = response.json()["data"]["data"]["results"][0]["rows"]
assert len(rows) == 1, f"expected only the large trace's LLM span, got {len(rows)} rows"
body = json.dumps(rows)
assert large[0].trace_id in body
assert small[0].trace_id not in body
# a threshold no trace meets: the empty qualification yields no spans, not an error
query = BuilderQuery(
signal="traces",
query_type="builder_ai_query",
name="A",
filter_expression=f"service.name = '{service}' AND trace.output_tokens > 1000",
limit=10,
)
response = make_query_request(signoz, token, start_ms, end_ms, [query.to_dict()], request_type=RequestType.RAW)
assert response.status_code == HTTPStatus.OK, response.text
assert not (response.json()["data"]["data"]["results"][0].get("rows") or [])
def test_ai_list_having_or_aggregates(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument

View File

@@ -2,10 +2,12 @@ from collections.abc import Callable
from datetime import UTC, datetime
from http import HTTPStatus
import pytest
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.metadata import get_field_keys, get_field_values
from fixtures.querierai import ai_trace
from fixtures.metadata import AttributesMetadata, get_field_keys, get_field_values
from fixtures.querierai import ai_trace, ai_trace_mixed_spans
from fixtures.traces import Traces
AI_KEYS_PATH = "/api/v1/ai_observability/fields/keys"
@@ -106,20 +108,94 @@ def test_ai_field_values_suggests_ingested_attribute_values(
assert values["stringValues"] == ["gpt-it-values"], values
def test_ai_field_values_reject_existing_query(
@pytest.mark.parametrize(
"existing_query,search_text,expected",
[
pytest.param(None, "", {"ai-rel-a", "ai-rel-b", "ai-rel-c"}, id="no_query_scopes_to_gen_ai_spans"),
pytest.param("gen_ai.user.id = 'alice'", "", {"ai-rel-a"}, id="span_filter_narrows_under_the_gate"),
pytest.param("llm_call_count > 0", "", {"ai-rel-a", "ai-rel-b", "ai-rel-c"}, id="pure_trace_aggregate_filter_is_stripped"),
pytest.param("llm_call_count > 0 AND gen_ai.user.id = 'alice'", "", {"ai-rel-a"}, id="mixed_filter_keeps_only_the_span_part"),
pytest.param(
"llm_call_count > 0 OR gen_ai.user.id = 'alice'",
"",
{"ai-rel-a", "ai-rel-b", "ai-rel-c"},
id="class_mixing_or_drops_the_filter_not_the_request",
),
pytest.param(
"gen_ai.user.id = ",
"",
{"ai-rel-a", "ai-rel-b", "ai-rel-c"},
id="unparseable_filter_falls_back_to_the_gate",
),
pytest.param(None, "ai-rel-a", {"ai-rel-a"}, id="search_text_narrows_related_values"),
# http.request.method lives on the root span's metadata row, gen_ai.* on
# the LLM/tool/agent rows; rows are per span-shape, so the gate AND a
# cross-span attribute filter can match no single row
pytest.param("http.request.method = 'POST'", "", set(), id="cross_span_attribute_filter_matches_no_row"),
],
)
def test_ai_field_values_related_values(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
insert_attributes_metadata: Callable[[list[AttributesMetadata]], None],
existing_query: str | None,
search_text: str,
expected: set[str],
) -> None:
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
# existingQuery key resolution reads the trace keys tables, not
# attributes_metadata; a mixed trace registers the gate keys (model/tool/
# agent) plus gen_ai.user.id and http.request.method
insert_traces(ai_trace_mixed_spans(now=now, service="ai-rel-a", user="alice"))
# related values are served from attributes_metadata; one row per gate key,
# the traces row without any gate attribute and the logs row (wrong
# data_source, gate attribute present) must never surface
insert_attributes_metadata(
[
AttributesMetadata(
data_source="traces",
resource_attributes={"service.name": "ai-rel-a"},
attributes={"gen_ai.request.model": "gpt-rel", "gen_ai.user.id": "alice"},
),
AttributesMetadata(
data_source="traces",
resource_attributes={"service.name": "ai-rel-b"},
attributes={"gen_ai.tool.name": "get_weather", "gen_ai.user.id": "bob"},
),
AttributesMetadata(
data_source="traces",
resource_attributes={"service.name": "ai-rel-c"},
attributes={"gen_ai.agent.name": "chat-agent"},
),
AttributesMetadata(
data_source="traces",
resource_attributes={"service.name": "plain-rel"},
attributes={"http.request.method": "POST"},
),
AttributesMetadata(
data_source="logs",
resource_attributes={"service.name": "ai-rel-logs"},
attributes={"gen_ai.request.model": "gpt-rel"},
),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = get_field_values(
signoz,
token,
{"name": "gen_ai.request.model", "existingQuery": "service.name = 'ai-it-values'"},
AI_VALUES_PATH,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
params = {"name": "service.name", "searchText": search_text}
if existing_query is not None:
params["existingQuery"] = existing_query
response = get_field_values(signoz, token, params, AI_VALUES_PATH)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["status"] == "success"
related = response.json()["data"]["values"].get("relatedValues") or []
assert set(related) == expected, related
def test_ai_field_values_of_computed_aggregate_are_empty(

View File

@@ -0,0 +1,276 @@
from collections.abc import Callable
from http import HTTPStatus
import requests
from sqlalchemy import sql
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
ALL_SOURCES = {
"traces",
"logs",
"api_monitoring",
"exceptions",
"meter",
"ai_observability",
}
def test_get_quick_filters_returns_defaults(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
data = response.json()["data"]
assert {source_filters["source"] for source_filters in data} == ALL_SOURCES
for source_filters in data:
assert source_filters["id"] != "00000000-0000-0000-0000-000000000000"
assert source_filters["orgId"] != "00000000-0000-0000-0000-000000000000"
assert source_filters["createdAt"] != ""
assert source_filters["updatedAt"] != ""
assert len(source_filters["filters"]) > 0
for field_key in source_filters["filters"]:
assert field_key["name"] != ""
assert "fieldContext" in field_key
assert "fieldDataType" in field_key
assert "key" not in field_key
def test_v1_get_serves_legacy_shape(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
data = response.json()["data"]
assert {source_filters["signal"] for source_filters in data} == ALL_SOURCES
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters/traces"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
filters = response.json()["data"]["filters"]
assert filters[0]["key"] == "duration_nano"
assert filters[0]["type"] == "tag"
assert filters[0]["dataType"] == "float64"
assert all("name" not in legacy_filter for legacy_filter in filters)
def test_v1_update_round_trips_to_v2(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.put(
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters"),
json={
"signal": "exceptions",
"filters": [
{"key": "service.name", "dataType": "string", "type": "resource"},
{"key": "http.method", "dataType": "string", "type": "tag"},
{"key": "code_line", "dataType": "int64", "type": "tag"},
],
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/exceptions"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
filters = response.json()["data"]["filters"]
assert [(field_key["name"], field_key["fieldContext"]) for field_key in filters] == [
("service.name", "resource"),
("http.method", "attribute"),
("code_line", "attribute"),
]
assert filters[2]["fieldDataType"] == "number"
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters/exceptions"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
filters = response.json()["data"]["filters"]
assert [(legacy_filter["key"], legacy_filter["type"]) for legacy_filter in filters] == [
("service.name", "resource"),
("http.method", "tag"),
("code_line", "tag"),
]
response = requests.put(
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters"),
json={
"signal": "meter",
"filters": [{"key": "host.name", "dataType": "string", "type": ""}],
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/meter"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
assert [(field_key["name"], field_key["signal"]) for field_key in response.json()["data"]["filters"]] == [("host.name", "metrics")]
def test_update_quick_filters_round_trip(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.put(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/logs"),
json={
"filters": [
{
"name": "k8s.pod.name",
"fieldContext": "resource",
"fieldDataType": "string",
},
{
"name": "body.status",
"fieldContext": "body",
"fieldDataType": "string",
},
],
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/logs"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
filters = response.json()["data"]["filters"]
assert [field_key["name"] for field_key in filters] == [
"k8s.pod.name",
"body.status",
]
assert filters[0]["fieldContext"] == "resource"
assert filters[1]["fieldContext"] == "body"
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters/logs"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
assert [(legacy_filter["key"], legacy_filter["type"]) for legacy_filter in response.json()["data"]["filters"]] == [
("k8s.pod.name", "resource"),
("body.status", ""),
]
def test_update_quick_filters_creates_row_for_source_without_one(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
with signoz.sqlstore.conn.connect() as conn:
conn.execute(
sql.text("DELETE FROM quick_filter WHERE source = :source"),
{"source": "api_monitoring"},
)
conn.commit()
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/api_monitoring"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
data = response.json()["data"]
assert data["source"] == "api_monitoring"
assert data["filters"] == []
assert data["id"] == "00000000-0000-0000-0000-000000000000"
response = requests.put(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/api_monitoring"),
json={
"filters": [
{
"name": "service.name",
"fieldContext": "resource",
"fieldDataType": "string",
},
],
},
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/api_monitoring"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.OK, response.text
data = response.json()["data"]
assert [field_key["name"] for field_key in data["filters"]] == ["service.name"]
assert data["id"] != "00000000-0000-0000-0000-000000000000"
def test_update_quick_filters_rejects_invalid_input(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
for source, invalid_body in [
(
"traces",
{"filters": [{"key": "service.name", "dataType": "string", "type": "resource"}]},
),
("invalid", {"filters": []}),
]:
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v2/quick_filters/{source}"),
json=invalid_body,
headers={"Authorization": f"Bearer {admin_token}"},
timeout=2,
)
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text