mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-28 14:20:42 +01:00
Compare commits
64 Commits
feat/ai-tr
...
feat/alert
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
34482d8cc1 | ||
|
|
a541d2687c | ||
|
|
6ab70762a7 | ||
|
|
8989e0a3a1 | ||
|
|
0d150554c7 | ||
|
|
bee75c8e34 | ||
|
|
eb54542ca1 | ||
|
|
c92fe4969a | ||
|
|
97a2fefe61 | ||
|
|
c059d3b8f6 | ||
|
|
37d342eb90 | ||
|
|
d60da7b8a7 | ||
|
|
8d5111bc7d | ||
|
|
e2578c15fb | ||
|
|
a446a832ae | ||
|
|
c89c19ce9e | ||
|
|
e543b3ef32 | ||
|
|
3f7531f444 | ||
|
|
abfa9bee1e | ||
|
|
2f497108df | ||
|
|
1fa1e292e9 | ||
|
|
a328988a15 | ||
|
|
7e53bd607d | ||
|
|
4d82eb8436 | ||
|
|
95c11d0550 | ||
|
|
b2810ec117 | ||
|
|
7556dd2868 | ||
|
|
339a10217e | ||
|
|
cd0eb62734 | ||
|
|
a667007991 | ||
|
|
208fc1e8ec | ||
|
|
54eb67392a | ||
|
|
0e0d97c16f | ||
|
|
659f9c7ecf | ||
|
|
ce8442c17a | ||
|
|
7ab21678b5 | ||
|
|
951bfb66cd | ||
|
|
113685fbdb | ||
|
|
b9c306dd26 | ||
|
|
7e71d4507c | ||
|
|
e13cb08097 | ||
|
|
1e7aaa9f14 | ||
|
|
09ab9f4081 | ||
|
|
f945b5f513 | ||
|
|
8b0ab0ff26 | ||
|
|
f2679a6866 | ||
|
|
253a117849 | ||
|
|
c9538be38f | ||
|
|
be51c37317 | ||
|
|
66796d1767 | ||
|
|
8e19b17855 | ||
|
|
8885497e14 | ||
|
|
626e047485 | ||
|
|
cc3d86c3ec | ||
|
|
0100083284 | ||
|
|
648e945471 | ||
|
|
7323d5ae0b | ||
|
|
578b9172ba | ||
|
|
2110006d17 | ||
|
|
01dda19d17 | ||
|
|
b519574b89 | ||
|
|
629ecabce7 | ||
|
|
131e302c55 | ||
|
|
239c92ba67 |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -68,7 +68,6 @@ jobs:
|
||||
- semconvfamilies
|
||||
- serviceaccount
|
||||
- spanmapper
|
||||
- tracedetail
|
||||
- querier_json_body
|
||||
- querier_skip_resource_fingerprint
|
||||
- ttl
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"},
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"`)
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
93
pkg/ruler/rulestore/sqlrulestore/rule_view.go
Normal file
93
pkg/ruler/rulestore/sqlrulestore/rule_view.go
Normal 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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -257,6 +257,7 @@ func NewSQLMigrationProviderFactories(
|
||||
sqlmigration.NewAddCloudIntegrationTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddNotificationChannelTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddAIObservabilityQuickFiltersFactory(sqlstore),
|
||||
sqlmigration.NewAddRuleViewFactory(sqlstore, sqlschema),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
78
pkg/sqlmigration/131_add_rule_view.go
Normal file
78
pkg/sqlmigration/131_add_rule_view.go
Normal 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
|
||||
}
|
||||
@@ -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 ""
|
||||
}
|
||||
@@ -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)}}
|
||||
}
|
||||
@@ -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))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
169
pkg/types/ruletypes/rule_view.go
Normal file
169
pkg/types/ruletypes/rule_view.go
Normal 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
|
||||
}
|
||||
210
pkg/types/ruletypes/rule_view_test.go
Normal file
210
pkg/types/ruletypes/rule_view_test.go
Normal 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))
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
30
tests/fixtures/alerts.py
vendored
30
tests/fixtures/alerts.py
vendored
@@ -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 []}
|
||||
|
||||
269
tests/integration/tests/ruler/06_rule_views.py
Normal file
269
tests/integration/tests/ruler/06_rule_views.py
Normal 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
|
||||
)
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user