Compare commits

..

64 Commits

Author SHA1 Message Date
nikhilmantri0902
34482d8cc1 refactor(ruletypes): type the rule view data column with a storage-owned struct 2026-09-28 11:18:31 +05:30
nikhilmantri0902
a541d2687c Merge branch 'main' into feat/alert_rule_views
# Conflicts:
#	docs/api/openapi.yml
#	tests/integration/tests/ruler/07_rule_views.py
2026-09-28 10:45:51 +05:30
nikhilmantri0902
6ab70762a7 fix(ruler): use a snake_case slog key for rule view decode errors 2026-09-24 13:42:12 +05:30
nikhilmantri0902
8989e0a3a1 chore(api): regenerate openapi spec and frontend client 2026-09-24 13:34:00 +05:30
nikhilmantri0902
0d150554c7 refactor(ruletypes): split rule view storage and wire shapes
StorableRuleView and GettableRuleView meet only in converters; the list
conversion is a pure types function with per-id errors the provider logs.
Also renames test cases to the PascalCase convention, drops the dead
ensure_notification_channel helper and restores the create_rule_view
fixture lost in the main merge.
2026-09-24 13:34:00 +05:30
nikhilmantri0902
bee75c8e34 Merge branch 'main' into feat/alert_rule_views
# Conflicts:
#	docs/api/openapi.yml
#	ee/sqlstore/postgressqlstore/formatter_test.go
#	frontend/src/api/generated/services/rules/index.ts
#	frontend/src/api/generated/services/sigNoz.schemas.ts
#	pkg/query-service/rules/manager.go
#	pkg/query-service/rules/manager_test.go
#	pkg/ruler/rulestore/rulestoretest/rule.go
#	pkg/ruler/rulestore/sqlrulestore/rule.go
#	pkg/signoz/provider.go
#	pkg/sqlstore/sqlitesqlstore/formatter_test.go
#	pkg/sqlstore/sqlstoretest/formatter_test.go
#	pkg/types/ruletypes/list.go
#	pkg/types/ruletypes/list_test.go
#	pkg/types/ruletypes/listable_rule_test.go
#	pkg/types/ruletypes/rule.go
#	tests/fixtures/alerts.py
2026-09-24 13:09:42 +05:30
Nikhil Mantri
eb54542ca1 Merge branch 'feat/alerts_listing_page_revamp' into feat/alert_rule_views 2026-09-22 12:17:03 +05:30
nikhilmantri0902
c92fe4969a fix(ruletypes): rank nodata above pending in the state display order 2026-09-22 12:00:13 +05:30
nikhilmantri0902
97a2fefe61 Merge branch 'feat/alerts_listing_page_revamp' into feat/alert_rule_views
# Conflicts:
#	pkg/signoz/provider.go
#	pkg/types/ruletypes/list.go
#	pkg/types/ruletypes/list_test.go
2026-09-22 11:32:54 +05:30
Nikhil Mantri
c059d3b8f6 Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-22 11:16:36 +05:30
nikhilmantri0902
37d342eb90 refactor(ruletypes): move the storable-to-listable loop into the types package 2026-09-22 11:16:08 +05:30
nikhilmantri0902
d60da7b8a7 docs(sqlrulestore): fix stale resolver comment on unknown keys 2026-09-22 11:03:52 +05:30
nikhilmantri0902
8d5111bc7d refactor(ruletypes): rename storable-to-listable converter to ToListableRule method 2026-09-22 10:59:50 +05:30
nikhilmantri0902
e2578c15fb revert(rules): drop the drive-by TriggeredAlerts read-lock fix 2026-09-21 16:59:07 +05:30
Nikhil Mantri
a446a832ae Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-21 16:48:30 +05:30
nikhilmantri0902
c89c19ce9e chore(api): regenerate openapi spec and frontend client 2026-09-21 16:43:44 +05:30
nikhilmantri0902
e543b3ef32 feat(sqlrulestore): match bare label keys and both sides of reserved-key collisions
A non-reserved key filters the rule labels directly. A key colliding with a
reserved keyword matches either interpretation, with negative operators
excluding both, mirroring the v5 querier's ambiguous-key semantics.
labels.<key> stays the explicit label-only form.
2026-09-21 16:43:43 +05:30
nikhilmantri0902
3f7531f444 refactor(ruletypes): co-locate state display rank with severity and pin exhaustiveness 2026-09-21 14:30:32 +05:30
nikhilmantri0902
abfa9bee1e fix(rules): validate list params in the manager for non-API callers 2026-09-21 14:13:27 +05:30
nikhilmantri0902
2f497108df fix(ruletypes): lower rules list max page size to 200 2026-09-21 13:54:44 +05:30
nikhilmantri0902
1fa1e292e9 test(sqlstore): cover JSONExtractMapValue in sqlite, postgres and test formatters 2026-09-21 13:47:25 +05:30
nikhilmantri0902
a328988a15 chore: alertStates -> GetAlertStates 2026-09-21 13:28:45 +05:30
nikhilmantri0902
7e53bd607d chore: added ruletypes layer in between structs 2026-09-21 13:26:23 +05:30
Nikhil Mantri
4d82eb8436 Merge branch 'feat/alerts_listing_page_revamp' into feat/alert_rule_views 2026-09-17 12:34:36 +05:30
Nikhil Mantri
95c11d0550 Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-17 12:34:23 +05:30
Nikhil Mantri
b2810ec117 Merge branch 'feat/alerts_listing_page_revamp' into feat/alert_rule_views 2026-09-16 15:48:59 +05:30
Naman Verma
7556dd2868 Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-16 11:28:20 +05:30
Nikhil Mantri
339a10217e Merge branch 'feat/alerts_listing_page_revamp' into feat/alert_rule_views 2026-09-10 18:57:09 +05:30
Nikhil Mantri
cd0eb62734 Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-10 18:56:27 +05:30
nikhilmantri0902
a667007991 Merge remote-tracking branch 'origin/feat/alerts_listing_page_revamp' into feat/alert_rule_views 2026-09-10 16:04:32 +05:30
nikhilmantri0902
208fc1e8ec Merge remote-tracking branch 'origin/feat/common_out_visitors_and_sql_parser' into feat/alerts_listing_page_revamp 2026-09-10 15:58:38 +05:30
nikhilmantri0902
54eb67392a refactor(filterquery): order the compiler main flow first
Visit methods follow Compile; builders, extractors and operator
spelling sit below with main's section separators.
2026-09-10 15:53:07 +05:30
nikhilmantri0902
0e0d97c16f refactor(filterquery): address review nits on the compiler
Verb-first helper names (BuildStringOperation etc., ResolveFreeText),
Sb and Formatter as exported fields instead of getters, and tests for
the dangling-backslash rejection.
2026-09-10 15:43:58 +05:30
nikhilmantri0902
659f9c7ecf test(ruletypes): assert UpdatedAt strictly advances on view update 2026-09-10 13:09:54 +05:30
nikhilmantri0902
ce8442c17a chore(ruler): drop rationale comments from rule view code 2026-09-10 13:09:53 +05:30
nikhilmantri0902
7ab21678b5 test(alerts): own rule view cleanup in a fixture
The lifecycle test filters lists by its own view names instead of wiping
the org's views and asserting global counts.
2026-09-10 13:02:47 +05:30
nikhilmantri0902
951bfb66cd fix(ruletypes): never return null states on a rule view
Validate normalizes nil states to an empty slice and the field is marked
non-nullable; regenerates the openapi spec and the frontend api client.
2026-09-10 13:02:42 +05:30
nikhilmantri0902
113685fbdb refactor(sqlmigration): create the rule_view index via bun builder
Matches the package's index precedent and idx_ naming.
2026-09-10 12:45:05 +05:30
nikhilmantri0902
b9c306dd26 refactor(ruletypes): share list filter validation between params and views
ListRulesParams and RuleViewData now embed one ListFilter with a single
Validate, mirroring dashboards; filter errors surface as rule_list_invalid.
2026-09-10 12:36:24 +05:30
nikhilmantri0902
7e71d4507c Merge remote-tracking branch 'origin/feat/common_out_visitors_and_sql_parser' into feat/alerts_listing_page_revamp 2026-09-09 20:10:09 +05:30
nikhilmantri0902
e13cb08097 refactor(filterquery): collapse the builder into the visitor
One Visitor struct now carries sb, fmter and errors like the old
per-feature visitors did; FieldResolver stays the only new concept.
2026-09-09 20:08:33 +05:30
nikhilmantri0902
1e7aaa9f14 feat(ruler): add saved views CRUD for the alerts list page
A saved view stores the v3 list params (query, states, sort, order) in a
new rule_view table and replays them; shapes mirror dashboard views.
2026-09-09 15:57:47 +05:30
nikhilmantri0902
09ab9f4081 Merge branch 'feat/common_out_visitors_and_sql_parser' into feat/alerts_listing_page_revamp 2026-09-09 11:41:05 +05:30
nikhilmantri0902
f945b5f513 refactor(dashboard): port the list filter to the shared sqlcompiler
The visitor moves to a key-policy resolver plus an error-code wrap;
emitted SQL is unchanged, pinned by the existing exact-SQL unit suite.
2026-09-09 11:33:35 +05:30
nikhilmantri0902
8b0ab0ff26 refactor(filterquery): add shared list filter SQL compiler
Extracted from the dashboards list visitor: grammar walk, operator
dispatch, value extraction, LIKE builders and the Compiled output type,
behind a per-feature FieldResolver. Scope is list pages over the
relational store; telemetry queries stay on querybuilder. Also rejects
LIKE and ILIKE patterns ending in an unescaped backslash, which never
match on sqlite and abort the query on Postgres.
2026-09-09 11:33:34 +05:30
Nikhil Mantri
f2679a6866 Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-08 17:47:29 +05:30
nikhilmantri0902
253a117849 refactor: address list API nits
Use fmt.Sprintf over string concatenation in the compiler, rule store
and JSON formatters; build the enum value lists once as package vars;
collapse the repeated integration seed blocks into a seed_alert_rules
fixture and drop name-restating fixture docstrings.
2026-09-08 17:42:46 +05:30
nikhilmantri0902
c9538be38f style: trim list API comments to one line per the repo comment rules 2026-09-08 17:16:18 +05:30
nikhilmantri0902
be51c37317 fix(sqlcompiler): reject LIKE patterns ending in an unescaped backslash
Such a pattern never matches on sqlite and aborts the query on Postgres
when the matcher consumes the dangling escape, turning user input into
a data-dependent 500. Reject it as a 400 on both list endpoints; a
literal trailing backslash stays expressible as an escaped backslash.
Also move the Compiled type into sqlcompiler, leaving each feature only
the error-code wrap.
2026-09-08 17:16:11 +05:30
nikhilmantri0902
66796d1767 refactor: extract shared filter query SQL compiler from dashboards and rules
The dashboards and rules list visitors were near-identical: boolean
composition, operator dispatch, value extraction and the LIKE family
builders. Move that core to pkg/parser/filterquery/sqlcompiler behind a
FieldResolver interface; each feature keeps only its key policy. Scope
is the bun-managed relational store, telemetry stays on querybuilder.
Verified by both features' unchanged exact-SQL unit suites, the
dashboard integration suite and the rules list suite on both providers.
2026-09-08 15:44:38 +05:30
nikhilmantri0902
8e19b17855 test(sqlrulestore): pin deep-nesting compile cases before visitor extraction
Four hand-derived cases covering mixed predicate kinds, parenthesized
OR groups, NOT over a group and three-level nesting with free text and
a timestamp, matching the depth of dashboards' ComplexExamples suite.
2026-09-08 15:06:11 +05:30
nikhilmantri0902
8885497e14 chore: further cleanup 2026-09-08 12:31:33 +05:30
nikhilmantri0902
626e047485 refactor: trim comments to constraint-only per repo comment rules
Rationale and restatements move out of source; dialect reasoning for
LIKE ESCAPE and ILIKE lowering is documented once in the PR body.
2026-09-08 12:13:16 +05:30
nikhilmantri0902
cc3d86c3ec fix(rules): take the read lock in TriggeredAlerts
The lock has been commented out since the method landed, leaving the
rules map read unguarded against concurrent manager writes. Same class
of race this branch fixed in ListRuleStates and GetRule.
2026-09-08 11:54:05 +05:30
nikhilmantri0902
0100083284 fix(sqlstore): drop unusable quote escape from sqlite JSON map paths
sqlite JSON paths have no backslash escapes, so escaping a double quote
in the key produced a path sqlite cannot parse. The character is also
unreachable: the filter grammar's KEY token cannot contain a quote.
Keep the backslash escape, which matches sqlite's raw key comparison.
2026-09-08 11:44:05 +05:30
Nikhil Mantri
648e945471 Merge branch 'main' into feat/alerts_listing_page_revamp 2026-09-08 11:07:26 +05:30
nikhilmantri0902
7323d5ae0b fix(ruletypes): make list sort deterministic on ties
Ties on the primary sort key kept arbitrary DB order, so rows could
shuffle between page requests causing overlaps or misses. Break ties on
name (case-insensitive) then id, always ascending, with the requested
order applied to the primary key only. Pin the tiebreak in unit tests
and paginate over state ties in the integration suite.
2026-09-08 01:16:39 +05:30
nikhilmantri0902
578b9172ba test(integration): cover the rules list v3 API on both sqlstore providers
Seeds five rules spanning every filter axis and pins the envelope, slim rows,
DSL filters (incl. dotted label keys and missing-label-as-empty semantics),
states param, sort ranks, pagination totals and the error contract; adds
delete_all_rules and an idempotent ensure_notification_channel fixture fn.
Verified against both sqlite and postgres providers.
2026-09-07 20:20:54 +05:30
nikhilmantri0902
2110006d17 feat(sqlrulestore): treat a missing label as empty string for all label value operators
COALESCE applies uniformly instead of only on negations, so labels.key = ''
also matches label-less rules; EXISTS/NOT EXISTS remain the presence checks.
2026-09-07 17:39:41 +05:30
nikhilmantri0902
01dda19d17 feat(sqlrulestore): align label filter semantics with the querier
severity takes the labels operator set (it is an alias for labels.severity,
so EXISTS/NOT EXISTS now work on it), and negative label operators evaluate a
missing label as the empty string instead of always matching, mirroring the
querier's AddDefaultExistsFilter map-attribute semantics.
2026-09-07 17:27:07 +05:30
nikhilmantri0902
b519574b89 feat(ruler): add GET /api/v3/rules route and regenerate API clients
Registers ListRulesV3 with the list params and envelope, marks the v2
ListRules operation deprecated, and regenerates the OpenAPI spec and the
frontend client.
2026-09-07 14:54:46 +05:30
nikhilmantri0902
629ecabce7 feat(ruler): list rules with SQL filter pushdown, state overlay and in-memory pagination
DSL query compiles into the store's WHERE; state is overlaid from a snapshot
of the rule manager map taken under RLock, then the states filter, total,
sort (state/severity rank comparators) and offset/limit run in code so total
always matches what is pageable. StorableRule gains alias:rule to match the
compiler's column references. Also guards the previously unlocked m.rules
reads in ListRuleStates and GetRule.
2026-09-07 14:43:40 +05:30
nikhilmantri0902
131e302c55 feat(sqlrulestore): compile rule list filter DSL to SQL
Adds JSONExtractMapValue to SQLFormatter so a labels map key is one path
segment (dotted label keys work on both dialects), and a visitor over the
shared filterquery grammar mapping rule list DSL keys to SQL.
2026-09-07 13:58:30 +05:30
nikhilmantri0902
239c92ba67 feat(ruletypes): add list params, filter allow-lists and listable rule types for rules list API 2026-09-07 12:48:49 +05:30
36 changed files with 1958 additions and 2706 deletions

View File

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

View File

@@ -1,48 +1,5 @@
components:
schemas:
AiobservabilitytypesMessage:
properties:
content:
items:
$ref: '#/components/schemas/AiobservabilitytypesPart'
type: array
finishReason:
type: string
role:
type: string
required:
- content
type: object
AiobservabilitytypesPart:
properties:
arguments: {}
content:
type: string
id:
type: string
isError:
type: boolean
name:
type: string
redacted:
type: boolean
server:
type: boolean
toolCallId:
type: string
type:
$ref: '#/components/schemas/AiobservabilitytypesPartType'
required:
- type
type: object
AiobservabilitytypesPartType:
enum:
- text
- thinking
- tool_call
- tool_result
- generic
type: string
AlertmanagertypesChannel:
properties:
createdAt:
@@ -8971,6 +8928,30 @@ components:
- kind
- spec
type: object
RuletypesGettableRuleView:
properties:
createdAt:
format: date-time
type: string
data:
$ref: '#/components/schemas/RuletypesRuleViewData'
id:
type: string
name:
type: string
orgId:
type: string
updatedAt:
format: date-time
type: string
required:
- id
- name
- data
- orgId
- createdAt
- updatedAt
type: object
RuletypesGettableTestRule:
properties:
alertCount:
@@ -9038,6 +9019,15 @@ components:
- alertType
- ruleType
type: object
RuletypesListableRuleViews:
properties:
views:
items:
$ref: '#/components/schemas/RuletypesGettableRuleView'
type: array
required:
- views
type: object
RuletypesListableRules:
properties:
labels:
@@ -9135,6 +9125,16 @@ components:
- ruleType
- condition
type: object
RuletypesPostableRuleView:
properties:
data:
$ref: '#/components/schemas/RuletypesRuleViewData'
name:
type: string
required:
- name
- data
type: object
RuletypesQueryType:
enum:
- builder
@@ -9273,6 +9273,23 @@ components:
- promql_rule
- anomaly_rule
type: string
RuletypesRuleViewData:
properties:
order:
$ref: '#/components/schemas/RuletypesListOrder'
query:
type: string
sort:
$ref: '#/components/schemas/RuletypesListSort'
states:
items:
type: string
type: array
version:
type: string
required:
- version
type: object
RuletypesScheduleType:
enum:
- hourly
@@ -9690,17 +9707,6 @@ components:
required:
- aggregations
type: object
SpantypesGettableTraceThread:
properties:
nextCursor:
type: string
spans:
items:
$ref: '#/components/schemas/SpantypesThreadSpan'
type: array
required:
- spans
type: object
SpantypesGettableWaterfallTrace:
properties:
endTimestampMillis:
@@ -10015,92 +10021,6 @@ components:
nullable: true
type: object
type: object
SpantypesThreadSpan:
properties:
attributes:
additionalProperties: {}
nullable: true
type: object
db_name:
type: string
db_operation:
type: string
duration_nano:
minimum: 0
type: integer
events:
items:
$ref: '#/components/schemas/SpantypesEvent'
nullable: true
type: array
external_http_method:
type: string
external_http_url:
type: string
flags:
minimum: 0
type: integer
formatted_input:
items:
$ref: '#/components/schemas/AiobservabilitytypesMessage'
type: array
formatted_output:
items:
$ref: '#/components/schemas/AiobservabilitytypesMessage'
type: array
has_children:
type: boolean
has_error:
type: boolean
http_host:
type: string
http_method:
type: string
http_url:
type: string
is_remote:
type: string
kind_string:
type: string
level:
minimum: 0
type: integer
name:
type: string
parent_span_id:
type: string
references:
items:
$ref: '#/components/schemas/SpantypesOtelSpanRef'
type: array
resource:
additionalProperties:
type: string
nullable: true
type: object
response_status_code:
type: string
span_id:
type: string
status_code:
type: integer
status_code_string:
type: string
status_message:
type: string
sub_tree_node_count:
minimum: 0
type: integer
time_unix:
minimum: 0
type: integer
trace_id:
type: string
trace_state:
type: string
required:
- references
type: object
SpantypesUpdatableSpanMapper:
properties:
config:
@@ -15825,81 +15745,6 @@ paths:
tags:
- tracedetail
x-signoz-stability: alpha
/api/v1/traces/{traceID}/thread:
get:
deprecated: false
description: Returns the spans carrying gen_ai input or output messages in timestamp
order, each with the messages normalised into formatted_input and formatted_output.
Pages are fetched with the returned nextCursor.
operationId: GetTraceThread
parameters:
- in: query
name: limit
schema:
type: integer
- in: query
name: cursor
schema:
type: string
- in: path
name: traceID
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/SpantypesGettableTraceThread'
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
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: Get thread view for a trace
tags:
- tracedetail
x-signoz-stability: alpha
/api/v1/user/me:
get:
deprecated: true
@@ -21362,6 +21207,235 @@ paths:
tags:
- users
x-signoz-stability: alpha
/api/v2/rule_views:
get:
deprecated: false
description: Returns every saved view in the calling user's org. Saved views
are shared org-wide.
operationId: ListRuleViews
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/RuletypesListableRuleViews'
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
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: List rule saved views
tags:
- rules
x-signoz-stability: alpha
post:
deprecated: false
description: Persists the calling user's rule listing state (query, states,
sort, order) as a named, reusable view shared across the org.
operationId: CreateRuleView
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/RuletypesPostableRuleView'
responses:
"201":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/RuletypesGettableRuleView'
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:
- VIEWER
- tokenizer:
- VIEWER
summary: Create rule saved view
tags:
- rules
x-signoz-stability: alpha
/api/v2/rule_views/{id}:
delete:
deprecated: false
description: Removes a saved view. Saved views are shared org-wide. Deleting
a non-existent view returns 404.
operationId: DeleteRuleView
parameters:
- in: path
name: id
required: true
schema:
type: string
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
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: Delete rule saved view
tags:
- rules
x-signoz-stability: alpha
put:
deprecated: false
description: Replaces a saved view's name and data. Saved views are shared org-wide.
operationId: UpdateRuleView
parameters:
- in: path
name: id
required: true
schema:
type: string
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/RuletypesPostableRuleView'
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/RuletypesGettableRuleView'
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
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: Update rule saved view
tags:
- rules
x-signoz-stability: alpha
/api/v2/rules:
get:
deprecated: true

View File

@@ -19,7 +19,9 @@ import type {
import type {
CreateRule201,
CreateRuleView201,
DeleteRuleByIDPathParameters,
DeleteRuleViewPathParameters,
GetRuleByID200,
GetRuleByIDPathParameters,
GetRuleHistoryFilterKeys200,
@@ -40,6 +42,7 @@ import type {
GetRuleHistoryTopContributors200,
GetRuleHistoryTopContributorsParams,
GetRuleHistoryTopContributorsPathParameters,
ListRuleViews200,
ListRules200,
ListRulesV3200,
ListRulesV3Params,
@@ -47,8 +50,11 @@ import type {
PatchRuleByIDPathParameters,
RenderErrorResponseDTO,
RuletypesPostableRuleDTO,
RuletypesPostableRuleViewDTO,
TestRule200,
UpdateRuleByIDPathParameters,
UpdateRuleView200,
UpdateRuleViewPathParameters,
} from '../sigNoz.schemas';
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
@@ -74,6 +80,351 @@ const withQueryKey = <T extends object, K>(
return result;
};
/**
* Returns every saved view in the calling user's org. Saved views are shared org-wide.
* @summary List rule saved views
*/
export const listRuleViews = (signal?: AbortSignal) => {
return GeneratedAPIInstance<ListRuleViews200>({
url: `/api/v2/rule_views`,
method: 'GET',
signal,
});
};
export const getListRuleViewsQueryKey = () => {
return [`/api/v2/rule_views`] as const;
};
export const getListRuleViewsQueryOptions = <
TData = Awaited<ReturnType<typeof listRuleViews>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listRuleViews>>,
TError,
TData
>;
}) => {
const { query: queryOptions } = options ?? {};
const queryKey = queryOptions?.queryKey ?? getListRuleViewsQueryKey();
const queryFn: QueryFunction<Awaited<ReturnType<typeof listRuleViews>>> = ({
signal,
}) => listRuleViews(signal);
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
Awaited<ReturnType<typeof listRuleViews>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type ListRuleViewsQueryResult = NonNullable<
Awaited<ReturnType<typeof listRuleViews>>
>;
export type ListRuleViewsQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary List rule saved views
*/
export function useListRuleViews<
TData = Awaited<ReturnType<typeof listRuleViews>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof listRuleViews>>,
TError,
TData
>;
}): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getListRuleViewsQueryOptions(options);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return withQueryKey(query, queryOptions.queryKey);
}
/**
* @summary List rule saved views
*/
export const invalidateListRuleViews = async (
queryClient: QueryClient,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getListRuleViewsQueryKey() },
options,
);
return queryClient;
};
/**
* Persists the calling user's rule listing state (query, states, sort, order) as a named, reusable view shared across the org.
* @summary Create rule saved view
*/
export const createRuleView = (
ruletypesPostableRuleViewDTO?: BodyType<RuletypesPostableRuleViewDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<CreateRuleView201>({
url: `/api/v2/rule_views`,
method: 'POST',
headers: { 'Content-Type': 'application/json' },
data: ruletypesPostableRuleViewDTO,
signal,
});
};
export const getCreateRuleViewMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof createRuleView>>,
TError,
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof createRuleView>>,
TError,
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
TContext
> => {
const mutationKey = ['createRuleView'];
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 createRuleView>>,
{ data?: BodyType<RuletypesPostableRuleViewDTO> }
> = (props) => {
const { data } = props ?? {};
return createRuleView(data);
};
return { mutationFn, ...mutationOptions };
};
export type CreateRuleViewMutationResult = NonNullable<
Awaited<ReturnType<typeof createRuleView>>
>;
export type CreateRuleViewMutationBody =
| BodyType<RuletypesPostableRuleViewDTO>
| undefined;
export type CreateRuleViewMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Create rule saved view
*/
export const useCreateRuleView = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof createRuleView>>,
TError,
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof createRuleView>>,
TError,
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
TContext
> => {
return useMutation(getCreateRuleViewMutationOptions(options));
};
/**
* Removes a saved view. Saved views are shared org-wide. Deleting a non-existent view returns 404.
* @summary Delete rule saved view
*/
export const deleteRuleView = (
{ id }: DeleteRuleViewPathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<void>({
url: `/api/v2/rule_views/${id}`,
method: 'DELETE',
signal,
});
};
export const getDeleteRuleViewMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof deleteRuleView>>,
TError,
{ pathParams: DeleteRuleViewPathParameters },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof deleteRuleView>>,
TError,
{ pathParams: DeleteRuleViewPathParameters },
TContext
> => {
const mutationKey = ['deleteRuleView'];
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 deleteRuleView>>,
{ pathParams: DeleteRuleViewPathParameters }
> = (props) => {
const { pathParams } = props ?? {};
return deleteRuleView(pathParams);
};
return { mutationFn, ...mutationOptions };
};
export type DeleteRuleViewMutationResult = NonNullable<
Awaited<ReturnType<typeof deleteRuleView>>
>;
export type DeleteRuleViewMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Delete rule saved view
*/
export const useDeleteRuleView = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof deleteRuleView>>,
TError,
{ pathParams: DeleteRuleViewPathParameters },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof deleteRuleView>>,
TError,
{ pathParams: DeleteRuleViewPathParameters },
TContext
> => {
return useMutation(getDeleteRuleViewMutationOptions(options));
};
/**
* Replaces a saved view's name and data. Saved views are shared org-wide.
* @summary Update rule saved view
*/
export const updateRuleView = (
{ id }: UpdateRuleViewPathParameters,
ruletypesPostableRuleViewDTO?: BodyType<RuletypesPostableRuleViewDTO>,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<UpdateRuleView200>({
url: `/api/v2/rule_views/${id}`,
method: 'PUT',
headers: { 'Content-Type': 'application/json' },
data: ruletypesPostableRuleViewDTO,
signal,
});
};
export const getUpdateRuleViewMutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof updateRuleView>>,
TError,
{
pathParams: UpdateRuleViewPathParameters;
data?: BodyType<RuletypesPostableRuleViewDTO>;
},
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof updateRuleView>>,
TError,
{
pathParams: UpdateRuleViewPathParameters;
data?: BodyType<RuletypesPostableRuleViewDTO>;
},
TContext
> => {
const mutationKey = ['updateRuleView'];
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 updateRuleView>>,
{
pathParams: UpdateRuleViewPathParameters;
data?: BodyType<RuletypesPostableRuleViewDTO>;
}
> = (props) => {
const { pathParams, data } = props ?? {};
return updateRuleView(pathParams, data);
};
return { mutationFn, ...mutationOptions };
};
export type UpdateRuleViewMutationResult = NonNullable<
Awaited<ReturnType<typeof updateRuleView>>
>;
export type UpdateRuleViewMutationBody =
| BodyType<RuletypesPostableRuleViewDTO>
| undefined;
export type UpdateRuleViewMutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Update rule saved view
*/
export const useUpdateRuleView = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof updateRuleView>>,
TError,
{
pathParams: UpdateRuleViewPathParameters;
data?: BodyType<RuletypesPostableRuleViewDTO>;
},
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof updateRuleView>>,
TError,
{
pathParams: UpdateRuleViewPathParameters;
data?: BodyType<RuletypesPostableRuleViewDTO>;
},
TContext
> => {
return useMutation(getUpdateRuleViewMutationOptions(options));
};
/**
* This endpoint lists all alert rules with their current evaluation state. Deprecated: use ListRulesV3, which supports filtering, sorting and pagination.
* @deprecated

View File

@@ -4,61 +4,6 @@
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
export enum AiobservabilitytypesPartTypeDTO {
text = 'text',
thinking = 'thinking',
tool_call = 'tool_call',
tool_result = 'tool_result',
generic = 'generic',
}
export interface AiobservabilitytypesPartDTO {
arguments?: unknown;
/**
* @type string
*/
content?: string;
/**
* @type string
*/
id?: string;
/**
* @type boolean
*/
isError?: boolean;
/**
* @type string
*/
name?: string;
/**
* @type boolean
*/
redacted?: boolean;
/**
* @type boolean
*/
server?: boolean;
/**
* @type string
*/
toolCallId?: string;
type: AiobservabilitytypesPartTypeDTO;
}
export interface AiobservabilitytypesMessageDTO {
/**
* @type array
*/
content: AiobservabilitytypesPartDTO[];
/**
* @type string
*/
finishReason?: string;
/**
* @type string
*/
role?: string;
}
export interface AlertmanagertypesChannelDTO {
/**
* @type string
@@ -10232,6 +10177,60 @@ export enum RuletypesEvaluationKindDTO {
rolling = 'rolling',
cumulative = 'cumulative',
}
export enum RuletypesListOrderDTO {
asc = 'asc',
desc = 'desc',
}
export enum RuletypesListSortDTO {
updated_at = 'updated_at',
created_at = 'created_at',
name = 'name',
state = 'state',
severity = 'severity',
}
export interface RuletypesRuleViewDataDTO {
order?: RuletypesListOrderDTO;
/**
* @type string
*/
query?: string;
sort?: RuletypesListSortDTO;
/**
* @type array
*/
states?: string[];
/**
* @type string
*/
version: string;
}
export interface RuletypesGettableRuleViewDTO {
/**
* @type string
* @format date-time
*/
createdAt: string;
data: RuletypesRuleViewDataDTO;
/**
* @type string
*/
id: string;
/**
* @type string
*/
name: string;
/**
* @type string
*/
orgId: string;
/**
* @type string
* @format date-time
*/
updatedAt: string;
}
export interface RuletypesGettableTestRuleDTO {
/**
* @type integer
@@ -10254,17 +10253,6 @@ export interface RuletypesLabelPairDTO {
value: string;
}
export enum RuletypesListOrderDTO {
asc = 'asc',
desc = 'desc',
}
export enum RuletypesListSortDTO {
updated_at = 'updated_at',
created_at = 'created_at',
name = 'name',
state = 'state',
severity = 'severity',
}
export type RuletypesListableRuleDTOLabels = { [key: string]: string };
export enum RuletypesRuleTypeDTO {
@@ -10316,6 +10304,13 @@ export interface RuletypesListableRuleDTO {
updatedBy?: string;
}
export interface RuletypesListableRuleViewsDTO {
/**
* @type array
*/
views: RuletypesGettableRuleViewDTO[];
}
export interface RuletypesListableRulesDTO {
/**
* @type array
@@ -10484,6 +10479,14 @@ export interface RuletypesPostableRuleDTO {
version?: string;
}
export interface RuletypesPostableRuleViewDTO {
data: RuletypesRuleViewDataDTO;
/**
* @type string
*/
name: string;
}
export type RuletypesRuleDTOAnnotations = { [key: string]: string };
export type RuletypesRuleDTOLabels = { [key: string]: string };
@@ -11197,165 +11200,6 @@ export interface SpantypesOtelSpanRefDTO {
traceId?: string;
}
export type SpantypesThreadSpanDTOAttributesAnyOf = { [key: string]: unknown };
/**
* @nullable
*/
export type SpantypesThreadSpanDTOAttributes =
SpantypesThreadSpanDTOAttributesAnyOf | null;
export type SpantypesThreadSpanDTOResourceAnyOf = { [key: string]: string };
/**
* @nullable
*/
export type SpantypesThreadSpanDTOResource =
SpantypesThreadSpanDTOResourceAnyOf | null;
export interface SpantypesThreadSpanDTO {
/**
* @type object,null
*/
attributes?: SpantypesThreadSpanDTOAttributes;
/**
* @type string
*/
db_name?: string;
/**
* @type string
*/
db_operation?: string;
/**
* @type integer
* @minimum 0
*/
duration_nano?: number;
/**
* @type array,null
*/
events?: SpantypesEventDTO[] | null;
/**
* @type string
*/
external_http_method?: string;
/**
* @type string
*/
external_http_url?: string;
/**
* @type integer
* @minimum 0
*/
flags?: number;
/**
* @type array
*/
formatted_input?: AiobservabilitytypesMessageDTO[];
/**
* @type array
*/
formatted_output?: AiobservabilitytypesMessageDTO[];
/**
* @type boolean
*/
has_children?: boolean;
/**
* @type boolean
*/
has_error?: boolean;
/**
* @type string
*/
http_host?: string;
/**
* @type string
*/
http_method?: string;
/**
* @type string
*/
http_url?: string;
/**
* @type string
*/
is_remote?: string;
/**
* @type string
*/
kind_string?: string;
/**
* @type integer
* @minimum 0
*/
level?: number;
/**
* @type string
*/
name?: string;
/**
* @type string
*/
parent_span_id?: string;
/**
* @type array
*/
references: SpantypesOtelSpanRefDTO[];
/**
* @type object,null
*/
resource?: SpantypesThreadSpanDTOResource;
/**
* @type string
*/
response_status_code?: string;
/**
* @type string
*/
span_id?: string;
/**
* @type integer
*/
status_code?: number;
/**
* @type string
*/
status_code_string?: string;
/**
* @type string
*/
status_message?: string;
/**
* @type integer
* @minimum 0
*/
sub_tree_node_count?: number;
/**
* @type integer
* @minimum 0
*/
time_unix?: number;
/**
* @type string
*/
trace_id?: string;
/**
* @type string
*/
trace_state?: string;
}
export interface SpantypesGettableTraceThreadDTO {
/**
* @type string
*/
nextCursor?: string;
/**
* @type array
*/
spans: SpantypesThreadSpanDTO[];
}
export type SpantypesWaterfallSpanDTOAttributesAnyOf = {
[key: string]: unknown;
};
@@ -13029,30 +12873,6 @@ export type GetTraceAggregations200 = {
status: string;
};
export type GetTraceThreadPathParameters = {
traceID: string;
};
export type GetTraceThreadParams = {
/**
* @type integer
* @description undefined
*/
limit?: number;
/**
* @type string
* @description undefined
*/
cursor?: string;
};
export type GetTraceThread200 = {
data: SpantypesGettableTraceThreadDTO;
/**
* @type string
*/
status: string;
};
export type ListUserPreferences200 = {
/**
* @type array
@@ -13954,6 +13774,36 @@ export type GetUsersByRoleID200 = {
status: string;
};
export type ListRuleViews200 = {
data: RuletypesListableRuleViewsDTO;
/**
* @type string
*/
status: string;
};
export type CreateRuleView201 = {
data: RuletypesGettableRuleViewDTO;
/**
* @type string
*/
status: string;
};
export type DeleteRuleViewPathParameters = {
id: string;
};
export type UpdateRuleViewPathParameters = {
id: string;
};
export type UpdateRuleView200 = {
data: RuletypesGettableRuleViewDTO;
/**
* @type string
*/
status: string;
};
export type ListRules200 = {
/**
* @type array

View File

@@ -4,17 +4,11 @@
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
import { useMutation, useQuery } from 'react-query';
import { useMutation } from 'react-query';
import type {
InvalidateOptions,
MutationFunction,
QueryClient,
QueryFunction,
QueryKey,
UseMutationOptions,
UseMutationResult,
UseQueryOptions,
UseQueryResult,
} from 'react-query';
import type {
@@ -22,9 +16,6 @@ import type {
GetFlamegraphPathParameters,
GetTraceAggregations200,
GetTraceAggregationsPathParameters,
GetTraceThread200,
GetTraceThreadParams,
GetTraceThreadPathParameters,
GetWaterfallV4200,
GetWaterfallV4PathParameters,
RenderErrorResponseDTO,
@@ -36,26 +27,6 @@ import type {
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
const withQueryKey = <T extends object, K>(
query: T,
queryKey: K,
): T & { queryKey: K } => {
const result = { queryKey } as T & { queryKey: K };
for (const key of Object.keys(query)) {
// The explicit queryKey always wins, matching the previous
// `{ ...query, queryKey }` spread where it was set last.
if (key === 'queryKey') {
continue;
}
Object.defineProperty(result, key, {
enumerable: true,
configurable: true,
get: () => (query as Record<string, unknown>)[key],
});
}
return result;
};
/**
* Computes span aggregations grouped by requested field.
* @summary Get aggregations for a trace
@@ -156,121 +127,6 @@ export const useGetTraceAggregations = <
> => {
return useMutation(getGetTraceAggregationsMutationOptions(options));
};
/**
* Returns the spans carrying gen_ai input or output messages in timestamp order, each with the messages normalised into formatted_input and formatted_output. Pages are fetched with the returned nextCursor.
* @summary Get thread view for a trace
*/
export const getTraceThread = (
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<GetTraceThread200>({
url: `/api/v1/traces/${traceID}/thread`,
method: 'GET',
params,
signal,
});
};
export const getGetTraceThreadQueryKey = (
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
) => {
return [
`/api/v1/traces/${traceID}/thread`,
...(params ? [params] : []),
] as const;
};
export const getGetTraceThreadQueryOptions = <
TData = Awaited<ReturnType<typeof getTraceThread>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getTraceThread>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getGetTraceThreadQueryKey({ traceID }, params);
const queryFn: QueryFunction<Awaited<ReturnType<typeof getTraceThread>>> = ({
signal,
}) => getTraceThread({ traceID }, params, signal);
return {
queryKey,
queryFn,
enabled: traceID !== null && traceID !== undefined,
...queryOptions,
} as UseQueryOptions<
Awaited<ReturnType<typeof getTraceThread>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type GetTraceThreadQueryResult = NonNullable<
Awaited<ReturnType<typeof getTraceThread>>
>;
export type GetTraceThreadQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Get thread view for a trace
*/
export function useGetTraceThread<
TData = Awaited<ReturnType<typeof getTraceThread>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getTraceThread>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getGetTraceThreadQueryOptions(
{ traceID },
params,
options,
);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return withQueryKey(query, queryOptions.queryKey);
}
/**
* @summary Get thread view for a trace
*/
export const invalidateGetTraceThread = async (
queryClient: QueryClient,
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getGetTraceThreadQueryKey({ traceID }, params) },
options,
);
return queryClient;
};
/**
* Returns the flamegraph view of spans for a given trace ID.
* @summary Get flamegraph view for a trace

View File

@@ -132,6 +132,64 @@ func (provider *provider) addRulerRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/rule_views", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.ListRuleViews), handler.OpenAPIDef{
ID: "ListRuleViews",
Tags: []string{"rules"},
Summary: "List rule saved views",
Description: "Returns every saved view in the calling user's org. Saved views are shared org-wide.",
Response: new(ruletypes.ListableRuleViews),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
})).Methods(http.MethodGet).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/rule_views", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.CreateRuleView), handler.OpenAPIDef{
ID: "CreateRuleView",
Tags: []string{"rules"},
Summary: "Create rule saved view",
Description: "Persists the calling user's rule listing state (query, states, sort, order) as a named, reusable view shared across the org.",
Request: new(ruletypes.PostableRuleView),
RequestContentType: "application/json",
Response: new(ruletypes.GettableRuleView),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
})).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/rule_views/{id}", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.UpdateRuleView), handler.OpenAPIDef{
ID: "UpdateRuleView",
Tags: []string{"rules"},
Summary: "Update rule saved view",
Description: "Replaces a saved view's name and data. Saved views are shared org-wide.",
Request: new(ruletypes.UpdatableRuleView),
RequestContentType: "application/json",
Response: new(ruletypes.GettableRuleView),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
})).Methods(http.MethodPut).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/rule_views/{id}", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.DeleteRuleView), handler.OpenAPIDef{
ID: "DeleteRuleView",
Tags: []string{"rules"},
Summary: "Delete rule saved view",
Description: "Removes a saved view. Saved views are shared org-wide. Deleting a non-existent view returns 404.",
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
})).Methods(http.MethodDelete).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v1/downtime_schedules", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.ListDowntimeSchedules), handler.OpenAPIDef{
ID: "ListDowntimeSchedules",
Tags: []string{"downtimeschedules"},

View File

@@ -67,23 +67,5 @@ func (provider *provider) addTraceDetailRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v1/traces/{traceID}/thread", handler.New(
provider.authzMiddleware.ViewAccess(provider.traceDetailHandler.GetThread),
handler.OpenAPIDef{
ID: "GetTraceThread",
Tags: []string{"tracedetail"},
Summary: "Get thread view for a trace",
Description: "Returns the spans carrying gen_ai input or output messages in timestamp order, each with the messages normalised into formatted_input and formatted_output. Pages are fetched with the returned nextCursor.",
RequestQuery: new(spantypes.PostableThreadQuery),
Response: new(spantypes.GettableTraceThread),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
},
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
return nil
}

View File

@@ -75,25 +75,3 @@ func (h *handler) GetFlamegraph(rw http.ResponseWriter, r *http.Request) {
render.Success(rw, http.StatusOK, result)
}
func (h *handler) GetThread(rw http.ResponseWriter, r *http.Request) {
req := new(spantypes.PostableThreadQuery)
if err := binding.Query.BindQuery(r.URL.Query(), req); err != nil {
render.Error(rw, err)
return
}
query, err := spantypes.NewThreadQuery(req)
if err != nil {
render.Error(rw, err)
return
}
result, err := h.module.GetThread(r.Context(), mux.Vars(r)["traceID"], query)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, result)
}

View File

@@ -173,19 +173,6 @@ func (m *module) getWindowedWaterfall(ctx context.Context, traceID, selectedSpan
), nil
}
func (m *module) GetThread(ctx context.Context, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error) {
summary, err := m.store.GetTraceSummary(ctx, traceID)
if err != nil {
return nil, err
}
spans, err := m.store.GetThreadSpans(ctx, traceID, summary, query.Cursor, query.Limit+1)
if err != nil {
return nil, err
}
return spantypes.NewGettableTraceThread(traceID, spans, query.Limit), nil
}
func (m *module) getFullFlamegraph(ctx context.Context, traceID string, summary *spantypes.TraceSummary, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error) {
fullSpans, err := m.store.GetFlamegraphSpans(ctx, traceID, summary.Start, summary.End, nil)
if err != nil {

View File

@@ -4,7 +4,6 @@ import (
"context"
"database/sql"
"fmt"
"strings"
"time"
sqlbuilder "github.com/huandu/go-sqlbuilder"
@@ -12,23 +11,12 @@ import (
"github.com/SigNoz/signoz/pkg/clickhousesql"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
const colServiceName = `resource_string_service$$$$name` // $ gets escaped so $$$$ converts to $$.
var fullSpanColumns = []string{
"duration_nano", "span_id", "has_error", "kind",
colServiceName, "name",
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
"flags", "is_remote", "trace_state", "status_code",
"db_name", "db_operation", "http_method", "http_url", "http_host",
"external_http_method", "external_http_url", "response_status_code", "links as references",
}
func buildFieldExpr(fieldKey telemetrytypes.TelemetryFieldKey) (string, error) {
switch fieldKey.FieldContext {
case telemetrytypes.FieldContextResource:
@@ -80,11 +68,18 @@ func (s *traceStore) GetTraceSummary(ctx context.Context, traceID string) (*span
func (s *traceStore) GetTraceSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary) ([]spantypes.StorableSpan, error) {
// DISTINCT ON (span_id) is ClickHouse-specific syntax not supported by sqlbuilder
query := fmt.Sprintf(`
SELECT DISTINCT ON (span_id) timestamp, %s
SELECT DISTINCT ON (span_id)
timestamp, duration_nano, span_id, has_error, kind,
resource_string_service$$name, name,
attributes_string, attributes_number, attributes_bool, resources_string,
events, status_message, status_code_string, kind_string, parent_span_id,
flags, is_remote, trace_state, status_code,
db_name, db_operation, http_method, http_url, http_host,
external_http_method, external_http_url, response_status_code, links as references
FROM %s.%s
WHERE trace_id=? AND ts_bucket_start>=? AND ts_bucket_start<=?
ORDER BY timestamp ASC, name ASC`,
strings.Join(fullSpanColumns, ", "), spantypes.TraceDB, spantypes.TraceTable,
spantypes.TraceDB, spantypes.TraceTable,
)
var spanItems []spantypes.StorableSpan
err := s.telemetryStore.ClickhouseDB().Select(
@@ -128,8 +123,16 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return []spantypes.StorableSpan{}, nil
}
sb := sqlbuilder.NewSelectBuilder()
sb.Select("DISTINCT ON (span_id) timestamp")
sb.SelectMore(fullSpanColumns...)
sb.Select(
"DISTINCT ON (span_id) timestamp",
"duration_nano", "span_id", "has_error", "kind",
colServiceName, "name",
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
"flags", "is_remote", "trace_state", "status_code",
"db_name", "db_operation", "http_method", "http_url", "http_host",
"external_http_method", "external_http_url", "response_status_code", "links as references",
)
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
ids := make([]any, len(spanIDs))
for i, id := range spanIDs {
@@ -152,36 +155,6 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return spans, nil
}
func (s *traceStore) GetThreadSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary, cursor *spantypes.ThreadCursor, limit int) ([]spantypes.StorableSpan, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("DISTINCT ON (span_id) timestamp")
sb.SelectMore(fullSpanColumns...)
sb.SelectMore("attributes")
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
sb.Where(
sb.E("trace_id", traceID),
sb.GE("ts_bucket_start", summary.Start.Unix()-1800),
sb.LE("ts_bucket_start", summary.End.Unix()),
sb.Or(
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIInputMessages))),
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIOutputMessages))),
),
)
if cursor != nil {
sb.Where(sb.GT("(toUnixTimestamp64Nano(timestamp), span_id)", sqlbuilder.Tuple(cursor.TimeUnixNano, cursor.SpanID)))
}
sb.OrderByAsc("timestamp")
sb.OrderByAsc("span_id")
sb.Limit(limit)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
var spans []spantypes.StorableSpan
if err := s.telemetryStore.ClickhouseDB().Select(ctx, &spans, query, args...); err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error querying thread spans")
}
return spans, nil
}
func (s *traceStore) GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]spantypes.StorableSpan, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select(

View File

@@ -13,7 +13,6 @@ type Handler interface {
GetWaterfallV4(http.ResponseWriter, *http.Request)
GetTraceAggregations(http.ResponseWriter, *http.Request)
GetFlamegraph(http.ResponseWriter, *http.Request)
GetThread(http.ResponseWriter, *http.Request)
}
// Module defines the business logic for trace detail operations.
@@ -21,5 +20,4 @@ type Module interface {
GetWaterfallV4(ctx context.Context, traceID string, selectedSpanID string, uncollapsedSpans []string) (*spantypes.GettableWaterfallTrace, error)
GetTraceAggregations(ctx context.Context, traceID string, req *spantypes.PostableTraceAggregations) (*spantypes.GettableTraceAggregations, error)
GetFlamegraph(ctx context.Context, traceID string, selectedSpanID string, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error)
GetThread(ctx context.Context, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error)
}

View File

@@ -566,6 +566,24 @@ func readAsRaw(rows driver.Rows, queryName string) (*qbtypes.RawData, error) {
}, nil
}
// flattenJSONPaths flattens a decoded JSON document into dotted keys, overwriting existing keys in out.
func flattenJSONPaths(prefix string, m map[string]any, out map[string]any) {
for k, v := range m {
key := k
if prefix != "" {
key = prefix + "." + k
}
switch child := v.(type) {
case map[string]any:
flattenJSONPaths(key, child, out)
case telemetrystoretypes.JSONValue:
flattenJSONPaths(key, child, out)
default:
out[key] = v
}
}
}
// mergeSpanAttributeColumns merges (attributes_string, attributes_number, attributes_bool, resources_string) into
// unified "attributes" and "resource" keys, and parses the stringified `events`
// and `links` columns into structured slices. Raw DB columns are removed.
@@ -580,7 +598,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
resStr, hasRes := data["resources_string"]
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {
attributes := make(map[string]any)
attrJSON.FlattenInto("", attributes)
flattenJSONPaths("", attrJSON, attributes)
if m, ok := attrStr.(map[string]string); ok {
for k, v := range m {
attributes[k] = v

View File

@@ -36,7 +36,7 @@ func TestManager_ListRules_ValidatesParams(t *testing.T) {
_, err = m.ListRules(context.Background(), &ruletypes.ListRulesParams{Limit: -1})
require.ErrorContains(t, err, "invalid limit")
_, err = m.ListRules(context.Background(), &ruletypes.ListRulesParams{States: []string{"bogus"}})
_, err = m.ListRules(context.Background(), &ruletypes.ListRulesParams{ListFilter: ruletypes.ListFilter{States: []string{"bogus"}}})
require.ErrorContains(t, err, `invalid state "bogus"`)
}

View File

@@ -12,6 +12,11 @@ type Handler interface {
PatchRuleByID(http.ResponseWriter, *http.Request)
TestRule(http.ResponseWriter, *http.Request)
ListRuleViews(http.ResponseWriter, *http.Request)
CreateRuleView(http.ResponseWriter, *http.Request)
UpdateRuleView(http.ResponseWriter, *http.Request)
DeleteRuleView(http.ResponseWriter, *http.Request)
ListDowntimeSchedules(http.ResponseWriter, *http.Request)
GetDowntimeScheduleByID(http.ResponseWriter, *http.Request)
CreateDowntimeSchedule(http.ResponseWriter, *http.Request)

View File

@@ -49,4 +49,9 @@ type Ruler interface {
// TODO: expose downtime CRUD as methods on Ruler directly instead of leaking the
// store interface. The handler should not call store methods directly.
MaintenanceStore() alertmanagertypes.MaintenanceStore
CreateRuleView(ctx context.Context, orgID valuer.UUID, postable ruletypes.PostableRuleView) (*ruletypes.GettableRuleView, error)
ListRuleViews(ctx context.Context, orgID valuer.UUID) (*ruletypes.ListableRuleViews, error)
UpdateRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID, updatable ruletypes.UpdatableRuleView) (*ruletypes.GettableRuleView, error)
DeleteRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) error
}

View File

@@ -0,0 +1,93 @@
package sqlrulestore
import (
"context"
"github.com/SigNoz/signoz/pkg/errors"
ruletypes "github.com/SigNoz/signoz/pkg/types/ruletypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
func (r *rule) CreateRuleView(ctx context.Context, view *ruletypes.StorableRuleView) error {
_, err := r.sqlstore.
BunDBCtx(ctx).
NewInsert().
Model(view).
Exec(ctx)
if err != nil {
return r.sqlstore.WrapAlreadyExistsErrf(err, errors.CodeAlreadyExists, "rule view with id %s already exists", view.ID)
}
return nil
}
func (r *rule) GetRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) (*ruletypes.StorableRuleView, error) {
view := new(ruletypes.StorableRuleView)
err := r.sqlstore.
BunDBCtx(ctx).
NewSelect().
Model(view).
Where("id = ?", id).
Where("org_id = ?", orgID).
Scan(ctx)
if err != nil {
return nil, r.sqlstore.WrapNotFoundErrf(err, ruletypes.ErrCodeRuleViewNotFound, "rule view with id %s doesn't exist", id)
}
return view, nil
}
func (r *rule) ListRuleViews(ctx context.Context, orgID valuer.UUID) ([]*ruletypes.StorableRuleView, error) {
views := make([]*ruletypes.StorableRuleView, 0)
err := r.sqlstore.
BunDBCtx(ctx).
NewSelect().
Model(&views).
Where("org_id = ?", orgID).
OrderExpr("updated_at DESC").
Scan(ctx)
if err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "couldn't list rule views")
}
return views, nil
}
func (r *rule) UpdateRuleView(ctx context.Context, view *ruletypes.StorableRuleView) error {
res, err := r.sqlstore.
BunDBCtx(ctx).
NewUpdate().
Model(view).
WherePK().
Where("org_id = ?", view.OrgID).
Exec(ctx)
if err != nil {
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't update rule view")
}
rows, err := res.RowsAffected()
if err != nil {
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't read rule view update result")
}
if rows == 0 {
return errors.Newf(errors.TypeNotFound, ruletypes.ErrCodeRuleViewNotFound, "rule view with id %s doesn't exist", view.ID)
}
return nil
}
func (r *rule) DeleteRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) error {
res, err := r.sqlstore.
BunDBCtx(ctx).
NewDelete().
Model(new(ruletypes.StorableRuleView)).
Where("id = ?", id).
Where("org_id = ?", orgID).
Exec(ctx)
if err != nil {
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't delete rule view")
}
rows, err := res.RowsAffected()
if err != nil {
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't read rule view delete result")
}
if rows == 0 {
return errors.Newf(errors.TypeNotFound, ruletypes.ErrCodeRuleViewNotFound, "rule view with id %s doesn't exist", id)
}
return nil
}

View File

@@ -345,3 +345,122 @@ func (handler *handler) DeleteDowntimeScheduleByID(rw http.ResponseWriter, req *
render.Success(rw, http.StatusNoContent, nil)
}
func (handler *handler) ListRuleViews(rw http.ResponseWriter, req *http.Request) {
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
defer cancel()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID, err := valuer.NewUUID(claims.OrgID)
if err != nil {
render.Error(rw, err)
return
}
views, err := handler.ruler.ListRuleViews(ctx, orgID)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, views)
}
func (handler *handler) CreateRuleView(rw http.ResponseWriter, req *http.Request) {
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
defer cancel()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID, err := valuer.NewUUID(claims.OrgID)
if err != nil {
render.Error(rw, err)
return
}
var postable ruletypes.PostableRuleView
if err := binding.JSON.BindBody(req.Body, &postable); err != nil {
render.Error(rw, err)
return
}
view, err := handler.ruler.CreateRuleView(ctx, orgID, postable)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusCreated, view)
}
func (handler *handler) UpdateRuleView(rw http.ResponseWriter, req *http.Request) {
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
defer cancel()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID, err := valuer.NewUUID(claims.OrgID)
if err != nil {
render.Error(rw, err)
return
}
id, err := valuer.NewUUID(mux.Vars(req)["id"])
if err != nil {
render.Error(rw, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "id is not a valid uuid-v7"))
return
}
var updatable ruletypes.UpdatableRuleView
if err := binding.JSON.BindBody(req.Body, &updatable); err != nil {
render.Error(rw, err)
return
}
view, err := handler.ruler.UpdateRuleView(ctx, orgID, id, updatable)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, view)
}
func (handler *handler) DeleteRuleView(rw http.ResponseWriter, req *http.Request) {
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
defer cancel()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID, err := valuer.NewUUID(claims.OrgID)
if err != nil {
render.Error(rw, err)
return
}
id, err := valuer.NewUUID(mux.Vars(req)["id"])
if err != nil {
render.Error(rw, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "id is not a valid uuid-v7"))
return
}
if err := handler.ruler.DeleteRuleView(ctx, orgID, id); err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusNoContent, nil)
}

View File

@@ -147,3 +147,41 @@ func (provider *provider) TestNotification(ctx context.Context, orgID valuer.UUI
func (provider *provider) MaintenanceStore() alertmanagertypes.MaintenanceStore {
return provider.manager.MaintenanceStore()
}
func (provider *provider) CreateRuleView(ctx context.Context, orgID valuer.UUID, postable ruletypes.PostableRuleView) (*ruletypes.GettableRuleView, error) {
if err := postable.Validate(); err != nil {
return nil, err
}
storable := postable.ToStorableRuleView(orgID)
if err := provider.ruleStore.CreateRuleView(ctx, storable); err != nil {
return nil, err
}
return storable.ToGettableRuleView(), nil
}
func (provider *provider) ListRuleViews(ctx context.Context, orgID valuer.UUID) (*ruletypes.ListableRuleViews, error) {
storables, err := provider.ruleStore.ListRuleViews(ctx, orgID)
if err != nil {
return nil, err
}
return &ruletypes.ListableRuleViews{Views: ruletypes.NewGettableRuleViewsFromStorableRuleViews(storables)}, nil
}
func (provider *provider) UpdateRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID, updatable ruletypes.UpdatableRuleView) (*ruletypes.GettableRuleView, error) {
if err := updatable.Validate(); err != nil {
return nil, err
}
storable, err := provider.ruleStore.GetRuleView(ctx, orgID, id)
if err != nil {
return nil, err
}
storable.Update(updatable)
if err := provider.ruleStore.UpdateRuleView(ctx, storable); err != nil {
return nil, err
}
return storable.ToGettableRuleView(), nil
}
func (provider *provider) DeleteRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) error {
return provider.ruleStore.DeleteRuleView(ctx, orgID, id)
}

View File

@@ -257,6 +257,7 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddCloudIntegrationTuplesFactory(sqlstore),
sqlmigration.NewAddNotificationChannelTuplesFactory(sqlstore),
sqlmigration.NewAddAIObservabilityQuickFiltersFactory(sqlstore),
sqlmigration.NewAddRuleViewFactory(sqlstore, sqlschema),
)
}

View File

@@ -0,0 +1,78 @@
package sqlmigration
import (
"context"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlschema"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
)
type addRuleView struct {
sqlstore sqlstore.SQLStore
sqlschema sqlschema.SQLSchema
}
func NewAddRuleViewFactory(sqlstore sqlstore.SQLStore, sqlschema sqlschema.SQLSchema) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(factory.MustNewName("add_rule_view"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addRuleView{
sqlstore: sqlstore,
sqlschema: sqlschema,
}, nil
})
}
func (migration *addRuleView) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func (migration *addRuleView) Up(ctx context.Context, db *bun.DB) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
sqls := migration.sqlschema.Operator().CreateTable(&sqlschema.Table{
Name: "rule_view",
Columns: []*sqlschema.Column{
{Name: "id", DataType: sqlschema.DataTypeText, Nullable: false},
{Name: "name", DataType: sqlschema.DataTypeText, Nullable: false},
{Name: "data", DataType: sqlschema.DataTypeText, Nullable: false},
{Name: "org_id", DataType: sqlschema.DataTypeText, Nullable: false},
{Name: "created_at", DataType: sqlschema.DataTypeTimestamp, Nullable: false},
{Name: "updated_at", DataType: sqlschema.DataTypeTimestamp, Nullable: false},
},
PrimaryKeyConstraint: &sqlschema.PrimaryKeyConstraint{ColumnNames: []sqlschema.ColumnName{"id"}},
ForeignKeyConstraints: []*sqlschema.ForeignKeyConstraint{
{
ReferencingColumnName: sqlschema.ColumnName("org_id"),
ReferencedTableName: sqlschema.TableName("organizations"),
ReferencedColumnName: sqlschema.ColumnName("id"),
},
},
})
for _, sql := range sqls {
if _, err := tx.ExecContext(ctx, string(sql)); err != nil {
return err
}
}
if _, err := tx.NewCreateIndex().
Table("rule_view").
Column("org_id").
Index("idx_rule_view_org_id").
IfNotExists().
Exec(ctx); err != nil {
return err
}
return tx.Commit()
}
func (migration *addRuleView) Down(_ context.Context, _ *bun.DB) error {
return nil
}

View File

@@ -1,102 +0,0 @@
package aiobservabilitytypes
import "strings"
const (
MessageRoleSystem MessageRole = "system"
MessageRoleUser MessageRole = "user"
MessageRoleAssistant MessageRole = "assistant"
MessageRoleTool MessageRole = "tool"
)
const (
FinishReasonStop FinishReason = "stop"
FinishReasonToolCall FinishReason = "tool_call"
FinishReasonLength FinishReason = "length"
FinishReasonContentFilter FinishReason = "content_filter"
FinishReasonError FinishReason = "error"
)
const (
PartTypeText PartType = "text"
PartTypeThinking PartType = "thinking"
PartTypeToolCall PartType = "tool_call"
PartTypeToolResult PartType = "tool_result"
PartTypeGeneric PartType = "generic"
)
type MessageRole string
type FinishReason string
type PartType string
// Part is one piece of a message. Which fields are set depends on Type:
//
// text Content
// thinking Content, Redacted
// tool_call ID, Name, Arguments, Server
// tool_result ToolCallID, Name, Content, IsError, Server
// generic Content (the original value, always a string)
type Part struct {
Type PartType `json:"type" required:"true"`
Content string `json:"content,omitempty"`
Redacted bool `json:"redacted,omitempty"`
ID string `json:"id,omitempty"`
Name string `json:"name,omitempty"`
Arguments any `json:"arguments,omitempty"`
Server bool `json:"server,omitempty"`
ToolCallID string `json:"toolCallId,omitempty"`
IsError bool `json:"isError,omitempty"`
}
type Message struct {
Role MessageRole `json:"role,omitempty"`
Content []Part `json:"content" required:"true" nullable:"false"`
FinishReason FinishReason `json:"finishReason,omitempty"`
}
func (PartType) Enum() []any {
return []any{PartTypeText, PartTypeThinking, PartTypeToolCall, PartTypeToolResult, PartTypeGeneric}
}
// normalizeRole keeps an unknown role, lowercased.
func normalizeRole(role string) MessageRole {
if known := knownRole(role); known != "" {
return known
}
return MessageRole(strings.ToLower(strings.TrimSpace(role)))
}
// normalizeFinishReason keeps an unknown reason, lowercased.
func normalizeFinishReason(reason string) FinishReason {
lowered := strings.ToLower(strings.TrimSpace(reason))
switch lowered {
case "stop", "end_turn", "stop_sequence", "completed", "complete", "eos", "finished":
return FinishReasonStop
case "tool_call", "tool_calls", "tool_use", "function_call":
return FinishReasonToolCall
case "length", "max_tokens", "max_output_tokens", "max_completion_tokens", "model_length":
return FinishReasonLength
case "content_filter", "content_filtered", "guardrail_intervened", "safety", "refusal", "recitation", "blocklist", "prohibited_content", "spii":
return FinishReasonContentFilter
case "error", "failed", "incomplete":
return FinishReasonError
}
return FinishReason(lowered)
}
// knownRole maps vendor role names onto MessageRole; anything else is "".
func knownRole(role string) MessageRole {
switch strings.ToLower(strings.TrimSpace(role)) {
case "system", "developer":
return MessageRoleSystem
case "user", "human":
return MessageRoleUser
case "assistant", "ai", "model":
return MessageRoleAssistant
case "tool", "function":
return MessageRoleTool
}
return ""
}

View File

@@ -1,971 +0,0 @@
package aiobservabilitytypes
import (
"encoding/json"
"strings"
)
// Ordered by specificity: earlier converters never match a later format.
var converters = []converter{
convertSemconvMessages,
convertChatMessageList,
convertToolCallList,
convertContentBlockList,
convertChatRequest,
convertChatResponse,
convertResponsesAPIResponse,
convertGeminiResponse,
convertCompletionObject,
convertLangChainGenerations,
convertSingleMessage,
}
var finishReasonKeys = []string{"finish_reason", "finishReason", "stop_reason", "stopReason", "done_reason"}
// NormalizeMessages converts a gen_ai.*.messages value, a JSON string or a
// decoded value; unknown formats become one generic part holding the original.
func NormalizeMessages(raw any) []Message {
var (
value any
original string
)
switch v := raw.(type) {
case nil:
return []Message{}
case string:
original = v
if err := json.Unmarshal([]byte(v), &value); err != nil {
return genericMessages(original)
}
default:
value = v
original = stringOf(v)
}
if list, ok := value.([]any); ok {
if len(list) == 0 {
return []Message{}
}
// [[...]]: some SDKs wrap the conversation in one more list
if _, nested := list[0].([]any); nested {
value = flattenOnce(list)
}
// ["{...}", "{...}"]: an array attribute holding one JSON message per element
if decoded, ok := decodeStringList(list); ok {
value = decoded
}
}
for _, convert := range converters {
if messages, ok := convert(value); ok {
return messages
}
}
return genericMessages(original)
}
type converter func(value any) (messages []Message, ok bool)
func genericMessages(content string) []Message {
return []Message{{Content: []Part{genericPart(content)}}}
}
// convertSemconvMessages handles [{role, parts, finish_reason}] and Gemini contents.
func convertSemconvMessages(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
first, ok := list[0].(map[string]any)
if !ok {
return nil, false
}
if _, ok := first["parts"]; !ok {
return nil, false
}
messages := make([]Message, 0, len(list))
for _, item := range list {
m, ok := item.(map[string]any)
if !ok {
messages = append(messages, genericMessages(stringOf(item))[0])
continue
}
messages = append(messages, partsMessage(m, "")...)
}
return messages, true
}
// partsMessage falls back to chatMessage when m has no parts.
func partsMessage(m map[string]any, defaultRole MessageRole) []Message {
parts, ok := m["parts"].([]any)
if !ok {
return chatMessage(m, defaultRole)
}
role := normalizeRole(stringOf(m["role"]))
if role == "" {
role = defaultRole
}
msg := Message{
Role: role,
Content: []Part{},
FinishReason: normalizeFinishReason(finishReasonOf(m)),
}
for _, p := range parts {
msg.Content = append(msg.Content, semconvPart(p))
}
return []Message{msg}
}
func semconvPart(value any) Part {
p, ok := value.(map[string]any)
if !ok {
if s, ok := value.(string); ok {
return textPart(s)
}
return genericPart(value)
}
switch stringOf(p["type"]) {
case "text":
if boolOf(p["thought"]) {
return Part{Type: PartTypeThinking, Content: stringOf(firstOf(p, "content", "text"))}
}
return Part{Type: PartTypeText, Content: stringOf(firstOf(p, "content", "text"))}
case "reasoning", "thinking":
return Part{Type: PartTypeThinking, Content: stringOf(firstOf(p, "content", "thinking", "text"))}
case "redacted_thinking", "redacted_reasoning":
return Part{Type: PartTypeThinking, Redacted: true}
case "tool_call":
return Part{
Type: PartTypeToolCall,
ID: idOf(p["id"]),
Name: stringOf(p["name"]),
Arguments: parseArguments(firstOf(p, "arguments", "args", "input")),
Server: boolOf(p["server"]),
}
case "tool_call_response":
return Part{
Type: PartTypeToolResult,
ToolCallID: idOf(p["id"]),
Name: stringOf(p["name"]),
Content: stringOf(firstOf(p, "response", "result", "content", "output")),
IsError: boolOf(firstOf(p, "is_error", "isError")),
Server: boolOf(p["server"]),
}
case "":
// Gemini parts carry no type; the field name is the type.
if text, ok := p["text"]; ok {
if boolOf(p["thought"]) {
return Part{Type: PartTypeThinking, Content: stringOf(text)}
}
return Part{Type: PartTypeText, Content: stringOf(text)}
}
if call, ok := firstOf(p, "functionCall", "function_call").(map[string]any); ok {
return Part{Type: PartTypeToolCall, ID: stringOf(call["id"]), Name: stringOf(call["name"]), Arguments: parseArguments(call["args"])}
}
if resp, ok := firstOf(p, "functionResponse", "function_response").(map[string]any); ok {
return Part{Type: PartTypeToolResult, ToolCallID: stringOf(resp["id"]), Name: stringOf(resp["name"]), Content: stringOf(resp["response"])}
}
}
return genericPart(p)
}
// convertChatMessageList handles OpenAI, Anthropic, Bedrock, Vercel and LangChain message lists.
func convertChatMessageList(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
if !isChatMessage(list[0]) {
return nil, false
}
messages := make([]Message, 0, len(list))
for _, item := range list {
messages = append(messages, chatMessage(item, "")...)
}
return messages, true
}
func isChatMessage(value any) bool {
m, ok := value.(map[string]any)
if !ok {
return false
}
if _, ok := m["role"]; ok {
return true
}
if _, ok := m["gen_ai.event.content"]; ok {
return true
}
switch typ := stringOf(m["type"]); typ {
case "message", "reasoning", "human", "ai", "tool", "system":
return true
case "constructor":
_, ok := m["kwargs"]
return ok
default:
return isResponsesItemType(typ)
}
}
// chatMessage returns nil for LangGraph tool definitions.
func chatMessage(value any, defaultRole MessageRole) []Message {
m, ok := value.(map[string]any)
if !ok {
return genericMessages(stringOf(value))
}
if inner, role, ok := unwrapChatEnvelope(m, defaultRole); ok {
return chatMessage(inner, role)
}
if isLangGraphToolDefinition(m) {
return nil
}
typ := stringOf(m["type"])
if isResponsesItemType(typ) {
return responsesItem(m, typ)
}
if typ == "reasoning" {
return []Message{{Role: MessageRoleAssistant, Content: reasoningParts(m)}}
}
role := normalizeRole(stringOf(m["role"]))
if role == "" {
role = knownRole(typ)
}
if role == "" {
role = defaultRole
}
msg := Message{
Role: role,
Content: append(chatContentParts(m, role), chatToolCallParts(m)...),
FinishReason: normalizeFinishReason(finishReasonOf(m)),
}
if refusal := stringOf(m["refusal"]); refusal != "" {
msg.Content = append(msg.Content, textPart(refusal))
}
return []Message{msg}
}
// unwrapChatEnvelope unwraps LangChain serialised messages and Semantic Kernel events.
func unwrapChatEnvelope(m map[string]any, defaultRole MessageRole) (map[string]any, MessageRole, bool) {
if kwargs, ok := m["kwargs"].(map[string]any); ok && stringOf(m["type"]) == "constructor" {
role := langChainRole(m["id"])
if role == "" {
role = knownRole(stringOf(kwargs["type"]))
}
return kwargs, role, true
}
event, ok := m["gen_ai.event.content"].(string)
if !ok {
return nil, "", false
}
var inner map[string]any
if err := json.Unmarshal([]byte(event), &inner); err != nil || inner == nil {
return nil, "", false
}
if message, ok := inner["message"].(map[string]any); ok {
if _, has := message["finish_reason"]; !has {
message["finish_reason"] = inner["finish_reason"]
}
inner = message
}
return inner, defaultRole, true
}
// chatContentParts turns a tool message's text into its result.
func chatContentParts(m map[string]any, role MessageRole) []Part {
parts := []Part{}
switch content := m["content"].(type) {
case nil:
case string:
if role == MessageRoleTool {
parts = append(parts, toolResultOf(m, content))
} else if content != "" {
parts = append(parts, textPart(content))
}
case []any:
for _, item := range content {
part := chatContentPart(item)
if role == MessageRoleTool && part.Type == PartTypeText {
part = toolResultOf(m, part.Content)
}
parts = append(parts, part)
}
case map[string]any:
if contentParts, ok := content["parts"].([]any); ok {
for _, p := range contentParts {
parts = append(parts, semconvPart(p))
}
} else if role == MessageRoleTool {
parts = append(parts, toolResultOf(m, stringOf(content)))
} else {
parts = append(parts, genericPart(content))
}
default:
parts = append(parts, genericPart(content))
}
return parts
}
func chatToolCallParts(m map[string]any) []Part {
calls, _ := firstOf(m, "tool_calls", "toolCalls").([]any)
if kwargs, ok := m["additional_kwargs"].(map[string]any); ok && len(calls) == 0 {
calls, _ = kwargs["tool_calls"].([]any)
}
parts := make([]Part, 0, len(calls)+1)
for _, call := range calls {
parts = append(parts, toolCallPart(call))
}
if call, ok := m["function_call"].(map[string]any); ok {
parts = append(parts, Part{Type: PartTypeToolCall, Name: stringOf(call["name"]), Arguments: parseArguments(call["arguments"])})
}
return parts
}
func toolResultOf(m map[string]any, content string) Part {
return Part{Type: PartTypeToolResult, ToolCallID: stringOf(firstOf(m, "tool_call_id", "toolCallId")), Name: stringOf(m["name"]), Content: content}
}
// isLangGraphToolDefinition matches {role: "tool", content: {type: "function"}} without tool_call_id.
func isLangGraphToolDefinition(m map[string]any) bool {
if normalizeRole(stringOf(m["role"])) != MessageRoleTool {
return false
}
if _, has := m["tool_call_id"]; has {
return false
}
content, ok := m["content"].(map[string]any)
if !ok || stringOf(content["type"]) != "function" {
return false
}
_, ok = content["function"]
return ok
}
func chatContentPart(value any) Part {
p, ok := value.(map[string]any)
if !ok {
if s, ok := value.(string); ok {
return textPart(s)
}
return genericPart(value)
}
typ := stringOf(p["type"])
switch typ {
case "text", "input_text", "output_text", "refusal", "summary_text":
return Part{Type: PartTypeText, Content: stringOf(firstOf(p, "text", "content", "refusal"))}
case "thinking", "reasoning":
return Part{Type: PartTypeThinking, Content: stringOf(firstOf(p, "thinking", "text", "content", "reasoning"))}
case "redacted_thinking":
return Part{Type: PartTypeThinking, Redacted: true}
case "tool-call", "tool_use", "tool_call", "function_call":
return toolCallPart(p)
case "tool-result", "tool_result", "function_call_output":
return toolResultPart(p, firstOf(p, "result", "output", "content"), false)
case "server_tool_use", "mcp_tool_use":
part := toolCallPart(p)
part.Server = true
return part
case "":
if text, ok := p["text"]; ok {
return Part{Type: PartTypeText, Content: stringOf(text)}
}
// Bedrock Converse blocks are typed by field name.
if use, ok := p["toolUse"].(map[string]any); ok {
return Part{Type: PartTypeToolCall, ID: stringOf(use["toolUseId"]), Name: stringOf(use["name"]), Arguments: parseArguments(use["input"])}
}
if result, ok := p["toolResult"].(map[string]any); ok {
part := toolResultPart(result, result["content"], false)
part.ToolCallID = stringOf(result["toolUseId"])
part.IsError = stringOf(result["status"]) == "error"
return part
}
default:
if strings.HasSuffix(typ, "_tool_result") {
return toolResultPart(p, p["content"], true)
}
}
return genericPart(p)
}
// toolCallPart reads the OpenAI, flat, Anthropic and Vercel tool call shapes.
func toolCallPart(value any) Part {
call, ok := value.(map[string]any)
if !ok {
return genericPart(value)
}
part := Part{
Type: PartTypeToolCall,
ID: idOf(firstOf(call, "toolCallId", "call_id", "id")),
Name: stringOf(firstOf(call, "toolName", "name")),
}
if fn, ok := call["function"].(map[string]any); ok {
part.Name = stringOf(fn["name"])
part.Arguments = parseArguments(fn["arguments"])
return part
}
part.Arguments = parseArguments(firstOf(call, "arguments", "args", "input"))
return part
}
// toolResultPart unwraps the Vercel {type, value} result wrapper.
func toolResultPart(p map[string]any, result any, server bool) Part {
if nested, ok := result.(map[string]any); ok && len(nested) <= 2 {
if v, ok := nested["value"]; ok {
result = v
}
}
return Part{
Type: PartTypeToolResult,
ToolCallID: idOf(firstOf(p, "toolCallId", "tool_use_id", "tool_call_id", "call_id", "id")),
Name: stringOf(firstOf(p, "toolName", "name")),
Content: textOf(result),
IsError: boolOf(firstOf(p, "isError", "is_error")),
Server: server,
}
}
// textOf joins a list of text blocks; anything else goes through stringOf.
func textOf(result any) string {
list, ok := result.([]any)
if !ok || len(list) == 0 {
return stringOf(result)
}
texts := make([]string, 0, len(list))
for _, item := range list {
block := asMap(item)
text, ok := block["text"].(string)
if !ok || (len(block) == 2 && stringOf(block["type"]) != "text") || len(block) > 2 {
return stringOf(result)
}
texts = append(texts, text)
}
return strings.Join(texts, "\n")
}
// isResponsesItemType matches role-less Responses API tool and MCP items.
func isResponsesItemType(typ string) bool {
switch typ {
case "":
return false
case "function_call", "function_call_output", "tool_call", "custom_tool_call", "custom_tool_call_output",
"mcp_call", "mcp_list_tools", "mcp_approval_request", "mcp_approval_response":
return true
}
return strings.HasSuffix(typ, "_call") || strings.HasSuffix(typ, "_call_output")
}
// responsesItem maps built-in tools to server tool parts.
func responsesItem(m map[string]any, typ string) []Message {
switch typ {
case "function_call", "tool_call", "custom_tool_call":
return []Message{{Role: MessageRoleAssistant, Content: []Part{{
Type: PartTypeToolCall,
ID: stringOf(firstOf(m, "call_id", "id")),
Name: stringOf(m["name"]),
Arguments: parseArguments(firstOf(m, "arguments", "args", "input")),
}}}}
case "function_call_output", "custom_tool_call_output":
return []Message{{Role: MessageRoleTool, Content: []Part{{
Type: PartTypeToolResult,
ToolCallID: stringOf(firstOf(m, "call_id", "id")),
Content: stringOf(firstOf(m, "output", "result")),
}}}}
}
id := stringOf(firstOf(m, "call_id", "id"))
if strings.HasSuffix(typ, "_output") || typ == "mcp_approval_response" {
return []Message{{Role: MessageRoleTool, Content: []Part{{
Type: PartTypeToolResult,
ToolCallID: id,
Name: strings.TrimSuffix(typ, "_output"),
Content: stringOf(firstOf(m, "output", "result", "results")),
Server: true,
}}}}
}
args := make(map[string]any, len(m))
for k, v := range m {
switch k {
case "type", "id", "call_id", "status", "name", "server_label", "output", "result", "results":
default:
args[k] = v
}
}
name := stringOf(firstOf(m, "name", "server_label"))
if name == "" {
name = typ
}
call := Part{Type: PartTypeToolCall, ID: id, Name: name, Server: true}
if len(args) > 0 {
call.Arguments = args
}
msg := Message{Role: MessageRoleAssistant, Content: []Part{call}}
if result := firstOf(m, "output", "result", "results"); result != nil {
msg.Content = append(msg.Content, Part{Type: PartTypeToolResult, ToolCallID: id, Name: name, Content: stringOf(result), Server: true})
}
return []Message{msg}
}
// reasoningParts marks encrypted reasoning without a summary as redacted.
func reasoningParts(m map[string]any) []Part {
parts := []Part{}
if summary, ok := m["summary"].([]any); ok {
for _, s := range summary {
parts = append(parts, Part{Type: PartTypeThinking, Content: stringOf(firstOf(asMap(s), "text", "content"))})
}
}
if content, ok := m["content"].([]any); ok {
for _, c := range content {
parts = append(parts, Part{Type: PartTypeThinking, Content: stringOf(firstOf(asMap(c), "text", "content"))})
}
}
if len(parts) == 0 {
parts = append(parts, Part{Type: PartTypeThinking, Redacted: true})
}
return parts
}
// convertToolCallList handles a bare tool call list, e.g. Vercel ai.response.toolCalls.
func convertToolCallList(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
msg := Message{Role: MessageRoleAssistant, Content: make([]Part, 0, len(list))}
for _, item := range list {
call, ok := item.(map[string]any)
if !ok || !isToolCall(call) {
return nil, false
}
msg.Content = append(msg.Content, toolCallPart(call))
}
return []Message{msg}, true
}
// isToolCall rejects tool definitions, which carry no arguments.
func isToolCall(call map[string]any) bool {
if _, has := call["toolName"]; has {
return true
}
if fn, ok := call["function"].(map[string]any); ok {
_, has := fn["arguments"]
return has
}
if _, has := call["name"]; !has {
return false
}
_, has := lookup(call, "arguments", "args")
return has
}
// convertContentBlockList handles a bare content block list; the role is unknown.
func convertContentBlockList(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
msg := Message{Content: make([]Part, 0, len(list))}
for _, item := range list {
block, ok := item.(map[string]any)
if !ok {
return nil, false
}
if _, has := block["type"].(string); !has {
return nil, false
}
part := chatContentPart(block)
if part.Type == PartTypeGeneric {
return nil, false
}
msg.Content = append(msg.Content, part)
}
return []Message{msg}, true
}
// convertChatRequest handles OpenAI, Anthropic, Vercel, Gemini and LangChain request objects.
func convertChatRequest(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
conversation, ok := lookup(m, "messages", "input", "contents", "prompt")
if !ok || !isConversation(m, conversation) {
return nil, false
}
messages := []Message{}
if system := systemMessage(firstOf(m, "system", "instructions", "system_instruction", "systemInstruction", "system_prompt")); system != nil {
messages = append(messages, *system)
} else if config, ok := m["config"].(map[string]any); ok {
if system := systemMessage(firstOf(config, "system_instruction", "systemInstruction")); system != nil {
messages = append(messages, *system)
}
}
// {messages: "[...]"}: the list arrives JSON-encoded once more from some SDKs
if s, isString := conversation.(string); isString {
var decoded any
if err := json.Unmarshal([]byte(s), &decoded); err == nil {
if _, isList := decoded.([]any); isList {
conversation = decoded
}
}
}
switch c := conversation.(type) {
case string:
messages = append(messages, textMessage(MessageRoleUser, c))
case []any:
for _, item := range flattenOnce(c) {
if s, isString := item.(string); isString {
messages = append(messages, textMessage(MessageRoleUser, s))
continue
}
messages = append(messages, partsMessage(asMap(item), MessageRoleUser)...)
}
case map[string]any:
messages = append(messages, partsMessage(c, MessageRoleUser)...)
default:
return nil, false
}
return messages, true
}
// isConversation rejects embeddings requests.
func isConversation(m map[string]any, conversation any) bool {
if _, isRequestInput := m["input"]; !isRequestInput {
return true
}
if _, hasChatKey := lookup(m, "instructions", "tools", "tool_choice", "parallel_tool_calls", "previous_response_id"); hasChatKey {
return true
}
list, ok := conversation.([]any)
if !ok {
return false
}
for _, item := range list {
if _, isMap := item.(map[string]any); !isMap {
return false
}
}
return true
}
// systemMessage returns nil when value is empty or unknown.
func systemMessage(value any) *Message {
msg := Message{Role: MessageRoleSystem, Content: []Part{}}
switch v := value.(type) {
case string:
if v == "" {
return nil
}
msg.Content = append(msg.Content, textPart(v))
case []any:
for _, item := range v {
msg.Content = append(msg.Content, chatContentPart(item))
}
case map[string]any:
parts, ok := v["parts"].([]any)
if !ok {
return nil
}
for _, p := range parts {
msg.Content = append(msg.Content, semconvPart(p))
}
default:
return nil
}
if len(msg.Content) == 0 {
return nil
}
return &msg
}
// convertChatResponse handles {choices}, {message} and {output: {message}} responses.
func convertChatResponse(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
wrapped, _ := m["message"].(map[string]any)
if output, ok := m["output"].(map[string]any); ok && wrapped == nil {
wrapped, _ = output["message"].(map[string]any)
}
if wrapped != nil {
if !isChatMessage(wrapped) {
return nil, false
}
return withFinishReason(chatMessage(wrapped, MessageRoleAssistant), finishReasonOf(m)), true
}
choices, ok := m["choices"].([]any)
if !ok {
return nil, false
}
messages := make([]Message, 0, len(choices))
for _, c := range choices {
choice := asMap(c)
var converted []Message
if message, ok := firstOf(choice, "message", "delta").(map[string]any); ok {
converted = chatMessage(message, MessageRoleAssistant)
} else if text, ok := choice["text"]; ok {
converted = []Message{textMessage(MessageRoleAssistant, stringOf(text))}
} else {
converted = []Message{{Role: MessageRoleAssistant, Content: []Part{genericPart(choice)}}}
}
messages = append(messages, withFinishReason(converted, stringOf(choice["finish_reason"]))...)
}
return messages, true
}
func convertResponsesAPIResponse(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
output, ok := m["output"].([]any)
if !ok {
return nil, false
}
messages := make([]Message, 0, len(output))
for _, item := range output {
messages = append(messages, chatMessage(item, MessageRoleAssistant)...)
}
if len(messages) == 0 {
return messages, true
}
last := &messages[len(messages)-1]
if last.FinishReason == "" {
if details, ok := m["incomplete_details"].(map[string]any); ok {
last.FinishReason = normalizeFinishReason(stringOf(details["reason"]))
} else if stringOf(m["status"]) == "completed" {
last.FinishReason = FinishReasonStop
}
}
return messages, true
}
// convertGeminiResponse handles Gemini {candidates} and Google ADK {content} responses.
func convertGeminiResponse(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
if content, ok := m["content"].(map[string]any); ok {
if _, hasParts := content["parts"]; hasParts {
return withFinishReason(partsMessage(content, MessageRoleAssistant), finishReasonOf(m)), true
}
}
candidates, ok := m["candidates"].([]any)
if !ok {
return nil, false
}
messages := make([]Message, 0, len(candidates))
for _, c := range candidates {
candidate := asMap(c)
content, ok := candidate["content"].(map[string]any)
if !ok {
messages = append(messages, Message{Role: MessageRoleAssistant, Content: []Part{genericPart(candidate)}})
continue
}
messages = append(messages, withFinishReason(partsMessage(content, MessageRoleAssistant), finishReasonOf(candidate))...)
}
return messages, true
}
func convertCompletionObject(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
completion, ok := m["completion"].(string)
if !ok {
return nil, false
}
msg := Message{Role: MessageRoleAssistant, Content: []Part{}}
if reasoning, ok := m["reasoning"].(string); ok && reasoning != "" {
msg.Content = append(msg.Content, Part{Type: PartTypeThinking, Content: reasoning})
}
msg.Content = append(msg.Content, textPart(completion))
return []Message{msg}, true
}
func convertLangChainGenerations(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
generations, ok := m["generations"].([]any)
if !ok {
return nil, false
}
messages := []Message{}
for _, g := range flattenOnce(generations) {
gen := asMap(g)
var converted []Message
if message, ok := gen["message"].(map[string]any); ok {
converted = chatMessage(message, MessageRoleAssistant)
} else {
converted = []Message{textMessage(MessageRoleAssistant, stringOf(gen["text"]))}
}
messages = append(messages, withFinishReason(converted, stringOf(asMap(gen["generation_info"])["finish_reason"]))...)
}
return messages, true
}
func convertSingleMessage(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
if _, ok := m["parts"]; ok {
return partsMessage(m, ""), true
}
if isChatMessage(m) {
return chatMessage(m, ""), true
}
return nil, false
}
// withFinishReason sets reason on the last message that has none.
func withFinishReason(messages []Message, reason string) []Message {
if len(messages) == 0 {
return messages
}
last := &messages[len(messages)-1]
if last.FinishReason == "" {
last.FinishReason = normalizeFinishReason(reason)
}
return messages
}
func langChainRole(id any) MessageRole {
path, ok := id.([]any)
if !ok || len(path) == 0 {
return ""
}
switch class := stringOf(path[len(path)-1]); {
case strings.HasPrefix(class, "System"):
return MessageRoleSystem
case strings.HasPrefix(class, "Human"):
return MessageRoleUser
case strings.HasPrefix(class, "AI"):
return MessageRoleAssistant
case strings.HasPrefix(class, "Tool"), strings.HasPrefix(class, "Function"):
return MessageRoleTool
}
return ""
}
// parseArguments decodes JSON-encoded arguments; anything else is returned as is.
func parseArguments(value any) any {
s, ok := value.(string)
if !ok {
return value
}
var decoded any
if err := json.Unmarshal([]byte(s), &decoded); err != nil {
return s
}
return decoded
}
// idOf picks the call_ entry, else the last, from list ids such as ["run_id", "call_id"].
func idOf(value any) string {
list, ok := value.([]any)
if !ok {
return stringOf(value)
}
if len(list) == 0 {
return ""
}
for _, item := range list {
if s, ok := item.(string); ok && strings.HasPrefix(s, "call_") {
return s
}
}
return stringOf(list[len(list)-1])
}
func decodeStringList(list []any) ([]any, bool) {
out := make([]any, 0, len(list))
for _, item := range list {
s, ok := item.(string)
if !ok {
return nil, false
}
var decoded map[string]any
if err := json.Unmarshal([]byte(s), &decoded); err != nil || decoded == nil {
return nil, false
}
out = append(out, decoded)
}
return out, true
}
func flattenOnce(list []any) []any {
out := make([]any, 0, len(list))
for _, item := range list {
if inner, ok := item.([]any); ok {
out = append(out, inner...)
continue
}
out = append(out, item)
}
return out
}
func lookup(m map[string]any, keys ...string) (any, bool) {
for _, k := range keys {
if v, ok := m[k]; ok && v != nil {
return v, true
}
}
return nil, false
}
func firstOf(m map[string]any, keys ...string) any {
v, _ := lookup(m, keys...)
return v
}
func asMap(value any) map[string]any {
m, _ := value.(map[string]any)
return m
}
func boolOf(value any) bool {
b, _ := value.(bool)
return b
}
// stringOf renders nil as "" and non-strings as compact JSON.
func stringOf(value any) string {
switch v := value.(type) {
case nil:
return ""
case string:
return v
}
data, err := json.Marshal(value)
if err != nil {
return ""
}
return string(data)
}
func finishReasonOf(m map[string]any) string {
return stringOf(firstOf(m, finishReasonKeys...))
}
func textPart(content string) Part {
return Part{Type: PartTypeText, Content: content}
}
func genericPart(value any) Part {
return Part{Type: PartTypeGeneric, Content: stringOf(value)}
}
func textMessage(role MessageRole, content string) Message {
return Message{Role: role, Content: []Part{textPart(content)}}
}

View File

@@ -1,387 +0,0 @@
package aiobservabilitytypes
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestNormalizeMessages(t *testing.T) {
text := func(role MessageRole, content string) Message {
return Message{Role: role, Content: []Part{{Type: PartTypeText, Content: content}}}
}
testCases := []struct {
name string
raw any
want []Message
}{
{
name: "SemconvInput_Litellm",
raw: `[{"role": "system", "parts": [{"type": "text", "content": "You are a concise assistant."}]}, {"role": "user", "parts": [{"type": "text", "content": "Give me a one-line definition of observability."}]}]`,
want: []Message{
text(MessageRoleSystem, "You are a concise assistant."),
text(MessageRoleUser, "Give me a one-line definition of observability."),
},
},
{
name: "SemconvOutput_FinishReason_Bifrost",
raw: `[{"role": "assistant", "parts": [{"content": "Observability is X.", "type": "text"}], "finish_reason": "stop"}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Observability is X."}}, FinishReason: FinishReasonStop}},
},
{
name: "SemconvToolCallAndResponse_Langchain",
raw: `[{"role": "user", "parts": [{"type": "text", "content": "What's the weather in Bengaluru?"}]}, {"role": "assistant", "parts": [{"type": "tool_call", "id": "call_1", "name": "get_weather", "arguments": {"city": "Bengaluru"}}]}, {"role": "tool", "parts": [{"type": "tool_call_response", "id": "call_1", "response": "{\"city\": \"Bengaluru\", \"temp_c\": 18}"}]}]`,
want: []Message{
text(MessageRoleUser, "What's the weather in Bengaluru?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_1", Name: "get_weather", Arguments: map[string]any{"city": "Bengaluru"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_1", Content: `{"city": "Bengaluru", "temp_c": 18}`}}},
},
},
{
name: "SemconvOutput_TwoToolCalls_Openllmetry",
raw: `[{"role": "assistant", "parts": [{"type": "tool_call", "name": "search_web", "id": "call_a", "arguments": {"query": "SigNoz"}}, {"type": "tool_call", "name": "get_weather", "id": "call_b", "arguments": {"city": "Bengaluru"}}], "finish_reason": "tool_call"}]`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonToolCall, Content: []Part{
{Type: PartTypeToolCall, ID: "call_a", Name: "search_web", Arguments: map[string]any{"query": "SigNoz"}},
{Type: PartTypeToolCall, ID: "call_b", Name: "get_weather", Arguments: map[string]any{"city": "Bengaluru"}},
}}},
},
{
name: "OpenAIChatList_FlattenedToolCalls_BifrostGateway",
raw: `[{"role":"user","content":"What's the weather in Bengaluru? Use the tool."},{"role":"assistant","content":"","tool_calls":[{"id":"call_y","type":"function","name":"get_current_weather","args":"{\"city\":\"Bengaluru\"}"}]},{"role":"tool","content":"{\"city\": \"Bengaluru\", \"temp_c\": 28}"}]`,
want: []Message{
text(MessageRoleUser, "What's the weather in Bengaluru? Use the tool."),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_y", Name: "get_current_weather", Arguments: map[string]any{"city": "Bengaluru"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, Content: `{"city": "Bengaluru", "temp_c": 28}`}}},
},
},
{
name: "OpenAIChatRequest_NestedFunctionToolCalls_OpenrouterGateway",
raw: `{"messages":[{"role":"user","content":"Weather?"},{"content":null,"refusal":null,"role":"assistant","tool_calls":[{"id":"call_A","function":{"arguments":"{\"city\":\"Bengaluru\"}","name":"get_current_weather"},"type":"function","index":0}]},{"role":"tool","tool_call_id":"call_A","content":"{\"temp_c\": 28}"}]}`,
want: []Message{
text(MessageRoleUser, "Weather?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_A", Name: "get_current_weather", Arguments: map[string]any{"city": "Bengaluru"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_A", Content: `{"temp_c": 28}`}}},
},
},
{
name: "OpenAIChatRequest_IgnoresModelAndTools_Openinference",
raw: `{"messages": [{"role": "user", "content": "What is the weather in Bengaluru in celsius?"}], "model": "gpt-4o-mini", "tool_choice": "auto", "tools": [{"type": "function", "function": {"name": "get_current_weather"}}]}`,
want: []Message{text(MessageRoleUser, "What is the weather in Bengaluru in celsius?")},
},
{
name: "OpenAIChatResponse_Openinference",
raw: `{"id":"chatcmpl-1","choices":[{"finish_reason":"stop","index":0,"logprobs":null,"message":{"content":"Hello! How are you today?","refusal":null,"role":"assistant","annotations":[]}}],"model":"gpt-4o-mini","object":"chat.completion","usage":{"total_tokens":28}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Hello! How are you today?"}}, FinishReason: FinishReasonStop}},
},
{
name: "OpenAIResponsesRequest_OpenAIAgents",
raw: `{"include": [], "input": [{"content": "What's the weather in Bangalore right now?", "role": "user"}], "instructions": "You are a concise weather assistant.", "model": "gpt-4o-mini", "tools": [{"name": "get_weather", "type": "function"}]}`,
want: []Message{
text(MessageRoleSystem, "You are a concise weather assistant."),
text(MessageRoleUser, "What's the weather in Bangalore right now?"),
},
},
{
name: "OpenAIResponsesResponse_FunctionCall_OpenAIAgents",
raw: `{"id":"resp_1","object":"response","status":"completed","output":[{"arguments":"{\"city\":\"Bangalore\"}","call_id":"call_M","name":"get_weather","type":"function_call","id":"fc_1","status":"completed"}],"usage":{"total_tokens":113}}`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonStop, Content: []Part{{Type: PartTypeToolCall, ID: "call_M", Name: "get_weather", Arguments: map[string]any{"city": "Bangalore"}}}}},
},
{
name: "OpenAIResponsesResponse_MessageAndReasoning",
raw: `{"object":"response","status":"completed","output":[{"type":"reasoning","id":"rs_1","summary":[{"type":"summary_text","text":"Thinking about it"}]},{"type":"message","role":"assistant","content":[{"type":"output_text","text":"Paris."}]}]}`,
want: []Message{
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeThinking, Content: "Thinking about it"}}},
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Paris."}}, FinishReason: FinishReasonStop},
},
},
{
name: "VercelPromptMessages_ToolParts",
raw: `[{"role":"user","content":[{"type":"text","text":"What is the weather in Bengaluru in celsius?"}]},{"role":"assistant","content":[{"type":"tool-call","toolCallId":"call_l","toolName":"get_current_weather","args":{"city":"Bengaluru","unit":"c"}}]},{"role":"tool","content":[{"type":"tool-result","toolCallId":"call_l","toolName":"get_current_weather","result":{"city":"Bengaluru","temperature":27}}]}]`,
want: []Message{
text(MessageRoleUser, "What is the weather in Bengaluru in celsius?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_l", Name: "get_current_weather", Arguments: map[string]any{"city": "Bengaluru", "unit": "c"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_l", Name: "get_current_weather", Content: `{"city":"Bengaluru","temperature":27}`}}},
},
},
{
name: "MastraPromptMessages_MixedContent",
raw: `[{"role":"system","content":"You are a weather assistant."},{"role":"user","content":[{"type":"text","text":"What is the weather in Bengaluru?"}]}]`,
want: []Message{
text(MessageRoleSystem, "You are a weather assistant."),
text(MessageRoleUser, "What is the weather in Bengaluru?"),
},
},
{
name: "AnthropicMessages_ThinkingAndToolUse",
raw: `[{"role":"user","content":"Hi"},{"role":"assistant","content":[{"type":"thinking","thinking":"Let me see"},{"type":"redacted_thinking","data":"x"},{"type":"tool_use","id":"toolu_1","name":"lookup","input":{"q":"a"}}],"stop_reason":"tool_use"},{"role":"user","content":[{"type":"tool_result","tool_use_id":"toolu_1","content":"found","is_error":true}]}]`,
want: []Message{
text(MessageRoleUser, "Hi"),
{Role: MessageRoleAssistant, FinishReason: FinishReasonToolCall, Content: []Part{
{Type: PartTypeThinking, Content: "Let me see"},
{Type: PartTypeThinking, Redacted: true},
{Type: PartTypeToolCall, ID: "toolu_1", Name: "lookup", Arguments: map[string]any{"q": "a"}},
}},
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "toolu_1", Content: "found", IsError: true}}},
},
},
{
name: "GeminiContents_Converted",
raw: `[{"role":"user","parts":[{"text":"Weather in Paris?"}]},{"role":"model","parts":[{"functionCall":{"name":"get_weather","args":{"city":"Paris"}}}]},{"role":"user","parts":[{"functionResponse":{"name":"get_weather","response":{"temp":20}}}]}]`,
want: []Message{
text(MessageRoleUser, "Weather in Paris?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}},
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, Name: "get_weather", Content: `{"temp":20}`}}},
},
},
{
name: "LangChainSerialisedPrompt_Langsmith",
raw: `{"messages":[[{"lc":1,"type":"constructor","id":["langchain","schema","messages","SystemMessage"],"kwargs":{"content":"You are concise.","type":"system"}},{"lc":1,"type":"constructor","id":["langchain","schema","messages","HumanMessage"],"kwargs":{"content":"Define observability.","type":"human"}}]]}`,
want: []Message{
text(MessageRoleSystem, "You are concise."),
text(MessageRoleUser, "Define observability."),
},
},
{
name: "LangChainGenerations_Langsmith",
raw: `{"generations":[[{"text":"Observability is Y.","generation_info":{"finish_reason":"stop","logprobs":null},"type":"ChatGeneration","message":{"lc":1,"type":"constructor","id":["langchain","schema","messages","AIMessage"],"kwargs":{"content":"Observability is Y.","type":"ai","tool_calls":[],"invalid_tool_calls":[]}}}]],"llm_output":{"model_name":"gpt-4o-mini"}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Observability is Y."}}, FinishReason: FinishReasonStop}},
},
{
name: "SemconvTextPartWithTextKey_Litellm",
raw: `[{"role": "user", "parts": [{"type": "text", "text": "What animal is in this image?"}, {"type": "image_url", "image_url": {"url": "data:image/jpeg;base64,AAAA"}}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{
{Type: PartTypeText, Content: "What animal is in this image?"},
{Type: PartTypeGeneric, Content: `{"image_url":{"url":"data:image/jpeg;base64,AAAA"},"type":"image_url"}`},
}}},
},
{
name: "CompletionObject_OpenrouterGateway",
raw: `{"completion":"Paris.","reasoning":"The user asks for a capital.","rawRequest":{"model":"openai/gpt-4o-mini"}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{
{Type: PartTypeThinking, Content: "The user asks for a capital."},
{Type: PartTypeText, Content: "Paris."},
}}},
},
{
name: "VercelPrompt_SystemAndPrompt",
raw: `{"system":"You are concise.","prompt":"Say hello in five words."}`,
want: []Message{text(MessageRoleSystem, "You are concise."), text(MessageRoleUser, "Say hello in five words.")},
},
{
name: "VercelResponseToolCalls_BareList",
raw: `[{"toolCallType":"function","toolCallId":"call_1","toolName":"getWeather","args":"{\"city\":\"Bengaluru\"}"},{"type":"tool-call","toolCallId":"call_2","toolName":"searchWeb","input":{"q":"SigNoz"}}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{
{Type: PartTypeToolCall, ID: "call_1", Name: "getWeather", Arguments: map[string]any{"city": "Bengaluru"}},
{Type: PartTypeToolCall, ID: "call_2", Name: "searchWeb", Arguments: map[string]any{"q": "SigNoz"}},
}}},
},
{
name: "LangChainTypeMessages_AdditionalKwargsToolCalls",
raw: `[{"type":"human","content":"Weather?"},{"type":"ai","content":"","additional_kwargs":{"tool_calls":[{"id":"call_1","type":"function","function":{"name":"get_weather","arguments":"{\"city\":\"Paris\"}"}}]}},{"type":"tool","content":"20C","tool_call_id":"call_1","name":"get_weather"}]`,
want: []Message{
text(MessageRoleUser, "Weather?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_1", Name: "get_weather", Content: "20C"}}},
},
},
{
name: "LangGraphToolDefinitionMessage_Skipped",
raw: `[{"role":"tool","content":{"type":"function","function":{"name":"get_weather","parameters":{}}}},{"role":"user","content":"Hi"}]`,
want: []Message{text(MessageRoleUser, "Hi")},
},
{
name: "GeminiResponse_CandidatesWithFinishReason",
raw: `{"candidates":[{"content":{"parts":[{"text":"Let me check","thought":true},{"function_call":{"name":"get_weather","args":{"city":"Paris"}}}],"role":"model"},"finishReason":"STOP"}],"usageMetadata":{}}`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonStop, Content: []Part{
{Type: PartTypeThinking, Content: "Let me check"},
{Type: PartTypeToolCall, Name: "get_weather", Arguments: map[string]any{"city": "Paris"}},
}}},
},
{
name: "GeminiRequest_ContentsWithSystemInstruction",
raw: `{"model":"gemini-2.0","config":{"system_instruction":"Be brief."},"contents":[{"role":"user","parts":[{"text":"Hi"}]},{"role":"user","parts":[{"function_response":{"name":"get_weather","response":{"temp":20}}}]}]}`,
want: []Message{
text(MessageRoleSystem, "Be brief."),
text(MessageRoleUser, "Hi"),
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, Name: "get_weather", Content: `{"temp":20}`}}},
},
},
{
name: "GeminiRequest_StringContents",
raw: `{"contents":"Hi there","model":"gemini-2.0"}`,
want: []Message{text(MessageRoleUser, "Hi there")},
},
{
name: "MicrosoftAgent_ArrayToolCallIDs",
raw: `[{"role":"assistant","parts":[{"type":"tool_call","id":["run_1","call_9"],"name":"lookup","arguments":{"q":"x"}}]},{"role":"tool","parts":[{"type":"tool_call_response","id":["run_1","call_9"],"response":"found"}]}]`,
want: []Message{
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_9", Name: "lookup", Arguments: map[string]any{"q": "x"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_9", Content: "found"}}},
},
},
{
name: "PydanticAI_ToolCallResponseResultKey",
raw: `[{"role":"user","parts":[{"type":"tool_call_response","id":"call_1","name":"lookup","result":{"ok":true}}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_1", Name: "lookup", Content: `{"ok":true}`}}}},
},
{
name: "SemanticKernel_EventContentWrapper",
raw: `[{"role":"system","gen_ai.event.content":"{\"role\":\"system\",\"content\":\"Be brief.\",\"tool_calls\":[]}","gen_ai.system":"openai"},{"gen_ai.event.content":"{\"index\":0,\"message\":{\"role\":\"Assistant\",\"content\":\"Paris.\"},\"finish_reason\":\"Stop\"}"}]`,
want: []Message{
text(MessageRoleSystem, "Be brief."),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Paris."}}, FinishReason: FinishReasonStop},
},
},
{
name: "BedrockConverse_ToolUseAndToolResult",
raw: `{"messages":[{"role":"user","content":[{"text":"Weather?"}]},{"role":"assistant","content":[{"toolUse":{"toolUseId":"t1","name":"get_weather","input":{"city":"Paris"}}}]},{"role":"user","content":[{"toolResult":{"toolUseId":"t1","content":[{"text":"20C"}],"status":"error"}}]}],"system":[{"text":"Be brief."}]}`,
want: []Message{
text(MessageRoleSystem, "Be brief."),
text(MessageRoleUser, "Weather?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "t1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}},
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "t1", Content: "20C", IsError: true}}},
},
},
{
name: "AnthropicRequest_SystemString",
raw: `{"model":"claude","system":"Be brief.","messages":[{"role":"user","content":"Hi"}],"max_tokens":100}`,
want: []Message{text(MessageRoleSystem, "Be brief."), text(MessageRoleUser, "Hi")},
},
{
name: "OpenAIResponses_BuiltInToolCallIsServer",
raw: `{"object":"response","status":"completed","output":[{"type":"web_search_call","id":"ws_1","status":"completed","action":{"type":"search","query":"SigNoz"}},{"type":"custom_tool_call","call_id":"c1","name":"grep","input":"foo"},{"type":"message","role":"assistant","content":[{"type":"output_text","text":"Found it."}]}]}`,
want: []Message{
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "ws_1", Name: "web_search_call", Arguments: map[string]any{"action": map[string]any{"type": "search", "query": "SigNoz"}}, Server: true}}},
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "grep", Arguments: "foo"}}},
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Found it."}}, FinishReason: FinishReasonStop},
},
},
{
name: "NestedMessageList_Unwrapped",
raw: `[[{"role":"user","content":"Hi"}]]`,
want: []Message{text(MessageRoleUser, "Hi")},
},
{
name: "StringifiedMessages_Decoded",
raw: `{"messages":"[{\"role\":\"user\",\"content\":\"Hi\"}]"}`,
want: []Message{text(MessageRoleUser, "Hi")},
},
{
name: "EmbeddingsRequest_Generic",
raw: `{"input": ["a", "b"], "model": "text-embedding-3-small"}`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `{"input": ["a", "b"], "model": "text-embedding-3-small"}`}}}},
},
{
name: "BedrockConverseResponse_OutputMessageWrapper",
raw: `{"output":{"message":{"role":"assistant","content":[{"text":"20C in Paris."}]}},"stopReason":"end_turn","usage":{"inputTokens":10}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "20C in Paris."}}, FinishReason: FinishReasonStop}},
},
{
name: "OllamaResponse_MessageWrapper",
raw: `{"model":"llama3","message":{"role":"assistant","content":"Hi!"},"done":true,"done_reason":"stop"}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Hi!"}}, FinishReason: FinishReasonStop}},
},
{
name: "CohereV2Response_MessageWrapper",
raw: `{"id":"x","message":{"role":"assistant","tool_calls":[{"id":"c1","type":"function","function":{"name":"get_weather","arguments":"{\"city\":\"Paris\"}"}}]},"finish_reason":"TOOL_CALL"}`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonToolCall, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}}},
},
{
name: "OpenAIToolMessage_TextBlocksBecomeToolResult",
raw: `[{"role":"tool","tool_call_id":"c1","content":[{"type":"text","text":"20C"},{"type":"text","text":"clear"}]}]`,
want: []Message{{Role: MessageRoleTool, Content: []Part{
{Type: PartTypeToolResult, ToolCallID: "c1", Content: "20C"},
{Type: PartTypeToolResult, ToolCallID: "c1", Content: "clear"},
}}},
},
{
name: "AnthropicToolResult_TextBlocksJoined",
raw: `[{"role":"user","content":[{"type":"tool_result","tool_use_id":"t1","content":[{"type":"text","text":"line one"},{"type":"text","text":"line two"}]}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "t1", Content: "line one\nline two"}}}},
},
{
name: "MessageWithContentParts_GeminiNested",
raw: `[{"role":"model","content":{"parts":[{"text":"Hi"}],"role":"model"}}]`,
want: []Message{text(MessageRoleAssistant, "Hi")},
},
{
name: "ToolCallTypedItem_PlainToolCall",
raw: `[{"type":"tool_call","id":"c1","name":"get_weather","args":{"city":"Paris"}}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}}},
},
{
name: "OpenAIToolCallList_Bare",
raw: `[{"id":"c1","type":"function","function":{"name":"get_weather","arguments":"{\"city\":\"Paris\"}"}}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}}},
},
{
name: "ToolDefinitionList_Generic",
raw: `[{"type":"function","function":{"name":"get_weather","description":"Weather","parameters":{"type":"object"}}}]`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `[{"type":"function","function":{"name":"get_weather","description":"Weather","parameters":{"type":"object"}}}]`}}}},
},
{
name: "ContentBlockList_RolelessMessage",
raw: `[{"type":"text","text":"Let me check."},{"type":"tool_use","id":"t1","name":"lookup","input":{"q":"x"}}]`,
want: []Message{{Content: []Part{
{Type: PartTypeText, Content: "Let me check."},
{Type: PartTypeToolCall, ID: "t1", Name: "lookup", Arguments: map[string]any{"q": "x"}},
}}},
},
{
name: "ListOfJSONStrings_Decoded",
raw: []any{`{"role":"user","content":"Hi"}`, `{"role":"assistant","content":"Hello"}`},
want: []Message{text(MessageRoleUser, "Hi"), text(MessageRoleAssistant, "Hello")},
},
{
name: "SingleMessageObject_Converted",
raw: `{"role":"assistant","content":"Done."}`,
want: []Message{text(MessageRoleAssistant, "Done.")},
},
{
name: "UnknownRoleAndFinishReason_KeptLowercased",
raw: `[{"role":"Narrator","parts":[{"type":"text","content":"x"}],"finish_reason":"Weird"}]`,
want: []Message{{Role: "narrator", Content: []Part{{Type: PartTypeText, Content: "x"}}, FinishReason: "weird"}},
},
{
name: "UnknownPartType_Generic",
raw: `[{"role":"user","parts":[{"type":"image","url":"http://x/y.png"}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{{Type: PartTypeGeneric, Content: `{"type":"image","url":"http://x/y.png"}`}}}},
},
{
name: "PlainText_Generic",
raw: "Let the cost of the ball be x dollars.",
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: "Let the cost of the ball be x dollars."}}}},
},
{
name: "UnknownJSONShape_GenericWithOriginal",
raw: `{"output": "{\"query\": \"SigNoz\"}", "kwargs": {"name": "search_web"}}`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `{"output": "{\"query\": \"SigNoz\"}", "kwargs": {"name": "search_web"}}`}}}},
},
{
name: "JSONEncodedString_Generic",
raw: `"{\"query\": \"SigNoz\"}"`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `"{\"query\": \"SigNoz\"}"`}}}},
},
{
name: "DecodedValue_Converted",
raw: []any{map[string]any{"role": "user", "content": "hi"}},
want: []Message{text(MessageRoleUser, "hi")},
},
{
name: "EmptyList_NoMessages",
raw: `[]`,
want: []Message{},
},
{
name: "Nil_NoMessages",
raw: nil,
want: []Message{},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.want, NormalizeMessages(testCase.raw))
})
}
}

View File

@@ -49,66 +49,55 @@ func (o ListOrder) IsValid() bool {
return slices.ContainsFunc(o.Enum(), func(v any) bool { return v == o })
}
type ListRulesParams struct {
Query string `query:"query"`
// gin cannot bind a slice of valuer enums; AlertStates converts these.
States []string `query:"states"`
Sort ListSort `query:"sort"`
Order ListOrder `query:"order"`
Limit int `query:"limit"`
Offset int `query:"offset"`
// ListFilter is the rule listing state shared by the v3 list params and saved views.
type ListFilter struct {
Query string `query:"query" json:"query"`
// gin cannot bind a slice of valuer enums; GetAlertStates converts these.
States []string `query:"states" json:"states" nullable:"false"`
Sort ListSort `query:"sort" json:"sort"`
Order ListOrder `query:"order" json:"order"`
}
// Validate normalizes in place; an over-max limit is clamped, not rejected.
func (p *ListRulesParams) Validate() error {
if n := utf8.RuneCountInString(p.Query); n > MaxListQueryLen {
// Validate normalizes in place; zero sort/order get the defaults, nil states an empty slice.
func (f *ListFilter) Validate() error {
if f.States == nil {
f.States = []string{}
}
if n := utf8.RuneCountInString(f.Query); n > MaxListQueryLen {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"query cannot be longer than %d characters, got %d", MaxListQueryLen, n)
}
if p.Sort.IsZero() {
p.Sort = ListSortUpdatedAt
} else if !p.Sort.IsValid() {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid sort %q, expected one of: `updated_at`, `created_at`, `name`, `state`, `severity`", p.Sort)
}
if p.Order.IsZero() {
p.Order = ListOrderDesc
} else if !p.Order.IsValid() {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid order %q, expected `asc` or `desc`", p.Order)
}
if p.Limit == 0 {
p.Limit = DefaultListLimit
} else if p.Limit < 0 {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid limit %d, must be a positive integer", p.Limit)
} else if p.Limit > MaxListLimit {
p.Limit = MaxListLimit
}
if p.Offset < 0 {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid offset %d, must be a non-negative integer", p.Offset)
}
if _, err := p.GetAlertStates(); err != nil {
if _, err := f.GetAlertStates(); err != nil {
return err
}
if f.Sort.IsZero() {
f.Sort = ListSortUpdatedAt
} else if !f.Sort.IsValid() {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid sort %q, expected one of: `updated_at`, `created_at`, `name`, `state`, `severity`", f.Sort)
}
if f.Order.IsZero() {
f.Order = ListOrderDesc
} else if !f.Order.IsValid() {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid order %q, expected `asc` or `desc`", f.Order)
}
return nil
}
// GetAlertStates parses States; empty means no state filtering.
func (p *ListRulesParams) GetAlertStates() ([]AlertState, error) {
if len(p.States) == 0 {
func (f *ListFilter) GetAlertStates() ([]AlertState, error) {
if len(f.States) == 0 {
return nil, nil
}
states := make([]AlertState, 0, len(p.States))
for _, raw := range p.States {
states := make([]AlertState, 0, len(f.States))
for _, raw := range f.States {
state, err := parseAlertState(raw)
if err != nil {
return nil, err
@@ -127,3 +116,32 @@ func parseAlertState(raw string) (AlertState, error) {
}
return state, nil
}
type ListRulesParams struct {
ListFilter
Limit int `query:"limit"`
Offset int `query:"offset"`
}
// Validate normalizes in place; an over-max limit is clamped, not rejected.
func (p *ListRulesParams) Validate() error {
if err := p.ListFilter.Validate(); err != nil {
return err
}
if p.Limit == 0 {
p.Limit = DefaultListLimit
} else if p.Limit < 0 {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid limit %d, must be a positive integer", p.Limit)
} else if p.Limit > MaxListLimit {
p.Limit = MaxListLimit
}
if p.Offset < 0 {
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
"invalid offset %d, must be a non-negative integer", p.Offset)
}
return nil
}

View File

@@ -27,7 +27,7 @@ func TestListRulesParamsValidate(t *testing.T) {
},
{
name: "ExplicitValues_Kept",
params: ListRulesParams{Sort: ListSortSeverity, Order: ListOrderAsc, Limit: 50, Offset: 100},
params: ListRulesParams{ListFilter: ListFilter{Sort: ListSortSeverity, Order: ListOrderAsc}, Limit: 50, Offset: 100},
wantSort: ListSortSeverity,
wantOrder: ListOrderAsc,
wantLimit: 50,
@@ -41,17 +41,17 @@ func TestListRulesParamsValidate(t *testing.T) {
},
{
name: "InvalidState_Rejected",
params: ListRulesParams{States: []string{"bogus"}},
params: ListRulesParams{ListFilter: ListFilter{States: []string{"bogus"}}},
wantErr: `invalid state "bogus"`,
},
{
name: "InvalidSort_Rejected",
params: ListRulesParams{Sort: ListSort{valuer.NewString("bogus")}},
params: ListRulesParams{ListFilter: ListFilter{Sort: ListSort{valuer.NewString("bogus")}}},
wantErr: "invalid sort",
},
{
name: "InvalidOrder_Rejected",
params: ListRulesParams{Order: ListOrder{valuer.NewString("bogus")}},
params: ListRulesParams{ListFilter: ListFilter{Order: ListOrder{valuer.NewString("bogus")}}},
wantErr: "invalid order",
},
{
@@ -66,7 +66,7 @@ func TestListRulesParamsValidate(t *testing.T) {
},
{
name: "OverLongQuery_Rejected",
params: ListRulesParams{Query: strings.Repeat("a", MaxListQueryLen+1)},
params: ListRulesParams{ListFilter: ListFilter{Query: strings.Repeat("a", MaxListQueryLen+1)}},
wantErr: "query cannot be longer",
},
}
@@ -112,7 +112,7 @@ func TestListRulesParamsAlertStates(t *testing.T) {
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
params := ListRulesParams{States: tc.states}
params := ListRulesParams{ListFilter: ListFilter{States: tc.states}}
states, err := params.GetAlertStates()
if tc.wantErr != "" {
require.Error(t, err)

View File

@@ -65,4 +65,10 @@ type RuleStore interface {
GetStoredRuleLabels(context.Context, string) ([]string, error)
GetStoredRule(context.Context, valuer.UUID, valuer.UUID) (*StorableRule, error)
GetStoredRulesByMetricName(context.Context, string, string) ([]RuleAlert, error)
CreateRuleView(context.Context, *StorableRuleView) error
GetRuleView(context.Context, valuer.UUID, valuer.UUID) (*StorableRuleView, error)
ListRuleViews(context.Context, valuer.UUID) ([]*StorableRuleView, error)
UpdateRuleView(context.Context, *StorableRuleView) error
DeleteRuleView(context.Context, valuer.UUID, valuer.UUID) error
}

View File

@@ -0,0 +1,169 @@
package ruletypes
import (
"bytes"
"encoding/json"
"strings"
"time"
"unicode/utf8"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/uptrace/bun"
)
const (
RuleViewSchemaVersion = "v1"
MaxRuleViewNameLen = 64
)
var (
ErrCodeRuleViewInvalidInput = errors.MustNewCode("rule_view_invalid_input")
ErrCodeRuleViewNotFound = errors.MustNewCode("rule_view_not_found")
)
type StorableRuleView struct {
bun.BaseModel `bun:"table:rule_view,alias:rule_view"`
types.Identifiable
types.TimeAuditable
Name string `bun:"name,type:text,notnull"`
Data storableRuleViewData `bun:"data,type:text,notnull"`
OrgID valuer.UUID `bun:"org_id,type:text,notnull"`
}
func (s *StorableRuleView) ToGettableRuleView() *GettableRuleView {
return &GettableRuleView{
ID: s.ID,
Name: s.Name,
Data: s.Data.toRuleViewData(),
OrgID: s.OrgID,
CreatedAt: s.CreatedAt,
UpdatedAt: s.UpdatedAt,
}
}
func (s *StorableRuleView) Update(updatable UpdatableRuleView) {
s.Name = updatable.Name
s.Data = newStorableRuleViewData(updatable.Data)
s.UpdatedAt = time.Now()
}
type GettableRuleView struct {
ID valuer.UUID `json:"id" required:"true"`
Name string `json:"name" required:"true"`
Data RuleViewData `json:"data" required:"true"`
OrgID valuer.UUID `json:"orgId" required:"true"`
CreatedAt time.Time `json:"createdAt" required:"true"`
UpdatedAt time.Time `json:"updatedAt" required:"true"`
}
// RuleViewData holds the rule listing state (ListRulesParams minus pagination) a view replays.
type RuleViewData struct {
Version string `json:"version" required:"true"`
ListFilter
}
func (d *RuleViewData) Validate() error {
if d.Version != RuleViewSchemaVersion {
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput,
"version must be %q, got %q", RuleViewSchemaVersion, d.Version)
}
return d.ListFilter.Validate()
}
type PostableRuleView struct {
Name string `json:"name" required:"true"`
Data RuleViewData `json:"data" required:"true"`
}
func (p *PostableRuleView) UnmarshalJSON(data []byte) error {
dec := json.NewDecoder(bytes.NewReader(data))
dec.DisallowUnknownFields()
type alias PostableRuleView
var tmp alias
if err := dec.Decode(&tmp); err != nil {
return errors.WrapInvalidInputf(err, ErrCodeRuleViewInvalidInput, "invalid saved view request body").WithAdditional(err.Error())
}
*p = PostableRuleView(tmp)
return p.Validate()
}
func (p *PostableRuleView) Validate() error {
if err := validateRuleViewName(p.Name); err != nil {
return err
}
return p.Data.Validate()
}
func (p PostableRuleView) ToStorableRuleView(orgID valuer.UUID) *StorableRuleView {
now := time.Now()
return &StorableRuleView{
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
TimeAuditable: types.TimeAuditable{CreatedAt: now, UpdatedAt: now},
Name: p.Name,
Data: newStorableRuleViewData(p.Data),
OrgID: orgID,
}
}
type UpdatableRuleView = PostableRuleView
func NewGettableRuleViewsFromStorableRuleViews(storables []*StorableRuleView) []*GettableRuleView {
views := make([]*GettableRuleView, 0, len(storables))
for _, storable := range storables {
views = append(views, storable.ToGettableRuleView())
}
return views
}
type ListableRuleViews struct {
Views []*GettableRuleView `json:"views" required:"true" nullable:"false"`
}
// storableRuleViewData owns the persisted blob format; wire tag changes must not affect stored rows.
type storableRuleViewData struct {
Version string `json:"version"`
Query string `json:"query"`
States []string `json:"states"`
Sort string `json:"sort"`
Order string `json:"order"`
}
func newStorableRuleViewData(data RuleViewData) storableRuleViewData {
return storableRuleViewData{
Version: data.Version,
Query: data.Query,
States: data.States,
Sort: data.Sort.StringValue(),
Order: data.Order.StringValue(),
}
}
func (d storableRuleViewData) toRuleViewData() RuleViewData {
return RuleViewData{
Version: d.Version,
ListFilter: ListFilter{
Query: d.Query,
States: d.States,
Sort: ListSort{valuer.NewString(d.Sort)},
Order: ListOrder{valuer.NewString(d.Order)},
},
}
}
func validateRuleViewName(name string) error {
if strings.TrimSpace(name) == "" {
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput, "name is required")
}
if name != strings.TrimSpace(name) {
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput, "name must not have leading or trailing whitespace")
}
if n := utf8.RuneCountInString(name); n > MaxRuleViewNameLen {
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput,
"name must be at most %d characters, got %d", MaxRuleViewNameLen, n)
}
return nil
}

View File

@@ -0,0 +1,210 @@
package ruletypes
import (
"encoding/json"
"strings"
"testing"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestRuleViewDataValidate(t *testing.T) {
testCases := []struct {
name string
data RuleViewData
expectError bool
}{
{
name: "AllFieldsSet_Valid",
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Query: "name CONTAINS 'prod'", States: []string{"firing", "pending"}, Sort: ListSortName, Order: ListOrderAsc}},
expectError: false,
},
{
name: "ZeroStatesSortOrder_Valid",
data: RuleViewData{Version: RuleViewSchemaVersion},
expectError: false,
},
{
name: "QueryOverCap_Rejected",
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Query: strings.Repeat("x", MaxListQueryLen+1)}},
expectError: true,
},
{
name: "WrongVersion_Rejected",
data: RuleViewData{Version: "v2"},
expectError: true,
},
{
name: "EmptyVersion_Rejected",
data: RuleViewData{},
expectError: true,
},
{
name: "UnknownState_Rejected",
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"exploding"}}},
expectError: true,
},
{
name: "UnknownSort_Rejected",
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Sort: ListSort{valuer.NewString("bogus")}}},
expectError: true,
},
{
name: "UnknownOrder_Rejected",
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Order: ListOrder{valuer.NewString("sideways")}}},
expectError: true,
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
err := testCase.data.Validate()
if testCase.expectError {
assert.Error(t, err)
} else {
assert.NoError(t, err)
}
})
}
}
func TestRuleViewDataValidateDefaults(t *testing.T) {
data := RuleViewData{Version: RuleViewSchemaVersion}
require.NoError(t, data.Validate())
assert.Equal(t, ListSortUpdatedAt, data.Sort)
assert.Equal(t, ListOrderDesc, data.Order)
assert.Equal(t, []string{}, data.States)
}
func TestPostableRuleViewUnmarshalJSON(t *testing.T) {
testCases := []struct {
name string
body string
expectError bool
expectedErrMsg string
expectedName string
}{
{
name: "ValidBody_NameKeptAsIs",
body: `{"name":"my view","data":{"version":"v1","query":"severity = 'critical'","states":["firing"],"sort":"name","order":"asc"}}`,
expectError: false,
expectedName: "my view",
},
{
name: "NameSurroundingWhitespace_Rejected",
body: `{"name":" my view ","data":{"version":"v1"}}`,
expectError: true,
expectedErrMsg: "name must not have leading or trailing whitespace",
},
{
name: "UnknownField_Rejected",
body: `{"name":"my view","data":{"version":"v1"},"extra":true}`,
expectError: true,
},
{
name: "BlankName_Rejected",
body: `{"name":" ","data":{"version":"v1"}}`,
expectError: true,
expectedErrMsg: "name is required",
},
{
name: "NameOverMaxLength_Rejected",
body: `{"name":"` + strings.Repeat("x", MaxRuleViewNameLen+1) + `","data":{"version":"v1"}}`,
expectError: true,
expectedErrMsg: "name must be at most",
},
{
name: "InvalidDataVersion_Rejected",
body: `{"name":"my view","data":{"version":"v9"}}`,
expectError: true,
},
{
name: "InvalidState_Rejected",
body: `{"name":"my view","data":{"version":"v1","states":["exploding"]}}`,
expectError: true,
expectedErrMsg: "invalid state",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
var p PostableRuleView
err := json.Unmarshal([]byte(testCase.body), &p)
if testCase.expectError {
assert.Error(t, err)
if testCase.expectedErrMsg != "" {
assert.ErrorContains(t, err, testCase.expectedErrMsg)
}
return
}
require.NoError(t, err)
assert.Equal(t, testCase.expectedName, p.Name)
})
}
}
func TestToStorableRuleView(t *testing.T) {
orgID := valuer.GenerateUUID()
postable := PostableRuleView{
Name: "my view",
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"firing"}, Sort: ListSortName, Order: ListOrderAsc}},
}
storable := postable.ToStorableRuleView(orgID)
assert.Equal(t, orgID, storable.OrgID)
assert.Equal(t, "my view", storable.Name)
assert.False(t, storable.ID.IsZero())
assert.False(t, storable.CreatedAt.IsZero())
assert.Equal(t, storable.CreatedAt, storable.UpdatedAt)
gettable := storable.ToGettableRuleView()
assert.Equal(t, storable.ID, gettable.ID)
assert.Equal(t, "my view", gettable.Name)
assert.Equal(t, postable.Data, gettable.Data)
assert.Equal(t, orgID, gettable.OrgID)
assert.Equal(t, storable.CreatedAt, gettable.CreatedAt)
}
func TestNewGettableRuleViewsFromStorableRuleViews(t *testing.T) {
orgID := valuer.GenerateUUID()
first := PostableRuleView{
Name: "first view",
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Sort: ListSortName, Order: ListOrderAsc}},
}.ToStorableRuleView(orgID)
second := PostableRuleView{
Name: "second view",
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"firing"}, Sort: ListSortState, Order: ListOrderDesc}},
}.ToStorableRuleView(orgID)
views := NewGettableRuleViewsFromStorableRuleViews([]*StorableRuleView{first, second})
require.Len(t, views, 2)
assert.Equal(t, "first view", views[0].Name)
assert.Equal(t, second.Data.States, views[1].Data.States)
assert.Equal(t, ListSortState, views[1].Data.Sort)
}
func TestStorableRuleViewUpdate(t *testing.T) {
orgID := valuer.GenerateUUID()
storable := PostableRuleView{
Name: "original",
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Sort: ListSortName, Order: ListOrderAsc}},
}.ToStorableRuleView(orgID)
createdAt := storable.CreatedAt
storable.Update(UpdatableRuleView{
Name: "renamed",
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"disabled"}, Sort: ListSortCreatedAt, Order: ListOrderDesc}},
})
gettable := storable.ToGettableRuleView()
assert.Equal(t, "renamed", gettable.Name)
assert.Equal(t, []string{"disabled"}, gettable.Data.States)
assert.Equal(t, ListSortCreatedAt, gettable.Data.Sort)
assert.Equal(t, ListOrderDesc, gettable.Data.Order)
assert.Equal(t, createdAt, storable.CreatedAt)
assert.True(t, storable.UpdatedAt.After(createdAt))
}

View File

@@ -36,7 +36,6 @@ type TraceStore interface {
GetMinimalSpans(ctx context.Context, traceID string, start, end time.Time) ([]MinimalSpan, error)
GetTraceSpansByIDs(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
GetThreadSpans(ctx context.Context, traceID string, summary *TraceSummary, cursor *ThreadCursor, limit int) ([]StorableSpan, error)
GetSpanCountByField(ctx context.Context, traceID string, summary *TraceSummary, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)
GetSpanDurationByField(ctx context.Context, traceID string, summary *TraceSummary, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)

View File

@@ -1,123 +0,0 @@
package spantypes
import (
"encoding/base64"
"encoding/json"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
)
const (
threadDefaultLimit = 100
threadMaxLimit = 1000
)
var (
ErrCodeThreadInvalidLimit = errors.MustNewCode("trace_thread_invalid_limit")
ErrCodeThreadInvalidCursor = errors.MustNewCode("trace_thread_invalid_cursor")
)
type PostableThreadQuery struct {
// Limit is the page size; 0 means 100.
Limit int `query:"limit"`
// Cursor is the nextCursor of the previous page; empty for the first page.
Cursor string `query:"cursor"`
}
type ThreadQuery struct {
Limit int
Cursor *ThreadCursor
}
func NewThreadQuery(postable *PostableThreadQuery) (*ThreadQuery, error) {
query := &ThreadQuery{Limit: postable.Limit}
if query.Limit < 0 {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidLimit, "limit cannot be negative, got %d", query.Limit)
}
if query.Limit == 0 {
query.Limit = threadDefaultLimit
}
if query.Limit > threadMaxLimit {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidLimit, "limit cannot exceed %d, got %d", threadMaxLimit, query.Limit)
}
if postable.Cursor != "" {
cursor, err := DecodeThreadCursor(postable.Cursor)
if err != nil {
return nil, err
}
query.Cursor = cursor
}
return query, nil
}
// ThreadCursor is the (TimeUnixNano, SpanID) of the last span of a page.
type ThreadCursor struct {
TimeUnixNano uint64 `json:"t"`
SpanID string `json:"s"`
}
func (c ThreadCursor) Encode() string {
data, _ := json.Marshal(c)
return base64.RawURLEncoding.EncodeToString(data)
}
func DecodeThreadCursor(cursor string) (*ThreadCursor, error) {
data, err := base64.RawURLEncoding.DecodeString(cursor)
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
}
c := new(ThreadCursor)
if err := json.Unmarshal(data, c); err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
}
if c.SpanID == "" {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidCursor, "invalid cursor: missing span id")
}
return c, nil
}
type GettableTraceThread struct {
Spans []*ThreadSpan `json:"spans" required:"true" nullable:"false"`
NextCursor string `json:"nextCursor,omitempty"`
}
// ThreadSpan sets the formatted fields only when the span has the matching gen_ai messages attribute.
type ThreadSpan struct {
WaterfallSpan
FormattedInput []aiobservabilitytypes.Message `json:"formatted_input,omitempty"`
FormattedOutput []aiobservabilitytypes.Message `json:"formatted_output,omitempty"`
}
// NewGettableTraceThread expects limit+1 spans; the extra one only signals a next page.
func NewGettableTraceThread(traceID string, spans []StorableSpan, limit int) *GettableTraceThread {
hasMore := len(spans) > limit
if hasMore {
spans = spans[:limit]
}
out := make([]*ThreadSpan, len(spans))
for i := range spans {
out[i] = newThreadSpan(traceID, &spans[i])
}
thread := &GettableTraceThread{Spans: out}
if hasMore {
last := spans[len(spans)-1]
thread.NextCursor = ThreadCursor{TimeUnixNano: uint64(last.StartTime.UnixNano()), SpanID: last.SpanID}.Encode()
}
return thread
}
func newThreadSpan(traceID string, storable *StorableSpan) *ThreadSpan {
span := &ThreadSpan{WaterfallSpan: *storable.ToWaterfallSpan(traceID)}
// client expects millis, as in the waterfall
span.TimeUnix = span.TimeUnix / 1_000_000
if v, ok := span.Attributes[aiobservabilitytypes.GenAIInputMessages]; ok {
span.FormattedInput = aiobservabilitytypes.NormalizeMessages(v)
}
if v, ok := span.Attributes[aiobservabilitytypes.GenAIOutputMessages]; ok {
span.FormattedOutput = aiobservabilitytypes.NormalizeMessages(v)
}
return span
}

View File

@@ -1,173 +0,0 @@
package spantypes
import (
"testing"
"time"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestNewThreadQuery(t *testing.T) {
cursor := ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}
testCases := []struct {
name string
postable PostableThreadQuery
want *ThreadQuery
wantErr bool
}{
{name: "ZeroLimit_UsesDefault", postable: PostableThreadQuery{}, want: &ThreadQuery{Limit: threadDefaultLimit}},
{name: "PositiveLimit_Kept", postable: PostableThreadQuery{Limit: 25}, want: &ThreadQuery{Limit: 25}},
{name: "MaxLimit_Kept", postable: PostableThreadQuery{Limit: threadMaxLimit}, want: &ThreadQuery{Limit: threadMaxLimit}},
{name: "AboveMaxLimit_Rejected", postable: PostableThreadQuery{Limit: threadMaxLimit + 1}, wantErr: true},
{name: "NegativeLimit_Rejected", postable: PostableThreadQuery{Limit: -1}, wantErr: true},
{name: "Cursor_Decoded", postable: PostableThreadQuery{Limit: 10, Cursor: cursor.Encode()}, want: &ThreadQuery{Limit: 10, Cursor: &cursor}},
{name: "InvalidCursor_Rejected", postable: PostableThreadQuery{Cursor: "not base64!"}, wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
got, err := NewThreadQuery(&testCase.postable)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, got)
})
}
}
func TestDecodeThreadCursor(t *testing.T) {
testCases := []struct {
name string
cursor string
want *ThreadCursor
wantErr bool
}{
{name: "EncodedCursor_RoundTrips", cursor: ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}.Encode(), want: &ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}},
{name: "NotBase64_Rejected", cursor: "not base64!", wantErr: true},
{name: "NotJSON_Rejected", cursor: "bm90IGpzb24", wantErr: true},
{name: "MissingSpanID_Rejected", cursor: "eyJ0IjogMX0", wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
got, err := DecodeThreadCursor(testCase.cursor)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, got)
})
}
}
func TestNewGettableTraceThread(t *testing.T) {
spans := []StorableSpan{
{SpanID: "a", StartTime: time.Unix(1, 500_000_000)},
{SpanID: "b", StartTime: time.Unix(2, 0)},
{SpanID: "c", StartTime: time.Unix(3, 0)},
}
testCases := []struct {
name string
spans []StorableSpan
limit int
wantSpanIDs []string
wantTimeUnix []uint64
wantNextCursor string
}{
{name: "MoreThanLimit_TrimsAndSetsCursor", spans: spans, limit: 2, wantSpanIDs: []string{"a", "b"}, wantTimeUnix: []uint64{1500, 2000}, wantNextCursor: ThreadCursor{TimeUnixNano: 2_000_000_000, SpanID: "b"}.Encode()},
{name: "WithinLimit_NoCursor", spans: spans, limit: 3, wantSpanIDs: []string{"a", "b", "c"}, wantTimeUnix: []uint64{1500, 2000, 3000}},
{name: "NoSpans_EmptyList", limit: 3, wantSpanIDs: []string{}, wantTimeUnix: []uint64{}},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
thread := NewGettableTraceThread("trace-1", testCase.spans, testCase.limit)
require.NotNil(t, thread.Spans)
spanIDs := make([]string, len(thread.Spans))
timeUnix := make([]uint64, len(thread.Spans))
for i, span := range thread.Spans {
spanIDs[i] = span.SpanID
timeUnix[i] = span.TimeUnix
assert.Equal(t, "trace-1", span.TraceID)
}
assert.Equal(t, testCase.wantSpanIDs, spanIDs)
assert.Equal(t, testCase.wantTimeUnix, timeUnix)
assert.Equal(t, testCase.wantNextCursor, thread.NextCursor)
})
}
}
func TestNewThreadSpan(t *testing.T) {
userHi := []aiobservabilitytypes.Message{{
Role: aiobservabilitytypes.MessageRoleUser,
Content: []aiobservabilitytypes.Part{{Type: aiobservabilitytypes.PartTypeText, Content: "hi"}},
}}
assistantHello := []aiobservabilitytypes.Message{{
Role: aiobservabilitytypes.MessageRoleAssistant,
Content: []aiobservabilitytypes.Part{{Type: aiobservabilitytypes.PartTypeText, Content: "hello"}},
FinishReason: aiobservabilitytypes.FinishReasonStop,
}}
testCases := []struct {
name string
span StorableSpan
wantInput []aiobservabilitytypes.Message
wantOutput []aiobservabilitytypes.Message
wantAttrs map[string]any
}{
{
name: "MessagesInLegacyMap",
span: StorableSpan{AttributesString: map[string]string{
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
}},
wantInput: userHi,
wantOutput: assistantHello,
wantAttrs: map[string]any{
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
},
},
{
name: "MessagesInJSONColumn_FlattenedToDottedKeys",
span: StorableSpan{AttributesJSON: telemetrystoretypes.JSONValue{
"gen_ai": map[string]any{
"input": map[string]any{"messages": `[{"role":"user","content":"hi"}]`},
"request": map[string]any{"model": "gpt-4o"},
},
}},
wantInput: userHi,
wantAttrs: map[string]any{
"gen_ai.input.messages": `[{"role":"user","content":"hi"}]`,
"gen_ai.request.model": "gpt-4o",
},
},
{
name: "LegacyMapWinsOverJSONColumn",
span: StorableSpan{
AttributesJSON: telemetrystoretypes.JSONValue{"gen_ai": map[string]any{"request": map[string]any{"model": "json"}}},
AttributesString: map[string]string{"gen_ai.request.model": "map"},
},
wantAttrs: map[string]any{"gen_ai.request.model": "map"},
},
{
name: "NoMessages_FieldsUnset",
span: StorableSpan{AttributesString: map[string]string{"http.method": "GET"}},
wantAttrs: map[string]any{"http.method": "GET"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
span := newThreadSpan("trace-1", &testCase.span)
assert.Equal(t, testCase.wantInput, span.FormattedInput)
assert.Equal(t, testCase.wantOutput, span.FormattedOutput)
assert.Equal(t, testCase.wantAttrs, span.Attributes)
})
}
}

View File

@@ -8,7 +8,6 @@ import (
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
@@ -94,36 +93,35 @@ type WaterfallSpan struct {
// StorableSpan is the ClickHouse scan struct for the v3 waterfall query.
type StorableSpan struct {
StartTime time.Time `ch:"timestamp"`
DurationNano uint64 `ch:"duration_nano"`
SpanID string `ch:"span_id"`
HasError bool `ch:"has_error"`
Kind int8 `ch:"kind"`
ServiceName string `ch:"resource_string_service$$name"`
Name string `ch:"name"`
AttributesString map[string]string `ch:"attributes_string"`
AttributesNumber map[string]float64 `ch:"attributes_number"`
AttributesBool map[string]bool `ch:"attributes_bool"`
AttributesJSON telemetrystoretypes.JSONValue `ch:"attributes"`
ResourcesString map[string]string `ch:"resources_string"`
Events []string `ch:"events"`
StatusMessage string `ch:"status_message"`
StatusCodeString string `ch:"status_code_string"`
SpanKind string `ch:"kind_string"`
ParentSpanID string `ch:"parent_span_id"`
Flags uint32 `ch:"flags"`
IsRemote string `ch:"is_remote"`
TraceState string `ch:"trace_state"`
StatusCode int16 `ch:"status_code"`
DBName string `ch:"db_name"`
DBOperation string `ch:"db_operation"`
HTTPMethod string `ch:"http_method"`
HTTPURL string `ch:"http_url"`
HTTPHost string `ch:"http_host"`
ExternalHTTPMethod string `ch:"external_http_method"`
ExternalHTTPURL string `ch:"external_http_url"`
ResponseStatusCode string `ch:"response_status_code"`
References string `ch:"references"`
StartTime time.Time `ch:"timestamp"`
DurationNano uint64 `ch:"duration_nano"`
SpanID string `ch:"span_id"`
HasError bool `ch:"has_error"`
Kind int8 `ch:"kind"`
ServiceName string `ch:"resource_string_service$$name"`
Name string `ch:"name"`
AttributesString map[string]string `ch:"attributes_string"`
AttributesNumber map[string]float64 `ch:"attributes_number"`
AttributesBool map[string]bool `ch:"attributes_bool"`
ResourcesString map[string]string `ch:"resources_string"`
Events []string `ch:"events"`
StatusMessage string `ch:"status_message"`
StatusCodeString string `ch:"status_code_string"`
SpanKind string `ch:"kind_string"`
ParentSpanID string `ch:"parent_span_id"`
Flags uint32 `ch:"flags"`
IsRemote string `ch:"is_remote"`
TraceState string `ch:"trace_state"`
StatusCode int16 `ch:"status_code"`
DBName string `ch:"db_name"`
DBOperation string `ch:"db_operation"`
HTTPMethod string `ch:"http_method"`
HTTPURL string `ch:"http_url"`
HTTPHost string `ch:"http_host"`
ExternalHTTPMethod string `ch:"external_http_method"`
ExternalHTTPURL string `ch:"external_http_url"`
ResponseStatusCode string `ch:"response_status_code"`
References string `ch:"references"`
}
// MinimalSpan with only the fields needed to build the parent-child tree.
@@ -279,10 +277,8 @@ func (item *StorableSpan) AttributeValue(name string) any {
return nil
}
// Attributes flattens the JSON column first, so the legacy maps win on collision.
func (item *StorableSpan) Attributes() map[string]any {
attributes := make(map[string]any, len(item.AttributesString)+len(item.AttributesNumber)+len(item.AttributesBool)+len(item.AttributesJSON))
item.AttributesJSON.FlattenInto("", attributes)
attributes := make(map[string]any, len(item.AttributesString)+len(item.AttributesNumber)+len(item.AttributesBool))
for k, v := range item.AttributesString {
attributes[k] = v
}

View File

@@ -35,21 +35,3 @@ func (v *JSONValue) Scan(src any) error {
*v = decoded
return nil
}
// FlattenInto writes v into out under dotted keys, overwriting existing keys.
func (v JSONValue) FlattenInto(prefix string, out map[string]any) {
for k, value := range v {
key := k
if prefix != "" {
key = prefix + "." + k
}
switch child := value.(type) {
case map[string]any:
JSONValue(child).FlattenInto(key, out)
case JSONValue:
child.FlattenInto(key, out)
default:
out[key] = value
}
}
}

View File

@@ -131,6 +131,36 @@ def seed_alert_rules(
return _seed_alert_rules
@pytest.fixture(name="create_rule_view", scope="function")
def create_rule_view(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> Callable[[dict], dict]:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
view_ids = []
def _create_rule_view(view: dict) -> dict:
response = requests.post(
signoz.self.host_configs["8080"].get("/api/v2/rule_views"),
json=view,
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.CREATED, f"Failed to create rule view, api returned {response.status_code} with response: {response.text}"
created = response.json()["data"]
view_ids.append(created["id"])
return created
yield _create_rule_view
# A view the test already deleted returns 404; only real failures are logged.
for view_id in view_ids:
response = requests.delete(
signoz.self.host_configs["8080"].get(f"/api/v2/rule_views/{view_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=5,
)
if response.status_code not in (HTTPStatus.NO_CONTENT, HTTPStatus.NOT_FOUND):
logger.error("Error deleting rule view: %s", {"view_id": view_id, "response": response.text})
def labels_to_map(labels: list[dict]) -> dict[str, str]:
"""Converts the label list shape of the v2 rule history APIs to a plain map."""
return {label["key"]["name"]: label["value"] for label in labels or []}

View File

@@ -0,0 +1,269 @@
import uuid
from collections.abc import Callable
from http import HTTPStatus
import pytest
import requests
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.types import Operation, SigNoz
BASE_URL = "/api/v2/rule_views"
@pytest.mark.parametrize(
("body", "expected_code", "expected_message"),
[
({"data": {"version": "v1"}}, "rule_view_invalid_input", "name is required"),
({"name": " ", "data": {"version": "v1"}}, "rule_view_invalid_input", "name is required"),
(
{"name": " Storage ", "data": {"version": "v1"}},
"rule_view_invalid_input",
"name must not have leading or trailing whitespace",
),
(
{"name": "x" * 65, "data": {"version": "v1"}},
"rule_view_invalid_input",
"name must be at most 64 characters, got 65",
),
(
{"name": "wrong-version", "data": {"version": "v2"}},
"rule_view_invalid_input",
'version must be "v1", got "v2"',
),
(
{"name": "missing-version", "data": {}},
"rule_view_invalid_input",
'version must be "v1", got ""',
),
(
{"name": "bad-state", "data": {"version": "v1", "states": ["exploding"]}},
"rule_list_invalid",
'invalid state "exploding"',
),
(
{"name": "bad-sort", "data": {"version": "v1", "sort": "bogus"}},
"rule_list_invalid",
"invalid sort",
),
(
{"name": "bad-order", "data": {"version": "v1", "order": "bogus"}},
"rule_list_invalid",
"invalid order",
),
(
{"name": "long-query", "data": {"version": "v1", "query": "x" * 1025}},
"rule_list_invalid",
"query cannot be longer than 1024 characters",
),
(
{"name": "rejects-unknown", "data": {"version": "v1"}, "unknownfield": "boom"},
"rule_view_invalid_input",
"invalid saved view request body",
),
],
ids=[
"missing_name",
"blank_name",
"whitespace_name",
"name_too_long",
"wrong_schema_version",
"missing_version",
"invalid_state",
"invalid_sort",
"invalid_order",
"query_too_long",
"unknown_field",
],
)
def test_create_rejects_invalid_body(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
body: dict,
expected_code: str,
expected_message: str,
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.post(
signoz.self.host_configs["8080"].get(BASE_URL),
json=body,
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.BAD_REQUEST
assert response.json()["error"]["code"] == expected_code
assert expected_message in response.json()["error"]["message"]
def test_update_rejects_malformed_id(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.put(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/not-a-uuid"),
json={"name": "x", "data": {"version": "v1"}},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.BAD_REQUEST
def test_update_missing_view_returns_not_found(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.put(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{uuid.uuid4()}"),
json={"name": "x", "data": {"version": "v1"}},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.NOT_FOUND
assert response.json()["error"]["code"] == "rule_view_not_found"
def test_delete_rejects_malformed_id(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.delete(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/not-a-uuid"),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.BAD_REQUEST
def test_delete_missing_view_returns_not_found(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.delete(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{uuid.uuid4()}"),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.NOT_FOUND
assert response.json()["error"]["code"] == "rule_view_not_found"
def test_rule_view_lifecycle(
signoz: SigNoz,
create_user_admin: Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_rule_view: Callable[[dict], dict],
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
# List assertions filter on this test's names so foreign views never interfere.
owned_names = {"Critical Prod", "Critical Staging", "Disabled"}
created = create_rule_view(
{
"name": "Critical Prod",
"data": {
"version": "v1",
"query": "name CONTAINS 'prod' AND severity = 'critical'",
"states": ["firing", "pending"],
"sort": "name",
"order": "asc",
},
}
)
view_id = created["id"]
assert created["name"] == "Critical Prod"
assert created["data"]["version"] == "v1"
assert created["data"]["query"] == "name CONTAINS 'prod' AND severity = 'critical'"
assert created["data"]["states"] == ["firing", "pending"]
# Omitted states, sort and order are normalized on save: [] and the list defaults, never null.
disabled = create_rule_view({"name": "Disabled", "data": {"version": "v1"}})
assert disabled["name"] == "Disabled"
assert disabled["data"]["states"] == []
assert disabled["data"]["sort"] == "updated_at"
assert disabled["data"]["order"] == "desc"
response = requests.get(
signoz.self.host_configs["8080"].get(BASE_URL),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK, response.text
views = [v for v in response.json()["data"]["views"] if v["name"] in owned_names]
assert {v["name"] for v in views} == {"Critical Prod", "Disabled"}
response = requests.put(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{view_id}"),
json={
"name": "Critical Staging",
"data": {
"version": "v1",
"query": "name CONTAINS 'staging'",
"states": ["firing"],
"sort": "created_at",
"order": "desc",
},
},
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK, response.text
updated = response.json()["data"]
assert updated["id"] == view_id
assert updated["name"] == "Critical Staging"
assert updated["data"]["query"] == "name CONTAINS 'staging'"
assert updated["data"]["states"] == ["firing"]
assert updated["data"]["sort"] == "created_at"
assert updated["data"]["order"] == "desc"
response = requests.get(
signoz.self.host_configs["8080"].get(BASE_URL),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
listed = {v["name"]: v for v in response.json()["data"]["views"] if v["name"] in owned_names}
assert set(listed) == {"Critical Staging", "Disabled"}
assert listed["Critical Staging"]["data"]["query"] == "name CONTAINS 'staging'"
assert (
requests.delete(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{view_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
).status_code
== HTTPStatus.NO_CONTENT
)
response = requests.get(
signoz.self.host_configs["8080"].get(BASE_URL),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert {v["name"] for v in response.json()["data"]["views"] if v["name"] in owned_names} == {"Disabled"}
assert (
requests.delete(
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{view_id}"),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
).status_code
== HTTPStatus.NOT_FOUND
)

View File

@@ -1,134 +0,0 @@
import json
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.traces import TraceIdGenerator, Traces, TracesKind
def test_thread_returns_message_spans_in_order(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
now = datetime.now(tz=UTC).replace(microsecond=0)
trace_id = TraceIdGenerator.trace_id()
root_id, first_llm_id, tool_id, second_llm_id, third_llm_id = (TraceIdGenerator.span_id() for _ in range(5))
resources = {"service.name": "tracedetail-thread"}
first_input = json.dumps([{"role": "user", "parts": [{"type": "text", "content": "weather in Bangalore?"}]}])
first_output = json.dumps([{"role": "assistant", "parts": [{"type": "tool_call", "id": "call_1", "name": "get_weather", "arguments": {"city": "Bangalore"}}], "finish_reason": "tool_call"}])
second_input = json.dumps([{"role": "tool", "content": "sunny", "tool_call_id": "call_1"}])
insert_traces(
[
Traces(timestamp=now - timedelta(seconds=10), duration=timedelta(seconds=9), trace_id=trace_id, span_id=root_id, name="POST /chat", kind=TracesKind.SPAN_KIND_SERVER, resources=resources, attribute_write_mode="json_only"),
Traces(
timestamp=now - timedelta(seconds=8), trace_id=trace_id, span_id=first_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.input.messages": first_input, "gen_ai.output.messages": first_output}, attribute_write_mode="json_only"
),
Traces(timestamp=now - timedelta(seconds=6), trace_id=trace_id, span_id=tool_id, parent_span_id=root_id, name="execute_tool get_weather", resources=resources, attributes={"gen_ai.tool.name": "get_weather"}, attribute_write_mode="json_only"),
Traces(timestamp=now - timedelta(seconds=4), trace_id=trace_id, span_id=second_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.input.messages": second_input}, attribute_write_mode="json_only"),
Traces(timestamp=now - timedelta(seconds=2), trace_id=trace_id, span_id=third_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.output.messages": "It is sunny in Bangalore."}, attribute_write_mode="json_only"),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread"), headers={"Authorization": f"Bearer {token}"}, timeout=10)
assert response.status_code == HTTPStatus.OK, response.text
thread = response.json()["data"]
assert [span["span_id"] for span in thread["spans"]] == [first_llm_id, second_llm_id, third_llm_id]
assert "nextCursor" not in thread
first, input_only, output_only = thread["spans"]
assert first["time_unix"] == int((now - timedelta(seconds=8)).timestamp() * 1000)
assert first["attributes"]["gen_ai.input.messages"] == first_input
assert first["attributes"]["gen_ai.request.model"] == "gpt-4o"
assert first["formatted_input"] == [{"role": "user", "content": [{"type": "text", "content": "weather in Bangalore?"}]}]
assert first["formatted_output"] == [
{
"role": "assistant",
"content": [{"type": "tool_call", "id": "call_1", "name": "get_weather", "arguments": {"city": "Bangalore"}}],
"finishReason": "tool_call",
}
]
assert input_only["formatted_input"] == [{"role": "tool", "content": [{"type": "tool_result", "toolCallId": "call_1", "content": "sunny"}]}]
assert "formatted_output" not in input_only
assert "formatted_input" not in output_only
assert output_only["formatted_output"] == [{"content": [{"type": "generic", "content": "It is sunny in Bangalore."}]}]
def test_thread_paginates_with_cursor(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
now = datetime.now(tz=UTC).replace(microsecond=0)
trace_id = TraceIdGenerator.trace_id()
span_ids = [TraceIdGenerator.span_id() for _ in range(3)]
# identical timestamps on the last two exercise the span_id tie-break
timestamps = [now - timedelta(seconds=6), now - timedelta(seconds=3), now - timedelta(seconds=3)]
insert_traces(
[
Traces(timestamp=timestamp, trace_id=trace_id, span_id=span_id, name="chat gpt-4o", resources={"service.name": "tracedetail-thread-pages"}, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": span_id}])}, attribute_write_mode="json_only")
for span_id, timestamp in zip(span_ids, timestamps, strict=True)
]
)
expected_order = [span_ids[0], *sorted(span_ids[1:])]
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
url = signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread")
headers = {"Authorization": f"Bearer {token}"}
first_page = requests.get(url, params={"limit": 2}, headers=headers, timeout=10)
assert first_page.status_code == HTTPStatus.OK, first_page.text
first = first_page.json()["data"]
assert [span["span_id"] for span in first["spans"]] == expected_order[:2]
assert first["nextCursor"]
second_page = requests.get(url, params={"limit": 2, "cursor": first["nextCursor"]}, headers=headers, timeout=10)
assert second_page.status_code == HTTPStatus.OK, second_page.text
second = second_page.json()["data"]
assert [span["span_id"] for span in second["spans"]] == expected_order[2:]
assert "nextCursor" not in second
def test_thread_without_messages_is_empty(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
trace_id = TraceIdGenerator.trace_id()
insert_traces([Traces(timestamp=datetime.now(tz=UTC) - timedelta(seconds=5), trace_id=trace_id, span_id=TraceIdGenerator.span_id(), name="GET /health", resources={"service.name": "tracedetail-thread-empty"}, attribute_write_mode="json_only")])
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread"), headers={"Authorization": f"Bearer {token}"}, timeout=10)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"] == {"spans": []}
def test_thread_rejects_invalid_requests(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
trace_id = TraceIdGenerator.trace_id()
insert_traces([Traces(timestamp=datetime.now(tz=UTC) - timedelta(seconds=5), trace_id=trace_id, span_id=TraceIdGenerator.span_id(), name="chat gpt-4o", resources={"service.name": "tracedetail-thread-invalid"}, attributes={"gen_ai.input.messages": "hi"}, attribute_write_mode="json_only")])
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
headers = {"Authorization": f"Bearer {token}"}
url = signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread")
for params in ({"limit": -1}, {"limit": 1001}, {"cursor": "not-a-cursor"}):
response = requests.get(url, params=params, headers=headers, timeout=10)
assert response.status_code == HTTPStatus.BAD_REQUEST, f"{params}: {response.text}"
missing = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{TraceIdGenerator.trace_id()}/thread"), headers=headers, timeout=10)
assert missing.status_code == HTTPStatus.NOT_FOUND, missing.text