Compare commits

...

2 Commits

Author SHA1 Message Date
Nityananda Gohain
ea8f95ee08 chore: enable ai 011y processors by default (#12912)
Some checks are pending
build-staging / js-build (push) Blocked by required conditions
build-staging / go-build (push) Blocked by required conditions
build-staging / prepare (push) Waiting to run
build-staging / staging (push) Blocked by required conditions
cacheci / tests (push) Waiting to run
Release Drafter / update_release_draft (push) Waiting to run
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description
Enable the processors by default so that metadata is populated

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR
No issue
2026-09-19 13:02:43 +00:00
Nityananda Gohain
b64116d67d feat: support for default attribute mapping (#12809)
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description

Until now every org started with no span mapper groups and we had to
create `llm`, `agent` and `tool` by hand. This PR ships them as
defaults.

- The three groups live as JSON in the binary. On startup and when a new
org is created, they get seeded. When we change a definition and bump
its version, the next release updates the org's copy in place.
- Anything SigNoz ships is marked `origin: system` and can only be
switched on or off. Users can't rename or delete these groups and
mappers, and can't take their names. Anything the user adds is theirs to
edit or remove, including new mappers in a shipped group or new sources
on a shipped mapper.
  - Upgrades keep the user's on/off choices and never touch their items.
- Condition substrings and sources now carry `enabled` and `origin`, so
a substring is an object instead of a plain string. Disabled ones are
left out of the collector config.
- - Migration 127 only adds the `origin` and `version` columns. Stored
JSON is not rewritten because these tables are empty on every instance.
- Fixes creating a group or mapper with `enabled: false` being saved as
true (the bun `default:true` tag turned false into SQL `DEFAULT`)
  
Frontend: no UI changes. Generated client regenerated; drafts carry
`enabled` and `origin` so saves from the existing screens round-trip
shipped items intact.

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


<!--Anything reviewers should keep in mind while reviewing -->
#### Additional Information
* Frontend follow-up: toggles for sources and substrings, "Default"
badge and read-only rows for shipped items, hide rename/delete on system
groups and mappers.
* Deferred: per-group upgrade changelog ( will come back to this later)

---------

Co-authored-by: Gaurav Tewari <gauravtewari111@gmail.com>
Co-authored-by: Gaurav Tewari <tewarig@users.noreply.github.com>
2026-09-19 11:22:46 +00:00
40 changed files with 2200 additions and 206 deletions

View File

@@ -9732,6 +9732,8 @@ components:
type: string
name:
type: string
origin:
$ref: '#/components/schemas/SpantypesSpanMapperOrigin'
updatedAt:
format: date-time
type: string
@@ -9744,6 +9746,7 @@ components:
- fieldContext
- config
- enabled
- origin
type: object
SpantypesSpanMapperConfig:
properties:
@@ -9772,48 +9775,75 @@ components:
type: string
orgId:
type: string
origin:
$ref: '#/components/schemas/SpantypesSpanMapperOrigin'
updatedAt:
format: date-time
type: string
updatedBy:
type: string
version:
type: integer
required:
- id
- orgId
- name
- condition
- enabled
- origin
- version
type: object
SpantypesSpanMapperGroupCondition:
nullable: true
properties:
attributes:
items:
type: string
$ref: '#/components/schemas/SpantypesSpanMapperGroupConditionKey'
nullable: true
type: array
resource:
items:
type: string
$ref: '#/components/schemas/SpantypesSpanMapperGroupConditionKey'
nullable: true
type: array
required:
- attributes
- resource
type: object
SpantypesSpanMapperGroupConditionKey:
properties:
enabled:
type: boolean
origin:
$ref: '#/components/schemas/SpantypesSpanMapperOrigin'
value:
type: string
required:
- value
- enabled
type: object
SpantypesSpanMapperOperation:
enum:
- move
- copy
type: string
SpantypesSpanMapperOrigin:
enum:
- user
- system
type: string
SpantypesSpanMapperSource:
properties:
context:
$ref: '#/components/schemas/SpantypesFieldContext'
enabled:
type: boolean
key:
type: string
operation:
$ref: '#/components/schemas/SpantypesSpanMapperOperation'
origin:
$ref: '#/components/schemas/SpantypesSpanMapperOrigin'
priority:
type: integer
required:
@@ -9821,6 +9851,7 @@ components:
- context
- operation
- priority
- enabled
type: object
SpantypesSpanMapperTestSpan:
properties:

View File

@@ -10817,6 +10817,22 @@ export interface SpantypesGettableFlamegraphTraceDTO {
startTimestampMillis: number;
}
export enum SpantypesSpanMapperOriginDTO {
user = 'user',
system = 'system',
}
export interface SpantypesSpanMapperGroupConditionKeyDTO {
/**
* @type boolean
*/
enabled: boolean;
origin?: SpantypesSpanMapperOriginDTO;
/**
* @type string
*/
value: string;
}
/**
* @nullable
*/
@@ -10824,11 +10840,11 @@ export type SpantypesSpanMapperGroupConditionDTO = {
/**
* @type array,null
*/
attributes: string[] | null;
attributes: SpantypesSpanMapperGroupConditionKeyDTO[] | null;
/**
* @type array,null
*/
resource: string[] | null;
resource: SpantypesSpanMapperGroupConditionKeyDTO[] | null;
} | null;
export interface SpantypesSpanMapperGroupDTO {
@@ -10858,6 +10874,7 @@ export interface SpantypesSpanMapperGroupDTO {
* @type string
*/
orgId: string;
origin: SpantypesSpanMapperOriginDTO;
/**
* @type string
* @format date-time
@@ -10867,6 +10884,10 @@ export interface SpantypesSpanMapperGroupDTO {
* @type string
*/
updatedBy?: string;
/**
* @type integer
*/
version: number;
}
export interface SpantypesGettableSpanMapperGroupsDTO {
@@ -10924,11 +10945,16 @@ export enum SpantypesSpanMapperOperationDTO {
}
export interface SpantypesSpanMapperSourceDTO {
context: SpantypesFieldContextDTO;
/**
* @type boolean
*/
enabled: boolean;
/**
* @type string
*/
key: string;
operation: SpantypesSpanMapperOperationDTO;
origin?: SpantypesSpanMapperOriginDTO;
/**
* @type integer
*/
@@ -10970,6 +10996,7 @@ export interface SpantypesSpanMapperDTO {
* @type string
*/
name: string;
origin: SpantypesSpanMapperOriginDTO;
/**
* @type string
* @format date-time

View File

@@ -333,6 +333,7 @@ describe('AttributeMappingsTab (integration)', () => {
context: FieldContext.attribute,
operation: MapperOperation.copy,
priority,
enabled: true,
})),
},
}),

View File

@@ -1,10 +1,12 @@
import { Typography } from '@signozhq/ui/typography';
import { ConditionKey } from 'container/LLMObservability/AttributeMapping/types';
import styles from './ConditionsTooltip.module.scss';
interface ConditionsTooltipProps {
attributes: string[];
resource: string[];
attributes: ConditionKey[];
resource: ConditionKey[];
}
function ConditionsTooltip({
@@ -33,8 +35,8 @@ function ConditionsTooltip({
</Typography.Text>
<div className={styles.keyList}>
{attributes.map((key) => (
<code key={key} className={styles.key}>
{key}
<code key={`${key.origin}-${key.value}`} className={styles.key}>
{key.value}
</code>
))}
</div>
@@ -47,8 +49,8 @@ function ConditionsTooltip({
</Typography.Text>
<div className={styles.keyList}>
{resource.map((key) => (
<code key={key} className={styles.key}>
{key}
<code key={`${key.origin}-${key.value}`} className={styles.key}>
{key.value}
</code>
))}
</div>

View File

@@ -3,6 +3,7 @@ import {
SpantypesSpanMapperDTO as Mapper,
SpantypesSpanMapperGroupDTO as MapperGroup,
SpantypesSpanMapperOperationDTO as MapperOperation,
SpantypesSpanMapperOriginDTO as MapperOrigin,
SpantypesSpanMapperTestSpanDTO as TestSpan,
} from 'api/generated/services/sigNoz.schemas';
@@ -21,9 +22,15 @@ export function makeGroup(overrides: Partial<MapperGroup> = {}): MapperGroup {
orgId: 'org-1',
name: 'demo',
enabled: true,
origin: MapperOrigin.user,
version: 0,
condition: {
attributes: ['ai.embeddings'],
resource: ['cloud.account.id'],
attributes: [
{ value: 'ai.embeddings', enabled: true, origin: MapperOrigin.user },
],
resource: [
{ value: 'cloud.account.id', enabled: true, origin: MapperOrigin.user },
],
},
...overrides,
};
@@ -35,6 +42,7 @@ export function makeMapper(overrides: Partial<Mapper> = {}): Mapper {
groupId: 'group-1',
name: 'gen_ai.request.model',
enabled: true,
origin: MapperOrigin.user,
fieldContext: FieldContext.attribute,
config: {
sources: [
@@ -43,12 +51,16 @@ export function makeMapper(overrides: Partial<Mapper> = {}): Mapper {
context: FieldContext.attribute,
operation: MapperOperation.copy,
priority: 2,
enabled: true,
origin: MapperOrigin.user,
},
{
key: 'llm.model',
context: FieldContext.attribute,
operation: MapperOperation.move,
priority: 1,
enabled: true,
origin: MapperOrigin.user,
},
],
},
@@ -85,8 +97,12 @@ export const mockGroups: MapperGroup[] = [
id: 'group-1',
name: 'demo',
condition: {
attributes: ['ai.embeddings'],
resource: ['cloud.account.id'],
attributes: [
{ value: 'ai.embeddings', enabled: true, origin: MapperOrigin.user },
],
resource: [
{ value: 'cloud.account.id', enabled: true, origin: MapperOrigin.user },
],
},
}),
makeGroup({

View File

@@ -1,19 +1,23 @@
import { Button } from '@signozhq/ui/button';
import { Plus, X } from '@signozhq/icons';
import { FieldContextValue } from 'container/LLMObservability/AttributeMapping/types';
import {
ConditionKey,
FieldContextValue,
} from 'container/LLMObservability/AttributeMapping/types';
import { createConditionKey } from 'container/LLMObservability/AttributeMapping/utils';
import KeySearchInput from '../../../KeySearchInput/KeySearchInput';
import styles from './ConditionKeyList.module.scss';
interface ConditionKeyListProps {
label: string;
labelHint?: string;
keys: string[];
keys: ConditionKey[];
placeholder: string;
addLabel: string;
testIdPrefix: string;
fieldContext: FieldContextValue;
onChange: (keys: string[]) => void;
onChange: (keys: ConditionKey[]) => void;
}
function ConditionKeyList({
@@ -27,11 +31,11 @@ function ConditionKeyList({
onChange,
}: ConditionKeyListProps): JSX.Element {
const updateKey = (index: number, value: string): void => {
onChange(keys.map((key, i) => (i === index ? value : key)));
onChange(keys.map((key, i) => (i === index ? { ...key, value } : key)));
};
const addKey = (): void => {
onChange([...keys, '']);
onChange([...keys, createConditionKey()]);
};
const removeKey = (index: number): void => {
@@ -53,7 +57,7 @@ function ConditionKeyList({
<KeySearchInput
className={styles.keyInput}
placeholder={placeholder}
value={key}
value={key.value}
fieldContext={fieldContext}
onChange={(next): void => updateKey(index, next)}
testId={`${testIdPrefix}-${index}`}

View File

@@ -42,7 +42,9 @@ function sourcesEqual(a: SourceConfig[], b: SourceConfig[]): boolean {
(source, index) =>
source.key === b[index].key &&
source.context === b[index].context &&
source.operation === b[index].operation,
source.operation === b[index].operation &&
source.enabled === b[index].enabled &&
source.origin === b[index].origin,
)
);
}

View File

@@ -1,8 +1,11 @@
import {
SpantypesFieldContextDTO,
SpantypesSpanMapperDTO,
SpantypesSpanMapperGroupConditionKeyDTO,
SpantypesSpanMapperGroupDTO,
SpantypesSpanMapperOperationDTO,
SpantypesSpanMapperOriginDTO,
SpantypesSpanMapperSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
export type MapperGroup = SpantypesSpanMapperGroupDTO;
@@ -11,21 +14,23 @@ export const FieldContext = SpantypesFieldContextDTO;
export type FieldContextValue = SpantypesFieldContextDTO;
export const MapperOperation = SpantypesSpanMapperOperationDTO;
export type MapperOperationValue = SpantypesSpanMapperOperationDTO;
export const MapperOrigin = SpantypesSpanMapperOriginDTO;
export type MapperOriginValue = SpantypesSpanMapperOriginDTO;
export type ConditionKey = SpantypesSpanMapperGroupConditionKeyDTO;
export type MapperDraftMode = 'add' | 'edit';
export interface SourceConfig {
key: string;
context: SpantypesFieldContextDTO;
operation: SpantypesSpanMapperOperationDTO;
}
// `priority` is left out: it is derived from list order when the draft is
// serialized.
export type SourceConfig = Omit<SpantypesSpanMapperSourceDTO, 'priority'>;
// Editable form state for a mapper. `sources` is ordered highest priority
// first; `fieldContext` is where the standardized target is written.
export interface MapperDraft {
id: string | null;
name: string;
fieldContext: SpantypesFieldContextDTO;
fieldContext: FieldContextValue;
sources: SourceConfig[];
enabled: boolean;
}
@@ -33,26 +38,20 @@ export interface MapperDraft {
export interface GroupDraft {
id: string | null;
name: string;
attributes: string[];
resource: string[];
attributes: ConditionKey[];
resource: ConditionKey[];
enabled: boolean;
}
export interface DraftMapper {
// The editor tree identifies rows by `localId` so unsaved ones are addressable;
// `serverId` is null until the row has been persisted.
export type DraftMapper = Omit<MapperDraft, 'id'> & {
localId: string;
serverId: string | null;
name: string;
fieldContext: SpantypesFieldContextDTO;
sources: SourceConfig[];
enabled: boolean;
}
};
export interface DraftGroup {
export type DraftGroup = Omit<GroupDraft, 'id'> & {
localId: string;
serverId: string | null;
name: string;
attributes: string[];
resource: string[];
enabled: boolean;
mappers: DraftMapper[];
}
};

View File

@@ -1,12 +1,14 @@
import {
SpantypesPostableSpanMapperDTO,
SpantypesPostableSpanMapperGroupDTO,
SpantypesSpanMapperGroupConditionKeyDTO,
SpantypesUpdatableSpanMapperDTO,
SpantypesUpdatableSpanMapperGroupDTO,
} from 'api/generated/services/sigNoz.schemas';
import { v4 as uuid } from 'uuid';
import {
ConditionKey,
DraftGroup,
DraftMapper,
FieldContext,
@@ -15,6 +17,7 @@ import {
MapperDraft,
MapperGroup,
MapperOperation,
MapperOrigin,
SourceConfig,
} from './types';
@@ -24,20 +27,36 @@ function genLocalId(prefix: 'group' | 'mapper'): string {
return `local-${prefix}-${uuid()}`;
}
// Trimmed, de-duplicated, non-empty keys preserving input order.
function cleanKeys(keys: string[]): string[] {
export function createConditionKey(value = ''): ConditionKey {
return { value, enabled: true, origin: MapperOrigin.user };
}
// Trimmed, de-duplicated, non-empty keys preserving input order. A shipped and
// a user key may share a value, so the origin is part of the identity.
function cleanKeys(keys: ConditionKey[]): ConditionKey[] {
const seen = new Set<string>();
const result: string[] = [];
const result: ConditionKey[] = [];
keys.forEach((raw) => {
const key = raw.trim();
if (key && !seen.has(key)) {
seen.add(key);
result.push(key);
const value = raw.value.trim();
const dedupeKey = `${raw.origin}:${value}`;
if (value && !seen.has(dedupeKey)) {
seen.add(dedupeKey);
result.push({ ...raw, value });
}
});
return result;
}
function fromConditionKeys(
keys: SpantypesSpanMapperGroupConditionKeyDTO[] | null | undefined,
): ConditionKey[] {
return (keys ?? []).map((key) => ({
value: key.value,
enabled: key.enabled,
origin: key.origin ?? MapperOrigin.user,
}));
}
// Source configs for a mapper, highest priority first (first match wins at
// evaluation time).
function getMapperSources(mapper: Mapper): SourceConfig[] {
@@ -48,6 +67,8 @@ function getMapperSources(mapper: Mapper): SourceConfig[] {
key: source.key,
context: source.context,
operation: source.operation,
enabled: source.enabled,
origin: source.origin ?? MapperOrigin.user,
}));
}
@@ -56,6 +77,8 @@ export function createEmptySource(): SourceConfig {
key: '',
context: FieldContext.attribute,
operation: MapperOperation.copy,
enabled: true,
origin: MapperOrigin.user,
};
}
@@ -72,7 +95,7 @@ function getCleanSources(draft: MapperDraft): SourceConfig[] {
const result: SourceConfig[] = [];
draft.sources.forEach((source) => {
const key = source.key.trim();
const dedupeKey = `${source.context}:${key}`;
const dedupeKey = `${source.origin}:${source.context}:${key}`;
if (key && !seen.has(dedupeKey)) {
seen.add(dedupeKey);
result.push({ ...source, key });
@@ -95,6 +118,8 @@ function buildSources(
context: source.context,
operation: source.operation,
priority: sources.length - index,
enabled: source.enabled,
origin: source.origin,
}));
}
@@ -123,7 +148,7 @@ export function buildUpdatableMapper(
export const EMPTY_GROUP_DRAFT: GroupDraft = {
id: null,
name: '',
attributes: [''],
attributes: [createConditionKey()],
resource: [],
enabled: true,
};
@@ -170,8 +195,8 @@ export function buildDraftGroup(
localId: group.id,
serverId: group.id,
name: group.name,
attributes: group.condition?.attributes ?? [],
resource: group.condition?.resource ?? [],
attributes: fromConditionKeys(group.condition?.attributes),
resource: fromConditionKeys(group.condition?.resource),
enabled: group.enabled,
mappers: mappers.map(buildDraftMapper),
};
@@ -182,7 +207,8 @@ export function groupDraftFromNode(group: DraftGroup): GroupDraft {
return {
id: group.localId,
name: group.name,
attributes: group.attributes.length > 0 ? group.attributes : [''],
attributes:
group.attributes.length > 0 ? group.attributes : [createConditionKey()],
resource: group.resource,
enabled: group.enabled,
};

View File

@@ -7,12 +7,10 @@ import (
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/modules/llmpricingrule"
"github.com/SigNoz/signoz/pkg/querier"
"github.com/SigNoz/signoz/pkg/query-service/agentConf"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
"github.com/SigNoz/signoz/pkg/types/llmpricingruletypes"
"github.com/SigNoz/signoz/pkg/types/opamptypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
@@ -26,11 +24,10 @@ const unmappedModelsLookback = time.Hour
type module struct {
store llmpricingruletypes.Store
querier querier.Querier
flagger flagger.Flagger
}
func NewModule(store llmpricingruletypes.Store, flagger flagger.Flagger, querier querier.Querier) llmpricingrule.Module {
return &module{store: store, flagger: flagger, querier: querier}
func NewModule(store llmpricingruletypes.Store, querier querier.Querier) llmpricingrule.Module {
return &module{store: store, querier: querier}
}
func (module *module) List(ctx context.Context, orgID valuer.UUID, offset, limit int, search string, isOverride *bool) ([]*llmpricingruletypes.LLMPricingRule, int, error) {
@@ -123,16 +120,6 @@ func (module *module) AgentFeatureType() agentConf.AgentFeatureType {
func (module *module) RecommendAgentConfig(orgID valuer.UUID, currentConfYaml []byte, configVersion *opamptypes.AgentConfigVersion) ([]byte, string, error) {
ctx := context.Background()
// Skip the llm pricing processor unless AI observability is enabled for the org.
evalCtx := featuretypes.NewFlaggerEvaluationContext(orgID)
enabled, err := module.flagger.Boolean(ctx, flagger.FeatureEnableAIObservability, evalCtx)
if err != nil {
return nil, "", err
}
if !enabled {
return currentConfYaml, "", nil
}
rules, err := module.getEnabledRules(ctx, orgID)
if err != nil {
return nil, "", err

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/dashboard"
"github.com/SigNoz/signoz/pkg/modules/organization"
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
"github.com/SigNoz/signoz/pkg/modules/spanmapper"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/valuer"
)
@@ -16,10 +17,11 @@ type setter struct {
alertmanager alertmanager.Alertmanager
quickfilter quickfilter.Module
dashboard dashboard.Module
spanMapper spanmapper.Module
}
func NewSetter(store types.OrganizationStore, alertmanager alertmanager.Alertmanager, quickfilter quickfilter.Module, dashboard dashboard.Module) organization.Setter {
return &setter{store: store, alertmanager: alertmanager, quickfilter: quickfilter, dashboard: dashboard}
func NewSetter(store types.OrganizationStore, alertmanager alertmanager.Alertmanager, quickfilter quickfilter.Module, dashboard dashboard.Module, spanMapper spanmapper.Module) organization.Setter {
return &setter{store: store, alertmanager: alertmanager, quickfilter: quickfilter, dashboard: dashboard, spanMapper: spanMapper}
}
func (module *setter) Create(ctx context.Context, organization *types.Organization, createManagedRoles func(context.Context, valuer.UUID) error) error {
@@ -43,6 +45,10 @@ func (module *setter) Create(ctx context.Context, organization *types.Organizati
return err
}
if err := module.spanMapper.ReconcileSystemGroups(ctx, organization.ID); err != nil {
return err
}
return nil
}

View File

@@ -0,0 +1,46 @@
package implspanmapper
import (
"embed"
"io/fs"
"path"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/spantypes"
)
const definitionsRoot = "fs/definitions"
//go:embed fs/definitions/*.json
var definitionFiles embed.FS
// NewSystemGroupRegistry parses every embedded definition. Definitions are
// build-time assets validated by a test, so a failure here means the binary
// shipped broken JSON.
func NewSystemGroupRegistry() (spantypes.SpanMapperGroupRegistry, error) {
entries, err := fs.ReadDir(definitionFiles, definitionsRoot)
if err != nil {
return spantypes.SpanMapperGroupRegistry{}, errors.WrapInternalf(err, errors.CodeInternal, "couldn't read span mapper group definitions")
}
definitions := make([]spantypes.SpanMapperGroupDefinition, 0, len(entries))
for _, entry := range entries {
if entry.IsDir() {
continue
}
file := path.Join(definitionsRoot, entry.Name())
raw, err := definitionFiles.ReadFile(file)
if err != nil {
return spantypes.SpanMapperGroupRegistry{}, errors.WrapInternalf(err, errors.CodeInternal, "couldn't read %s", file)
}
definition, err := spantypes.NewSpanMapperGroupDefinition(raw)
if err != nil {
return spantypes.SpanMapperGroupRegistry{}, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "couldn't parse %s", file)
}
definitions = append(definitions, definition)
}
return spantypes.NewSpanMapperGroupRegistry(definitions)
}

View File

@@ -0,0 +1,79 @@
{
"version": 1,
"definition": {
"name": "gen_ai.agent",
"condition": {
"attributes": [
{
"value": "agent"
}
],
"resource": []
},
"enabled": true,
"mappers": [
{
"name": "gen_ai.agent.name",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "agent.name",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "agent_name",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.agent.id",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "agent.id",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.agent.description",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "agent.description",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.output.messages",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "final_result",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
}
]
}
}

View File

@@ -0,0 +1,347 @@
{
"version": 1,
"definition": {
"name": "gen_ai.llm",
"condition": {
"attributes": [
{
"value": "model"
}
],
"resource": []
},
"enabled": true,
"mappers": [
{
"name": "gen_ai.request.model",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "llm.model_name",
"context": "attribute",
"operation": "copy",
"priority": 60
},
{
"key": "llm.request.model",
"context": "attribute",
"operation": "copy",
"priority": 50
},
{
"key": "ai.model.id",
"context": "attribute",
"operation": "copy",
"priority": 40
},
{
"key": "langfuse.observation.model.name",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "embedding.model_name",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "model",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.response.model",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "llm.response.model",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.response.model",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.provider.name",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.system",
"context": "attribute",
"operation": "copy",
"priority": 50
},
{
"key": "llm.vendor",
"context": "attribute",
"operation": "copy",
"priority": 40
},
{
"key": "llm.provider",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "llm.system",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.model.provider",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.operation.name",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "llm.request.type",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.usage.input_tokens",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.usage.prompt_tokens",
"context": "attribute",
"operation": "copy",
"priority": 50
},
{
"key": "llm.usage.prompt_tokens",
"context": "attribute",
"operation": "copy",
"priority": 40
},
{
"key": "llm.token_count.prompt",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "ai.usage.inputTokens",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.usage.promptTokens",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.usage.output_tokens",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.usage.completion_tokens",
"context": "attribute",
"operation": "copy",
"priority": 50
},
{
"key": "llm.usage.completion_tokens",
"context": "attribute",
"operation": "copy",
"priority": 40
},
{
"key": "llm.token_count.completion",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "ai.usage.outputTokens",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.usage.completionTokens",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.usage.cache_read.input_tokens",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.usage.cache_read_input_tokens",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "llm.token_count.prompt_details.cache_read",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.usage.cachedInputTokens",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.usage.cache_creation.input_tokens",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.usage.cache_write.input_tokens",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "gen_ai.usage.cache_creation_input_tokens",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "llm.token_count.prompt_details.cache_write",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.input.messages",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.prompt",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "ai.prompt.messages",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "input.value",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.output.messages",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "gen_ai.completion",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "ai.response.text",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "output.value",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.conversation.id",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "session.id",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "langfuse.session.id",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.response.finish_reason",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "ai.response.finishReason",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
}
]
}
}

View File

@@ -0,0 +1,135 @@
{
"version": 1,
"definition": {
"name": "gen_ai.tool",
"condition": {
"attributes": [
{
"value": "tool"
}
],
"resource": []
},
"enabled": true,
"mappers": [
{
"name": "gen_ai.tool.name",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "tool.name",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.toolCall.name",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.tool.call.id",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "tool.id",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "ai.toolCall.id",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.tool.description",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "tool.description",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.tool.call.arguments",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "ai.toolCall.args",
"context": "attribute",
"operation": "copy",
"priority": 40
},
{
"key": "traceloop.entity.input",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "gcp.vertex.agent.tool_call_args",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "input.value",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
},
{
"name": "gen_ai.tool.call.result",
"fieldContext": "attribute",
"config": {
"sources": [
{
"key": "ai.toolCall.result",
"context": "attribute",
"operation": "copy",
"priority": 40
},
{
"key": "traceloop.entity.output",
"context": "attribute",
"operation": "copy",
"priority": 30
},
{
"key": "gcp.vertex.agent.tool_response",
"context": "attribute",
"operation": "copy",
"priority": 20
},
{
"key": "output.value",
"context": "attribute",
"operation": "copy",
"priority": 10
}
]
}
}
]
}
}

View File

@@ -69,6 +69,10 @@ func (h *handler) CreateGroup(rw http.ResponseWriter, r *http.Request) {
render.Error(rw, err)
return
}
if err := req.Validate(); err != nil {
render.Error(rw, err)
return
}
group := spantypes.NewSpanMapperGroup(orgID, claims.Email, req)
@@ -191,6 +195,10 @@ func (h *handler) CreateMapper(rw http.ResponseWriter, r *http.Request) {
render.Error(rw, err)
return
}
if err := req.Validate(); err != nil {
render.Error(rw, err)
return
}
mapper := spantypes.NewSpanMapper(groupID, claims.Email, req)
if err := h.module.CreateMapper(ctx, orgID, groupID, mapper); err != nil {

View File

@@ -5,22 +5,30 @@ import (
"encoding/json"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/modules/spanmapper"
"github.com/SigNoz/signoz/pkg/query-service/agentConf"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
"github.com/SigNoz/signoz/pkg/types/opamptypes"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
// maxTestSpans bounds the input size: every test request boots a full
// in-memory collector pipeline and is reachable with viewer access.
const maxTestSpans = 100
type module struct {
store spantypes.SpanMapperStore
flagger flagger.Flagger
store spantypes.SpanMapperStore
registry spantypes.SpanMapperGroupRegistry
settings factory.ScopedProviderSettings
}
func NewModule(store spantypes.SpanMapperStore, flagger flagger.Flagger) spanmapper.Module {
return &module{store: store, flagger: flagger}
func NewModule(store spantypes.SpanMapperStore, registry spantypes.SpanMapperGroupRegistry, providerSettings factory.ProviderSettings) spanmapper.Module {
return &module{
store: store,
registry: registry,
settings: factory.NewScopedProviderSettings(providerSettings, "github.com/SigNoz/signoz/pkg/modules/spanmapper/implspanmapper"),
}
}
func (module *module) ListGroups(ctx context.Context, orgID valuer.UUID, q *spantypes.ListSpanMapperGroupsQuery) ([]*spantypes.SpanMapperGroup, error) {
@@ -32,6 +40,9 @@ func (module *module) GetGroup(ctx context.Context, orgID, id valuer.UUID) (*spa
}
func (module *module) CreateGroup(ctx context.Context, orgID valuer.UUID, group *spantypes.SpanMapperGroup) error {
if module.registry.IsReserved(group.Name) {
return errors.Newf(errors.TypeInvalidInput, spantypes.ErrCodeMappingGroupNameReserved, "group name %q is reserved for a default group", group.Name)
}
return module.store.CreateGroup(ctx, group)
}
@@ -40,10 +51,14 @@ func (module *module) UpdateGroup(ctx context.Context, orgID, id valuer.UUID, na
if err != nil {
return err
}
group.Update(name, condition, enabled, updatedBy)
if name != nil && *name != group.Name && module.registry.IsReserved(*name) {
return errors.Newf(errors.TypeInvalidInput, spantypes.ErrCodeMappingGroupNameReserved, "group name %q is reserved for a default group", *name)
}
if err := group.Update(name, condition, enabled, updatedBy); err != nil {
return err
}
err = module.store.UpdateGroup(ctx, group)
if err != nil {
if err := module.store.UpdateGroup(ctx, group); err != nil {
return err
}
agentConf.NotifyConfigUpdate(ctx)
@@ -51,8 +66,7 @@ func (module *module) UpdateGroup(ctx context.Context, orgID, id valuer.UUID, na
}
func (module *module) DeleteGroup(ctx context.Context, orgID, id valuer.UUID) error {
err := module.store.DeleteGroup(ctx, orgID, id)
if err != nil {
if err := module.store.DeleteGroup(ctx, orgID, id); err != nil {
return err
}
agentConf.NotifyConfigUpdate(ctx)
@@ -81,14 +95,13 @@ func (module *module) CreateMapper(ctx context.Context, orgID, groupID valuer.UU
}
func (module *module) UpdateMapper(ctx context.Context, orgID, groupID, id valuer.UUID, fieldContext spantypes.FieldContext, config *spantypes.SpanMapperConfig, enabled *bool, updatedBy string) error {
if _, err := module.store.GetGroup(ctx, orgID, groupID); err != nil {
return err
}
mapper, err := module.store.GetMapper(ctx, orgID, groupID, id)
if err != nil {
return err
}
mapper.Update(fieldContext, config, enabled, updatedBy)
if err := mapper.Update(fieldContext, config, enabled, updatedBy); err != nil {
return err
}
err = module.store.UpdateMapper(ctx, mapper)
if err != nil {
return err
@@ -98,18 +111,13 @@ func (module *module) UpdateMapper(ctx context.Context, orgID, groupID, id value
}
func (module *module) DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID) error {
err := module.store.DeleteMapper(ctx, orgID, groupID, id)
if err != nil {
if err := module.store.DeleteMapper(ctx, orgID, groupID, id, spantypes.SpanMapperOriginUser); err != nil {
return err
}
agentConf.NotifyConfigUpdate(ctx)
return nil
}
// maxTestSpans bounds the input size: every test request boots a full
// in-memory collector pipeline and is reachable with viewer access.
const maxTestSpans = 100
func (module *module) TestMappers(ctx context.Context, orgID valuer.UUID, spans []spantypes.SpanMapperTestSpan, groups []*spantypes.SpanMapperGroupWithMappers) ([]spantypes.SpanMapperTestSpan, []string, error) {
if len(spans) == 0 {
return nil, nil, errors.New(errors.TypeInvalidInput, spantypes.ErrCodeMappingInvalidInput, "'spans' must contain at least one span")
@@ -130,6 +138,31 @@ func (module *module) TestMappers(ctx context.Context, orgID valuer.UUID, spans
return out, collectorLogs, nil
}
func (module *module) AgentFeatureType() agentConf.AgentFeatureType {
return spantypes.SpanAttrMappingFeatureType
}
func (module *module) RecommendAgentConfig(orgID valuer.UUID, currentConfYaml []byte, configVersion *opamptypes.AgentConfigVersion) ([]byte, string, error) {
ctx := context.Background()
enabledMappers, err := module.listEnabledGroupsWithMappers(ctx, orgID)
if err != nil {
return nil, "", err
}
updatedConf, err := spantypes.GenerateCollectorConfigWithSpanMapperProcessor(currentConfYaml, enabledMappers)
if err != nil {
return nil, "", err
}
serialized, err := json.Marshal(enabledMappers)
if err != nil {
return nil, "", err
}
return updatedConf, string(serialized), nil
}
// backfillMappers loads saved mappers for any enabled group whose Mappers is
// nil. Disabled groups are skipped: the simulation filters them out anyway,
// so there is no point loading their mappers or failing on their names.
@@ -161,41 +194,6 @@ func (module *module) backfillMappers(ctx context.Context, orgID valuer.UUID, gr
return groups, nil
}
func (module *module) AgentFeatureType() agentConf.AgentFeatureType {
return spantypes.SpanAttrMappingFeatureType
}
func (module *module) RecommendAgentConfig(orgID valuer.UUID, currentConfYaml []byte, configVersion *opamptypes.AgentConfigVersion) ([]byte, string, error) {
ctx := context.Background()
// Skip the llm pricing processor unless AI observability is enabled for the org.
evalCtx := featuretypes.NewFlaggerEvaluationContext(orgID)
enabled, err := module.flagger.Boolean(ctx, flagger.FeatureEnableAIObservability, evalCtx)
if err != nil {
return nil, "", err
}
if !enabled {
return currentConfYaml, "", nil
}
enabledMappers, err := module.listEnabledGroupsWithMappers(ctx, orgID)
if err != nil {
return nil, "", err
}
updatedConf, err := spantypes.GenerateCollectorConfigWithSpanMapperProcessor(currentConfYaml, enabledMappers)
if err != nil {
return nil, "", err
}
serialized, err := json.Marshal(enabled)
if err != nil {
return nil, "", err
}
return updatedConf, string(serialized), nil
}
// listEnabledGroupsWithMappers returns groups with their mappers.
func (module *module) listEnabledGroupsWithMappers(ctx context.Context, orgID valuer.UUID) ([]*spantypes.SpanMapperGroupWithMappers, error) {
enabled := true

View File

@@ -0,0 +1,228 @@
package implspanmapper
import (
"context"
"path/filepath"
"strconv"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/factory/factorytest"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/sqlstore/sqlitesqlstore"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
const testUser = "user@signoz.io"
func newTestSQLStore(t *testing.T) sqlstore.SQLStore {
t.Helper()
store, err := sqlitesqlstore.New(context.Background(), factorytest.NewSettings(), sqlstore.Config{
Provider: "sqlite",
Connection: sqlstore.ConnectionConfig{MaxOpenConns: 10},
Sqlite: sqlstore.SqliteConfig{
Path: filepath.Join(t.TempDir(), "test.db"),
Mode: "wal",
BusyTimeout: 5 * time.Second,
TransactionMode: "deferred",
},
})
require.NoError(t, err)
for _, model := range []any{
(*spantypes.StorableSpanMapperGroup)(nil),
(*spantypes.StorableSpanMapper)(nil),
} {
_, err := store.BunDB().NewCreateTable().Model(model).IfNotExists().Exec(context.Background())
require.NoError(t, err)
}
_, err = store.BunDB().Exec(`CREATE UNIQUE INDEX IF NOT EXISTS uq_span_mapper_group_org_name ON span_mapper_group (org_id, name)`)
require.NoError(t, err)
_, err = store.BunDB().Exec(`CREATE UNIQUE INDEX IF NOT EXISTS uq_span_mapper_group_name ON span_mapper (group_id, name)`)
require.NoError(t, err)
return store
}
func newTestModule(t *testing.T, sqlStore sqlstore.SQLStore, definitions ...spantypes.SpanMapperGroupDefinition) *module {
t.Helper()
registry, err := spantypes.NewSpanMapperGroupRegistry(definitions)
require.NoError(t, err)
return NewModule(NewStore(sqlStore), registry, factorytest.NewSettings()).(*module)
}
func newTestDefinition(t *testing.T, version int, body string) spantypes.SpanMapperGroupDefinition {
t.Helper()
definition, err := spantypes.NewSpanMapperGroupDefinition([]byte(`{"version": ` + strconv.Itoa(version) + `, "definition": ` + body + `}`))
require.NoError(t, err)
return definition
}
// llmV1 ships two mappers; llmV2 renames a source, adds a mapper and drops one.
const llmV1 = `{
"name": "llm",
"condition": {"attributes": [{"value": "model"}], "resource": []},
"enabled": true,
"mappers": [
{"name": "gen_ai.request.model", "fieldContext": "attribute", "config": {"sources": [
{"key": "llm.model_name", "context": "attribute", "operation": "copy", "priority": 20},
{"key": "ai.model.id", "context": "attribute", "operation": "copy", "priority": 10}
]}},
{"name": "gen_ai.input.messages", "fieldContext": "attribute", "config": {"sources": [
{"key": "gen_ai.prompt", "context": "attribute", "operation": "copy", "priority": 10}
]}}
]
}`
const llmV2 = `{
"name": "llm",
"condition": {"attributes": [{"value": "model"}, {"value": "llm."}], "resource": []},
"enabled": true,
"mappers": [
{"name": "gen_ai.request.model", "fieldContext": "attribute", "config": {"sources": [
{"key": "llm.model_name", "context": "attribute", "operation": "copy", "priority": 20},
{"key": "langfuse.observation.model.name", "context": "attribute", "operation": "copy", "priority": 10}
]}},
{"name": "gen_ai.provider.name", "fieldContext": "attribute", "config": {"sources": [
{"key": "llm.vendor", "context": "attribute", "operation": "copy", "priority": 10}
]}}
]
}`
func findMapper(t *testing.T, mappers []*spantypes.SpanMapper, name string) *spantypes.SpanMapper {
t.Helper()
for _, m := range mappers {
if m.Name == name {
return m
}
}
require.Failf(t, "mapper not found", "no mapper named %q", name)
return nil
}
func findSource(t *testing.T, sources []spantypes.SpanMapperSource, key string, origin spantypes.SpanMapperOrigin) spantypes.SpanMapperSource {
t.Helper()
for _, s := range sources {
if s.Key == key && s.Origin == origin {
return s
}
}
require.Failf(t, "source not found", "no %s source with key %q", origin.StringValue(), key)
return spantypes.SpanMapperSource{}
}
func TestReconcileUpgradeKeepsTogglesAndUserItems(t *testing.T) {
ctx := context.Background()
orgID := valuer.GenerateUUID()
sqlStore := newTestSQLStore(t)
v1 := newTestModule(t, sqlStore, newTestDefinition(t, 1, llmV1))
require.NoError(t, v1.ReconcileSystemGroups(ctx, orgID))
group, err := v1.store.GetGroupByName(ctx, orgID, "llm")
require.NoError(t, err)
mappers, err := v1.ListMappers(ctx, orgID, group.ID)
require.NoError(t, err)
// Switch the shipped substring off and add a user one.
off := false
require.NoError(t, v1.UpdateGroup(ctx, orgID, group.ID, nil, &spantypes.SpanMapperGroupCondition{
Attributes: []spantypes.SpanMapperGroupConditionKey{
{Value: "model", Enabled: false, Origin: spantypes.SpanMapperOriginSystem},
{Value: "gen_ai.request.model", Enabled: true, Origin: spantypes.SpanMapperOriginUser},
},
Resource: []spantypes.SpanMapperGroupConditionKey{},
}, &off, testUser))
// Switch a shipped source off, add a user override, and switch the mapper off.
model := findMapper(t, mappers, "gen_ai.request.model")
require.NoError(t, v1.UpdateMapper(ctx, orgID, group.ID, model.ID, spantypes.FieldContext{}, &spantypes.SpanMapperConfig{Sources: []spantypes.SpanMapperSource{
{Key: "llm.model_name", Context: spantypes.FieldContextSpanAttribute, Operation: spantypes.SpanMapperOperationCopy, Priority: 20, Enabled: false, Origin: spantypes.SpanMapperOriginSystem},
{Key: "llm.model_name", Context: spantypes.FieldContextSpanAttribute, Operation: spantypes.SpanMapperOperationMove, Priority: 1, Enabled: true, Origin: spantypes.SpanMapperOriginUser},
}}, &off, testUser))
// Add a user source to the mapper v2 stops shipping, so it must survive.
messages := findMapper(t, mappers, "gen_ai.input.messages")
require.NoError(t, v1.UpdateMapper(ctx, orgID, group.ID, messages.ID, spantypes.FieldContext{}, &spantypes.SpanMapperConfig{Sources: []spantypes.SpanMapperSource{
{Key: "input.value", Context: spantypes.FieldContextSpanAttribute, Operation: spantypes.SpanMapperOperationCopy, Priority: 1, Enabled: true},
}}, nil, testUser))
// A user mapper in the shipped group.
require.NoError(t, v1.CreateMapper(ctx, orgID, group.ID, spantypes.NewSpanMapper(group.ID, testUser, &spantypes.PostableSpanMapper{
Name: "gen_ai.custom", FieldContext: spantypes.FieldContextSpanAttribute, Enabled: true,
Config: spantypes.SpanMapperConfig{Sources: []spantypes.SpanMapperSource{{Key: "custom", Context: spantypes.FieldContextSpanAttribute, Operation: spantypes.SpanMapperOperationCopy, Priority: 1, Enabled: true}}},
})))
v2 := newTestModule(t, sqlStore, newTestDefinition(t, 2, llmV2))
require.NoError(t, v2.ReconcileSystemGroups(ctx, orgID))
upgraded, err := v2.GetGroup(ctx, orgID, group.ID)
require.NoError(t, err)
assert.Equal(t, 2, upgraded.Version)
assert.False(t, upgraded.Enabled)
assert.Equal(t, spantypes.ProvisionerIdentity, upgraded.UpdatedBy)
assert.Equal(t, []spantypes.SpanMapperGroupConditionKey{
{Value: "model", Enabled: false, Origin: spantypes.SpanMapperOriginSystem},
{Value: "llm.", Enabled: true, Origin: spantypes.SpanMapperOriginSystem},
{Value: "gen_ai.request.model", Enabled: true, Origin: spantypes.SpanMapperOriginUser},
}, upgraded.Condition.Attributes)
mappers, err = v2.ListMappers(ctx, orgID, group.ID)
require.NoError(t, err)
require.Len(t, mappers, 4)
model = findMapper(t, mappers, "gen_ai.request.model")
assert.False(t, model.Enabled)
assert.Equal(t, spantypes.SpanMapperOriginSystem, model.Origin)
assert.False(t, findSource(t, model.Config.Sources, "llm.model_name", spantypes.SpanMapperOriginSystem).Enabled)
assert.True(t, findSource(t, model.Config.Sources, "langfuse.observation.model.name", spantypes.SpanMapperOriginSystem).Enabled)
assert.Equal(t, spantypes.SpanMapperOperationMove, findSource(t, model.Config.Sources, "llm.model_name", spantypes.SpanMapperOriginUser).Operation)
assert.Len(t, model.Config.Sources, 3)
messages = findMapper(t, mappers, "gen_ai.input.messages")
assert.Equal(t, spantypes.SpanMapperOriginUser, messages.Origin)
require.Len(t, messages.Config.Sources, 1)
assert.Equal(t, "input.value", messages.Config.Sources[0].Key)
assert.Equal(t, spantypes.SpanMapperOriginSystem, findMapper(t, mappers, "gen_ai.provider.name").Origin)
assert.Equal(t, spantypes.SpanMapperOriginUser, findMapper(t, mappers, "gen_ai.custom").Origin)
// Shipping v1 again drops provider.name outright (no user sources) and
// re-adopts the surviving user mapper input.messages as a shipped one.
v3 := newTestModule(t, sqlStore, newTestDefinition(t, 3, llmV1))
require.NoError(t, v3.ReconcileSystemGroups(ctx, orgID))
mappers, err = v3.ListMappers(ctx, orgID, group.ID)
require.NoError(t, err)
require.Len(t, mappers, 3)
for _, m := range mappers {
assert.NotEqual(t, "gen_ai.provider.name", m.Name)
}
messages = findMapper(t, mappers, "gen_ai.input.messages")
assert.Equal(t, spantypes.SpanMapperOriginSystem, messages.Origin)
assert.Len(t, messages.Config.Sources, 2)
}
func TestReconcileDoesNotDowngrade(t *testing.T) {
ctx := context.Background()
orgID := valuer.GenerateUUID()
sqlStore := newTestSQLStore(t)
require.NoError(t, newTestModule(t, sqlStore, newTestDefinition(t, 2, llmV2)).ReconcileSystemGroups(ctx, orgID))
older := newTestModule(t, sqlStore, newTestDefinition(t, 1, llmV1))
require.NoError(t, older.ReconcileSystemGroups(ctx, orgID))
group, err := older.store.GetGroupByName(ctx, orgID, "llm")
require.NoError(t, err)
assert.Equal(t, 2, group.Version)
mappers, err := older.ListMappers(ctx, orgID, group.ID)
require.NoError(t, err)
findMapper(t, mappers, "gen_ai.provider.name")
}

View File

@@ -0,0 +1,187 @@
package implspanmapper
import (
"context"
"log/slog"
"slices"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/query-service/agentConf"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
func (module *module) ReconcileSystemGroups(ctx context.Context, orgID valuer.UUID) error {
for _, definition := range module.registry.List() {
if err := module.reconcileSystemGroup(ctx, orgID, definition); err != nil {
return err
}
}
agentConf.NotifyConfigUpdate(ctx)
return nil
}
// reconcileSystemGroup brings one org's copy of a definition to the shipped
// version in a single transaction. A concurrent provisioner (another replica,
// or the org-creation hook racing the startup sweep) loses on the group's
// unique (org_id, name) index and is treated as a no-op.
func (module *module) reconcileSystemGroup(ctx context.Context, orgID valuer.UUID, definition spantypes.SpanMapperGroupDefinition) error {
err := module.store.RunInTx(ctx, func(ctx context.Context) error {
group, err := module.store.GetGroupByName(ctx, orgID, definition.Name())
if err != nil && errors.Ast(err, errors.TypeNotFound) {
group = newSystemGroup(orgID, definition)
err = module.store.CreateGroup(ctx, group)
}
if err != nil {
return err
}
if group.Origin != spantypes.SpanMapperOriginSystem {
module.settings.Logger().WarnContext(ctx, "skipping default span mapper group: a user group holds its name", slog.String("name", definition.Name()), slog.String("org_id", orgID.StringValue()))
return nil
}
if group.Version >= definition.Version {
return nil
}
return module.applyDefinition(ctx, orgID, group, definition)
})
if err != nil && errors.Ast(err, errors.TypeAlreadyExists) {
module.settings.Logger().DebugContext(ctx, "default span mapper group provisioned concurrently", slog.String("name", definition.Name()), slog.String("org_id", orgID.StringValue()))
return nil
}
return err
}
// applyDefinition replaces every shipped item with the definition, carrying each
// enabled flag over by identity, and leaves user items untouched. A mapper that
// is no longer shipped is deleted unless the user added sources to it, in which
// case it survives as a user mapper.
func (module *module) applyDefinition(ctx context.Context, orgID valuer.UUID, group *spantypes.SpanMapperGroup, definition spantypes.SpanMapperGroupDefinition) error {
mappers, err := module.store.ListMappers(ctx, orgID, group.ID)
if err != nil {
return err
}
byName := make(map[string]*spantypes.SpanMapper, len(mappers))
for _, m := range mappers {
byName[m.Name] = m
}
now := time.Now()
for i := range definition.Definition.Mappers {
pm := &definition.Definition.Mappers[i]
mapper, exists := byName[pm.Name]
delete(byName, pm.Name)
if !exists {
if err := module.store.CreateMapper(ctx, newSystemMapper(group.ID, pm)); err != nil {
return err
}
continue
}
mapper.Config.Sources = mergeShippedSources(mapper.Config.Sources, pm.Config.Sources)
mapper.FieldContext = pm.FieldContext
mapper.Origin = spantypes.SpanMapperOriginSystem
mapper.UpdatedAt = now
mapper.UpdatedBy = spantypes.ProvisionerIdentity
if err := module.store.UpdateMapper(ctx, mapper); err != nil {
return err
}
}
// Whatever is left in byName is not shipped any more.
for _, mapper := range byName {
if mapper.Origin != spantypes.SpanMapperOriginSystem {
continue
}
mapper.Config.Sources = mergeShippedSources(mapper.Config.Sources, nil)
if len(mapper.Config.Sources) == 0 {
if err := module.store.DeleteMapper(ctx, orgID, group.ID, mapper.ID, spantypes.SpanMapperOriginSystem); err != nil {
return err
}
continue
}
mapper.Origin = spantypes.SpanMapperOriginUser
mapper.UpdatedAt = now
mapper.UpdatedBy = spantypes.ProvisionerIdentity
if err := module.store.UpdateMapper(ctx, mapper); err != nil {
return err
}
}
shipped := definition.Definition.Condition
group.Condition = spantypes.SpanMapperGroupCondition{
Attributes: mergeShippedConditionKeys(group.Condition.Attributes, shipped.Attributes),
Resource: mergeShippedConditionKeys(group.Condition.Resource, shipped.Resource),
}
group.Version = definition.Version
group.UpdatedAt = now
group.UpdatedBy = spantypes.ProvisionerIdentity
if err := module.store.UpdateGroup(ctx, group); err != nil {
return err
}
module.settings.Logger().InfoContext(ctx, "applied default span mapper group", slog.String("name", definition.Name()), slog.Int("version", definition.Version), slog.String("org_id", orgID.StringValue()))
return nil
}
// newSystemGroup is the empty shell applyDefinition fills: version 0 so the
// definition is applied right after the row exists.
func newSystemGroup(orgID valuer.UUID, definition spantypes.SpanMapperGroupDefinition) *spantypes.SpanMapperGroup {
group := spantypes.NewSpanMapperGroup(orgID, spantypes.ProvisionerIdentity, &definition.Definition.PostableSpanMapperGroup)
group.Condition = definition.Definition.Condition
group.Enabled = true
group.Origin = spantypes.SpanMapperOriginSystem
return group
}
func newSystemMapper(groupID valuer.UUID, pm *spantypes.PostableSpanMapper) *spantypes.SpanMapper {
mapper := spantypes.NewSpanMapper(groupID, spantypes.ProvisionerIdentity, pm)
mapper.Config = pm.Config
mapper.Enabled = true
mapper.Origin = spantypes.SpanMapperOriginSystem
return mapper
}
// mergeShippedConditionKeys returns the shipped keys, each keeping the enabled
// flag of the stored system key with the same value, followed by the stored
// user keys.
func mergeShippedConditionKeys(stored, shipped []spantypes.SpanMapperGroupConditionKey) []spantypes.SpanMapperGroupConditionKey {
out := make([]spantypes.SpanMapperGroupConditionKey, 0, len(stored)+len(shipped))
for _, k := range shipped {
idx := slices.IndexFunc(stored, func(s spantypes.SpanMapperGroupConditionKey) bool {
return s.Origin == spantypes.SpanMapperOriginSystem && s.Value == k.Value
})
if idx != -1 {
k.Enabled = stored[idx].Enabled
}
out = append(out, k)
}
for _, k := range stored {
if k.Origin != spantypes.SpanMapperOriginSystem {
out = append(out, k)
}
}
return out
}
// mergeShippedSources returns the shipped sources, each keeping the enabled
// flag of the stored system source with the same key and context, followed by
// the stored user sources.
func mergeShippedSources(stored, shipped []spantypes.SpanMapperSource) []spantypes.SpanMapperSource {
out := make([]spantypes.SpanMapperSource, 0, len(stored)+len(shipped))
for _, s := range shipped {
idx := slices.IndexFunc(stored, func(o spantypes.SpanMapperSource) bool {
return o.Origin == spantypes.SpanMapperOriginSystem && o.Key == s.Key && o.Context == s.Context
})
if idx != -1 {
s.Enabled = stored[idx].Enabled
}
out = append(out, s)
}
for _, s := range stored {
if s.Origin != spantypes.SpanMapperOriginSystem {
out = append(out, s)
}
}
return out
}

View File

@@ -0,0 +1,81 @@
package implspanmapper
import (
"context"
"log/slog"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/modules/organization"
"github.com/SigNoz/signoz/pkg/modules/spanmapper"
)
const reconcileRetryInterval = 30 * time.Second
type service struct {
settings factory.ScopedProviderSettings
module spanmapper.Module
orgGetter organization.Getter
stopC chan struct{}
healthyC chan struct{}
}
// NewService reconciles every org's default mapping groups once at startup.
// Orgs created later are reconciled by the organization setter instead.
func NewService(providerSettings factory.ProviderSettings, module spanmapper.Module, orgGetter organization.Getter) factory.Service {
return &service{
settings: factory.NewScopedProviderSettings(providerSettings, "github.com/SigNoz/signoz/pkg/modules/spanmapper/implspanmapper"),
module: module,
orgGetter: orgGetter,
stopC: make(chan struct{}),
healthyC: make(chan struct{}),
}
}
func (service *service) Start(ctx context.Context) error {
ticker := time.NewTicker(reconcileRetryInterval)
defer ticker.Stop()
for {
err := service.reconcile(ctx)
if err == nil {
close(service.healthyC)
<-service.stopC
return nil
}
service.settings.Logger().WarnContext(ctx, "default span mapper group reconciliation failed, retrying", errors.Attr(err))
select {
case <-service.stopC:
return nil
case <-ticker.C:
}
}
}
func (service *service) Healthy() <-chan struct{} {
return service.healthyC
}
func (service *service) Stop(_ context.Context) error {
close(service.stopC)
return nil
}
func (service *service) reconcile(ctx context.Context) error {
orgs, err := service.orgGetter.ListByOwnedKeyRange(ctx)
if err != nil {
return err
}
for _, org := range orgs {
if err := service.module.ReconcileSystemGroups(ctx, org.ID); err != nil {
return errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "couldn't reconcile default span mapper groups for org %s", org.ID.StringValue())
}
}
service.settings.Logger().InfoContext(ctx, "default span mapper group reconciliation completed", slog.Int("orgs", len(orgs)))
return nil
}

View File

@@ -17,6 +17,10 @@ func NewStore(sqlstore sqlstore.SQLStore) spantypes.SpanMapperStore {
return &store{sqlstore: sqlstore}
}
func (s *store) RunInTx(ctx context.Context, cb func(ctx context.Context) error) error {
return s.sqlstore.RunInTxCtx(ctx, nil, cb)
}
func (s *store) CreateGroup(ctx context.Context, group *spantypes.SpanMapperGroup) error {
storable := group.ToStorable()
_, err := s.sqlstore.
@@ -34,7 +38,7 @@ func (s *store) GetGroup(ctx context.Context, orgID, id valuer.UUID) (*spantypes
storable := new(spantypes.StorableSpanMapperGroup)
err := s.sqlstore.
BunDB().
BunDBCtx(ctx).
NewSelect().
Model(storable).
Where("org_id = ?", orgID).
@@ -46,11 +50,27 @@ func (s *store) GetGroup(ctx context.Context, orgID, id valuer.UUID) (*spantypes
return storable.ToSpanMapperGroup(), nil
}
func (s *store) GetGroupByName(ctx context.Context, orgID valuer.UUID, name string) (*spantypes.SpanMapperGroup, error) {
storable := new(spantypes.StorableSpanMapperGroup)
err := s.sqlstore.
BunDBCtx(ctx).
NewSelect().
Model(storable).
Where("org_id = ?", orgID).
Where("name = ?", name).
Scan(ctx)
if err != nil {
return nil, s.sqlstore.WrapNotFoundErrf(err, spantypes.ErrCodeMappingGroupNotFound, "span mapper group %q not found", name)
}
return storable.ToSpanMapperGroup(), nil
}
func (s *store) ListGroups(ctx context.Context, orgID valuer.UUID, q *spantypes.ListSpanMapperGroupsQuery) ([]*spantypes.SpanMapperGroup, error) {
storables := make([]*spantypes.StorableSpanMapperGroup, 0)
sel := s.sqlstore.
BunDB().
BunDBCtx(ctx).
NewSelect().
Model(&storables).
Where("org_id = ?", orgID)
@@ -91,38 +111,43 @@ func (s *store) UpdateGroup(ctx context.Context, group *spantypes.SpanMapperGrou
}
func (s *store) DeleteGroup(ctx context.Context, orgID, id valuer.UUID) error {
tx, err := s.sqlstore.BunDBCtx(ctx).BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
return s.RunInTx(ctx, func(ctx context.Context) error {
db := s.sqlstore.BunDBCtx(ctx)
// Cascade: remove mappers belonging to this group first.
if _, err := tx.NewDelete().
Model((*spantypes.StorableSpanMapper)(nil)).
Where("group_id = ?", id).
Exec(ctx); err != nil {
return err
}
// Cascade: remove mappers belonging to this group first.
if _, err := db.NewDelete().
Model((*spantypes.StorableSpanMapper)(nil)).
Where("group_id = ?", id).
Exec(ctx); err != nil {
return err
}
res, err := tx.NewDelete().
Model((*spantypes.StorableSpanMapperGroup)(nil)).
Where("org_id = ?", orgID).
Where("id = ?", id).
Exec(ctx)
if err != nil {
return err
}
res, err := db.NewDelete().
Model((*spantypes.StorableSpanMapperGroup)(nil)).
Where("org_id = ?", orgID).
Where("id = ?", id).
Where("origin = ?", spantypes.SpanMapperOriginUser).
Exec(ctx)
if err != nil {
return err
}
rowsAffected, err := res.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
return errors.Newf(errors.TypeNotFound, spantypes.ErrCodeMappingGroupNotFound, "span mapper group %s not found", id)
}
return tx.Commit()
rowsAffected, err := res.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
group, err := s.GetGroup(ctx, orgID, id)
if err != nil {
return err
}
if err := group.ErrIfNotDeletable(); err != nil {
return err
}
return errors.Newf(errors.TypeNotFound, spantypes.ErrCodeMappingGroupNotFound, "span mapper group %s not found", id)
}
return nil
})
}
func (s *store) CreateMapper(ctx context.Context, mapper *spantypes.SpanMapper) error {
@@ -146,7 +171,7 @@ func (s *store) GetMapper(ctx context.Context, orgID, groupID, id valuer.UUID) (
storable := new(spantypes.StorableSpanMapper)
err := s.sqlstore.
BunDB().
BunDBCtx(ctx).
NewSelect().
Model(storable).
Where("group_id = ?", groupID).
@@ -166,7 +191,7 @@ func (s *store) ListMappers(ctx context.Context, orgID, groupID valuer.UUID) ([]
storables := make([]*spantypes.StorableSpanMapper, 0)
if err := s.sqlstore.
BunDB().
BunDBCtx(ctx).
NewSelect().
Model(&storables).
Where("group_id = ?", groupID).
@@ -200,7 +225,7 @@ func (s *store) UpdateMapper(ctx context.Context, mapper *spantypes.SpanMapper)
return nil
}
func (s *store) DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID) error {
func (s *store) DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID, origin spantypes.SpanMapperOrigin) error {
if _, err := s.GetGroup(ctx, orgID, groupID); err != nil {
return err
}
@@ -211,6 +236,7 @@ func (s *store) DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID
Model((*spantypes.StorableSpanMapper)(nil)).
Where("group_id = ?", groupID).
Where("id = ?", id).
Where("origin = ?", origin).
Exec(ctx)
if err != nil {
return err
@@ -221,6 +247,13 @@ func (s *store) DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID
return err
}
if rowsAffected == 0 {
mapper, err := s.GetMapper(ctx, orgID, groupID, id)
if err != nil {
return err
}
if err := mapper.ErrIfNotDeletable(); err != nil {
return err
}
return errors.Newf(errors.TypeNotFound, spantypes.ErrCodeMapperNotFound, "span mapper %s not found", id)
}
return nil

View File

@@ -28,6 +28,10 @@ type Module interface {
UpdateMapper(ctx context.Context, orgID, groupID, id valuer.UUID, fieldContext spantypes.FieldContext, config *spantypes.SpanMapperConfig, enabled *bool, updatedBy string) error
DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID) error
TestMappers(ctx context.Context, orgID valuer.UUID, spans []spantypes.SpanMapperTestSpan, groups []*spantypes.SpanMapperGroupWithMappers) ([]spantypes.SpanMapperTestSpan, []string, error)
// ReconcileSystemGroups provisions or upgrades the shipped mapping groups
// for one org. It runs at startup for every org and again on org creation.
ReconcileSystemGroups(ctx context.Context, orgID valuer.UUID) error
}
// Handler defines the HTTP handler interface for mapping group and mapper endpoints.

View File

@@ -18,6 +18,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/dashboard/impldashboard"
"github.com/SigNoz/signoz/pkg/modules/organization/implorganization"
"github.com/SigNoz/signoz/pkg/modules/retention/implretention"
"github.com/SigNoz/signoz/pkg/modules/spanmapper/implspanmapper"
"github.com/SigNoz/signoz/pkg/modules/tag/impltag"
"github.com/SigNoz/signoz/pkg/modules/user/impluser"
"github.com/SigNoz/signoz/pkg/querier"
@@ -61,7 +62,10 @@ func TestNewHandlers(t *testing.T) {
userGetter := impluser.NewGetter(impluser.NewStore(sqlstore, providerSettings), userRoleStore, flagger)
retentionGetter := implretention.NewGetter(implretention.NewStore(sqlstore))
modules := NewModules(sqlstore, tokenizer, emailing, providerSettings, orgGetter, alertmanager, nil, nil, nil, nil, nil, nil, nil, queryParser, Config{}, dashboardModule, userGetter, userRoleStore, nil, nil, nil, retentionGetter, flagger, tagModule, nil)
spanMapperRegistry, err := implspanmapper.NewSystemGroupRegistry()
require.NoError(t, err)
spanMapperModule := implspanmapper.NewModule(implspanmapper.NewStore(sqlstore), spanMapperRegistry, providerSettings)
modules := NewModules(sqlstore, tokenizer, emailing, providerSettings, orgGetter, alertmanager, nil, nil, nil, nil, nil, nil, nil, queryParser, Config{}, dashboardModule, userGetter, userRoleStore, nil, nil, nil, retentionGetter, flagger, tagModule, nil, spanMapperModule)
querierHandler := querier.NewHandler(providerSettings, nil, nil)
registryHandler := factory.NewHandler(nil)

View File

@@ -45,7 +45,6 @@ import (
"github.com/SigNoz/signoz/pkg/modules/session"
"github.com/SigNoz/signoz/pkg/modules/session/implsession"
"github.com/SigNoz/signoz/pkg/modules/spanmapper"
"github.com/SigNoz/signoz/pkg/modules/spanmapper/implspanmapper"
"github.com/SigNoz/signoz/pkg/modules/spanpercentile"
"github.com/SigNoz/signoz/pkg/modules/spanpercentile/implspanpercentile"
"github.com/SigNoz/signoz/pkg/modules/tag"
@@ -124,9 +123,10 @@ func NewModules(
fl flagger.Flagger,
tagModule tag.Module,
metricReductionRule metricreductionrule.Module,
spanMapper spanmapper.Module,
) Modules {
quickfilter := implquickfilter.NewModule(implquickfilter.NewStore(sqlstore))
orgSetter := implorganization.NewSetter(implorganization.NewStore(sqlstore), alertmanager, quickfilter, dashboard)
orgSetter := implorganization.NewSetter(implorganization.NewStore(sqlstore), alertmanager, quickfilter, dashboard, spanMapper)
// Cleanup callbacks from other modules, invoked when a user is deleted.
onDeleteUser := []user.OnDeleteUser{
dashboard.DeletePreferencesForUser,
@@ -162,8 +162,8 @@ func NewModules(
RuleStateHistory: implrulestatehistory.NewModule(implrulestatehistory.NewStore(telemetryStore, telemetryMetadataStore, providerSettings.Logger), ruleStore),
CloudIntegration: cloudIntegrationModule,
TraceDetail: impltracedetail.NewModule(impltracedetail.NewTraceStore(telemetryStore), providerSettings, config.TraceDetail),
SpanMapper: implspanmapper.NewModule(implspanmapper.NewStore(sqlstore), fl),
LLMPricingRule: impllmpricingrule.NewModule(impllmpricingrule.NewStore(sqlstore), fl, querier),
SpanMapper: spanMapper,
LLMPricingRule: impllmpricingrule.NewModule(impllmpricingrule.NewStore(sqlstore), querier),
Tag: tagModule,
}
}

View File

@@ -21,6 +21,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/retention/implretention"
"github.com/SigNoz/signoz/pkg/modules/serviceaccount"
"github.com/SigNoz/signoz/pkg/modules/serviceaccount/implserviceaccount"
"github.com/SigNoz/signoz/pkg/modules/spanmapper/implspanmapper"
"github.com/SigNoz/signoz/pkg/modules/tag/impltag"
"github.com/SigNoz/signoz/pkg/modules/user/impluser"
"github.com/SigNoz/signoz/pkg/queryparser"
@@ -68,7 +69,11 @@ func TestNewModules(t *testing.T) {
retentionGetter := implretention.NewGetter(implretention.NewStore(sqlstore))
modules := NewModules(sqlstore, tokenizer, emailing, providerSettings, orgGetter, alertmanager, nil, nil, nil, nil, nil, nil, nil, queryParser, Config{}, dashboardModule, userGetter, userRoleStore, serviceAccount, serviceAccountGetter, implcloudintegration.NewModule(), retentionGetter, flagger, tagModule, implmetricreductionrule.NewModule())
spanMapperRegistry, err := implspanmapper.NewSystemGroupRegistry()
require.NoError(t, err)
spanMapperModule := implspanmapper.NewModule(implspanmapper.NewStore(sqlstore), spanMapperRegistry, providerSettings)
modules := NewModules(sqlstore, tokenizer, emailing, providerSettings, orgGetter, alertmanager, nil, nil, nil, nil, nil, nil, nil, queryParser, Config{}, dashboardModule, userGetter, userRoleStore, serviceAccount, serviceAccountGetter, implcloudintegration.NewModule(), retentionGetter, flagger, tagModule, implmetricreductionrule.NewModule(), spanMapperModule)
reflectVal := reflect.ValueOf(modules)
for i := 0; i < reflectVal.NumField(); i++ {

View File

@@ -253,6 +253,7 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddIngestionTuplesFactory(sqlstore),
sqlmigration.NewAddSubscriptionTuplesFactory(sqlstore),
sqlmigration.NewNormalizeQuickFilterFieldsFactory(sqlstore),
sqlmigration.NewAddSpanMapperOriginFactory(sqlstore, sqlschema),
)
}

View File

@@ -36,6 +36,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/rulestatehistory"
"github.com/SigNoz/signoz/pkg/modules/serviceaccount"
"github.com/SigNoz/signoz/pkg/modules/serviceaccount/implserviceaccount"
"github.com/SigNoz/signoz/pkg/modules/spanmapper/implspanmapper"
"github.com/SigNoz/signoz/pkg/modules/tag"
"github.com/SigNoz/signoz/pkg/modules/tag/impltag"
"github.com/SigNoz/signoz/pkg/modules/user/impluser"
@@ -526,7 +527,15 @@ func New(
metricReductionRuleModule := metricReductionRuleModuleCallback(sqlstore, telemetrystore, dashboard, queryParser, licensing, flagger, telemetryMetadataStore, providerSettings, config.MetricsExplorer.TelemetryStore.Threads)
// Initialize all modules
modules := NewModules(sqlstore, tokenizer, emailing, providerSettings, orgGetter, alertmanager, analytics, querier, telemetrystore, telemetryMetadataStore, authNs, authz, cache, queryParser, config, dashboard, userGetter, userRoleStore, serviceAccount, serviceAccountGetter, cloudIntegrationModule, retentionGetter, flagger, tagModule, metricReductionRuleModule)
// The default mapping group registry is parsed here so a malformed embedded
// definition fails startup instead of a request.
spanMapperRegistry, err := implspanmapper.NewSystemGroupRegistry()
if err != nil {
return nil, err
}
spanMapperModule := implspanmapper.NewModule(implspanmapper.NewStore(sqlstore), spanMapperRegistry, providerSettings)
modules := NewModules(sqlstore, tokenizer, emailing, providerSettings, orgGetter, alertmanager, analytics, querier, telemetrystore, telemetryMetadataStore, authNs, authz, cache, queryParser, config, dashboard, userGetter, userRoleStore, serviceAccount, serviceAccountGetter, cloudIntegrationModule, retentionGetter, flagger, tagModule, metricReductionRuleModule, spanMapperModule)
// Initialize ruler from the variant-specific provider factories
rulerInstance, err := factory.NewProviderFromNamedMap(ctx, providerSettings, config.Ruler, rulerProviderFactories(cache, alertmanager, sqlstore, telemetrystore, telemetryMetadataStore, prometheus, orgGetter, modules.RuleStateHistory, querier, queryParser), "signoz")
@@ -596,6 +605,7 @@ func New(
factory.NewNamedService(factory.MustNewName("meterreporter"), meterReporter, factory.MustNewName("licensing")),
factory.NewNamedService(factory.MustNewName("ruler"), rulerInstance),
factory.NewNamedService(factory.MustNewName("systemdashboard"), impldashboard.NewService(providerSettings, dashboard, orgGetter)),
factory.NewNamedService(factory.MustNewName("spanmappergroup"), implspanmapper.NewService(providerSettings, spanMapperModule, orgGetter)),
)
if err != nil {
return nil, err

View File

@@ -0,0 +1,93 @@
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 addSpanMapperOrigin struct {
sqlstore sqlstore.SQLStore
sqlschema sqlschema.SQLSchema
}
func NewAddSpanMapperOriginFactory(sqlstore sqlstore.SQLStore, sqlschema sqlschema.SQLSchema) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(
factory.MustNewName("add_span_mapper_origin"),
func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &addSpanMapperOrigin{sqlstore: sqlstore, sqlschema: sqlschema}, nil
},
)
}
func (migration *addSpanMapperOrigin) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
// Up adds the ownership columns that let SigNoz ship default mapping groups
// alongside user ones.
func (migration *addSpanMapperOrigin) Up(ctx context.Context, db *bun.DB) error {
// span_mapper references span_mapper_group and both have foreign keys, so
// enforcement must be off for the SQLite recreate-table fallback.
if err := migration.sqlschema.ToggleFKEnforcement(ctx, db, false); err != nil {
return err
}
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() {
_ = tx.Rollback()
}()
groupTable, groupUniqueConstraints, err := migration.sqlschema.GetTable(ctx, sqlschema.TableName("span_mapper_group"))
if err != nil {
return err
}
sqls := migration.sqlschema.Operator().AddColumn(groupTable, groupUniqueConstraints, &sqlschema.Column{
Name: sqlschema.ColumnName("origin"),
DataType: sqlschema.DataTypeText,
Nullable: false,
Default: "'user'",
}, "user")
sqls = append(sqls, migration.sqlschema.Operator().AddColumn(groupTable, groupUniqueConstraints, &sqlschema.Column{
Name: sqlschema.ColumnName("version"),
DataType: sqlschema.DataTypeBigInt,
Nullable: false,
Default: "0",
}, 0)...)
mapperTable, mapperUniqueConstraints, err := migration.sqlschema.GetTable(ctx, sqlschema.TableName("span_mapper"))
if err != nil {
return err
}
sqls = append(sqls, migration.sqlschema.Operator().AddColumn(mapperTable, mapperUniqueConstraints, &sqlschema.Column{
Name: sqlschema.ColumnName("origin"),
DataType: sqlschema.DataTypeText,
Nullable: false,
Default: "'user'",
}, "user")...)
for _, sql := range sqls {
if _, err := tx.ExecContext(ctx, string(sql)); err != nil {
return err
}
}
if err := tx.Commit(); err != nil {
return err
}
return migration.sqlschema.ToggleFKEnforcement(ctx, db, true)
}
func (migration *addSpanMapperOrigin) Down(context.Context, *bun.DB) error {
return nil
}

View File

@@ -61,8 +61,12 @@ type LLMPricingRuleProcessorOutputAttrs struct {
func buildProcessorConfig(rules []*LLMPricingRule) *LLMPricingRuleProcessorConfig {
pricingRules := make([]LLMPricingRuleProcessor, 0, len(rules))
for _, r := range rules {
// The collector rejects negative prices (providers use them for "unknown").
if r.Pricing.Input < 0 || r.Pricing.Output < 0 {
continue
}
var cache *LLMPricingRuleProcessorCache
if r.Pricing.Cache != nil {
if r.Pricing.Cache != nil && r.Pricing.Cache.Read >= 0 && r.Pricing.Cache.Write >= 0 {
mode := r.Pricing.Cache.Mode.StringValue()
if mode != LLMPricingRuleCacheModeSubtract.StringValue() && mode != LLMPricingRuleCacheModeAdditive.StringValue() {
mode = ""

View File

@@ -90,6 +90,15 @@ func TestGenerateCollectorConfigWithLLMPricingProcessor(t *testing.T) {
},
expectedFile: "collector_rule_cache_mode_unknown.yaml",
},
// Negative in/out drops the rule; a negative cache price drops only the cache block.
{
name: "negative_prices",
rules: []*LLMPricingRule{
makePricingRule("openrouter/auto", []string{"openrouter/auto", "auto"}, LLMPricingRuleCacheModeSubtract, -1, -1, -1, -1),
makePricingRule("gpt-4o", []string{"gpt-4o*"}, LLMPricingRuleCacheModeSubtract, 5.0, 15.0, -1, 0),
},
expectedFile: "collector_rule_without_cache.yaml",
},
}
input, err := os.ReadFile(filepath.Join("testdata", "collector_baseline.yaml"))

View File

@@ -1,6 +1,8 @@
package spantypes
import (
"slices"
"strings"
"time"
"github.com/SigNoz/signoz/pkg/errors"
@@ -11,6 +13,7 @@ import (
var (
ErrCodeMapperNotFound = errors.MustNewCode("span_attribute_mapper_not_found")
ErrCodeMapperAlreadyExists = errors.MustNewCode("span_attribute_mapper_already_exists")
ErrCodeMapperNotDeletable = errors.MustNewCode("span_attribute_mapper_not_deletable")
ErrCodeMappingInvalidInput = errors.MustNewCode("span_attribute_mapping_invalid_input")
)
@@ -34,12 +37,25 @@ var (
SpanMapperOperationCopy = SpanMapperOperation{valuer.NewString("copy")}
)
// SpanMapperOrigin tells shipped (system) items apart from user-created ones.
// System items are read-only apart from their enabled toggle.
type SpanMapperOrigin struct {
valuer.String
}
var (
SpanMapperOriginUser = SpanMapperOrigin{valuer.NewString("user")}
SpanMapperOriginSystem = SpanMapperOrigin{valuer.NewString("system")}
)
// MapperSource describes one candidate source for a target attribute.
type SpanMapperSource struct {
Key string `json:"key" required:"true"`
Context FieldContext `json:"context" required:"true"`
Operation SpanMapperOperation `json:"operation" required:"true"`
Priority int `json:"priority" required:"true"`
Enabled bool `json:"enabled" required:"true"`
Origin SpanMapperOrigin `json:"origin"`
}
// MapperConfig holds the mapping logic for a single target attribute.
@@ -59,6 +75,7 @@ type SpanMapper struct {
FieldContext FieldContext `json:"fieldContext" required:"true"`
Config SpanMapperConfig `json:"config" required:"true"`
Enabled bool `json:"enabled" required:"true"`
Origin SpanMapperOrigin `json:"origin" required:"true"`
}
type PostableSpanMapper struct {
@@ -90,6 +107,63 @@ func (SpanMapperOperation) Enum() []any {
return []any{SpanMapperOperationMove, SpanMapperOperationCopy}
}
func (SpanMapperOrigin) Enum() []any {
return []any{SpanMapperOriginUser, SpanMapperOriginSystem}
}
func (p *PostableSpanMapper) Validate() error {
if strings.TrimSpace(p.Name) == "" {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "mapper name must not be blank")
}
if err := p.FieldContext.Validate(); err != nil {
return err
}
return p.Config.Validate()
}
func (f FieldContext) Validate() error {
if f != FieldContextSpanAttribute && f != FieldContextResource {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "field context must be one of %q or %q, got %q", FieldContextSpanAttribute, FieldContextResource, f.StringValue())
}
return nil
}
// Validate checks every source and rejects duplicate priorities within an
// origin. Shipped and user sources are never compared with each other: a user
// re-adding a shipped key with another operation is the supported override.
func (c *SpanMapperConfig) Validate() error {
if len(c.Sources) == 0 {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "config.sources must contain at least one source")
}
seen := map[SpanMapperOrigin]map[int]struct{}{}
for _, s := range c.Sources {
if strings.TrimSpace(s.Key) == "" {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "source key must not be blank")
}
if err := s.Context.Validate(); err != nil {
return err
}
if s.Operation != SpanMapperOperationCopy && s.Operation != SpanMapperOperationMove {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "source operation must be one of %q or %q, got %q", SpanMapperOperationCopy, SpanMapperOperationMove, s.Operation.StringValue())
}
if !s.Origin.IsZero() && s.Origin != SpanMapperOriginUser && s.Origin != SpanMapperOriginSystem {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "source origin must be one of %q or %q, got %q", SpanMapperOriginUser, SpanMapperOriginSystem, s.Origin.StringValue())
}
origin := s.Origin
if origin.IsZero() {
origin = SpanMapperOriginUser
}
if seen[origin] == nil {
seen[origin] = map[int]struct{}{}
}
if _, dup := seen[origin][s.Priority]; dup {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "source priority %d is used more than once", s.Priority)
}
seen[origin][s.Priority] = struct{}{}
}
return nil
}
func NewSpanMapper(groupID valuer.UUID, createdBy string, p *PostableSpanMapper) *SpanMapper {
now := time.Now()
return &SpanMapper{
@@ -97,8 +171,9 @@ func NewSpanMapper(groupID valuer.UUID, createdBy string, p *PostableSpanMapper)
GroupID: groupID,
Name: p.Name,
FieldContext: p.FieldContext,
Config: p.Config,
Config: SpanMapperConfig{Sources: withOrigin(p.Config.Sources, SpanMapperOriginUser)},
Enabled: p.Enabled,
Origin: SpanMapperOriginUser,
TimeAuditable: types.TimeAuditable{
CreatedAt: now,
UpdatedAt: now,
@@ -110,16 +185,42 @@ func NewSpanMapper(groupID valuer.UUID, createdBy string, p *PostableSpanMapper)
}
}
func (m *SpanMapper) Update(fieldContext FieldContext, config *SpanMapperConfig, enabled *bool, updatedBy string) {
m.FieldContext = fieldContext
// Update applies a user edit; a zero fieldContext means it was omitted. On a
// system mapper the field context is fixed and the stored system sources are
// kept; see nextSources.
func (m *SpanMapper) Update(fieldContext FieldContext, config *SpanMapperConfig, enabled *bool, updatedBy string) error {
if !fieldContext.IsZero() {
if m.Origin == SpanMapperOriginSystem && fieldContext != m.FieldContext {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "field context of system mapper %q cannot be changed", m.Name)
}
if err := fieldContext.Validate(); err != nil {
return err
}
m.FieldContext = fieldContext
}
if config != nil {
m.Config = *config
sources, err := m.nextSources(config.Sources)
if err != nil {
return err
}
m.Config = SpanMapperConfig{Sources: sources}
if err := m.Config.Validate(); err != nil {
return err
}
}
if enabled != nil {
m.Enabled = *enabled
}
m.UpdatedAt = time.Now()
m.UpdatedBy = updatedBy
return nil
}
func (m *SpanMapper) ErrIfNotDeletable() error {
if m.Origin == SpanMapperOriginSystem {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMapperNotDeletable, "system mapper %q cannot be deleted, disable it instead", m.Name)
}
return nil
}
func (m *SpanMapper) ToStorable() *StorableSpanMapper {
@@ -132,6 +233,7 @@ func (m *SpanMapper) ToStorable() *StorableSpanMapper {
FieldContext: m.FieldContext,
Config: m.Config,
Enabled: m.Enabled,
Origin: m.Origin,
}
}
@@ -145,6 +247,7 @@ func (s *StorableSpanMapper) ToSpanMapper() *SpanMapper {
FieldContext: s.FieldContext,
Config: s.Config,
Enabled: s.Enabled,
Origin: s.Origin,
}
}
@@ -159,3 +262,39 @@ func NewSpanMappersFromStorable(ss []*StorableSpanMapper) []*SpanMapper {
func NewGettableSpanMappers(m []*SpanMapper) *GettableSpanMappers {
return &GettableSpanMappers{Items: m}
}
// nextSources builds the source list from an edit: user sources are taken from
// the edit as sent, system sources stay as stored and the edit can only flip
// their enabled flag.
func (m *SpanMapper) nextSources(edit []SpanMapperSource) ([]SpanMapperSource, error) {
var systemSources, userSources []SpanMapperSource
for _, s := range m.Config.Sources {
if s.Origin == SpanMapperOriginSystem {
systemSources = append(systemSources, s)
}
}
for _, s := range edit {
if s.Origin != SpanMapperOriginSystem {
s.Origin = SpanMapperOriginUser
userSources = append(userSources, s)
continue
}
idx := slices.IndexFunc(systemSources, func(o SpanMapperSource) bool { return o.Key == s.Key && o.Context == s.Context })
if idx == -1 {
return nil, errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "system source %q does not exist on this mapper; only its enabled flag can change", s.Key)
}
systemSources[idx].Enabled = s.Enabled
}
return append(systemSources, userSources...), nil
}
func withOrigin(sources []SpanMapperSource, origin SpanMapperOrigin) []SpanMapperSource {
out := make([]SpanMapperSource, len(sources))
for i, s := range sources {
s.Origin = origin
out[i] = s
}
return out
}

View File

@@ -1,6 +1,8 @@
package spantypes
import (
"slices"
"strings"
"time"
"github.com/SigNoz/signoz/pkg/errors"
@@ -11,17 +13,28 @@ import (
var (
ErrCodeMappingGroupNotFound = errors.MustNewCode("span_attribute_mapping_group_not_found")
ErrCodeMappingGroupAlreadyExists = errors.MustNewCode("span_attribute_mapping_group_already_exists")
ErrCodeMappingGroupNameReserved = errors.MustNewCode("span_attribute_mapping_group_name_reserved")
ErrCodeMappingGroupNotDeletable = errors.MustNewCode("span_attribute_mapping_group_not_deletable")
)
// SpanMapperGroupConditionKey is one substring a span's attribute or resource
// keys are matched against.
type SpanMapperGroupConditionKey struct {
Value string `json:"value" required:"true"`
Enabled bool `json:"enabled" required:"true"`
Origin SpanMapperOrigin `json:"origin"`
}
// SpanMapperGroupCondition gates whether a group's rules run for a given span.
// A group runs when any attribute or resource key on the span CONTAINS one of
// the listed substrings (plain substring match — no glob syntax).
type SpanMapperGroupCondition struct {
Attributes []string `json:"attributes" required:"true" nullable:"true"`
Resource []string `json:"resource" required:"true" nullable:"true"`
Attributes []SpanMapperGroupConditionKey `json:"attributes" required:"true" nullable:"true"`
Resource []SpanMapperGroupConditionKey `json:"resource" required:"true" nullable:"true"`
}
// SpanMapperGroup is the domain model for a span attribute mapping group.
// Version is the shipped definition version for system groups and 0 otherwise.
type SpanMapperGroup struct {
types.TimeAuditable
types.UserAuditable
@@ -31,6 +44,8 @@ type SpanMapperGroup struct {
Name string `json:"name" required:"true"`
Condition SpanMapperGroupCondition `json:"condition" required:"true"`
Enabled bool `json:"enabled" required:"true"`
Origin SpanMapperOrigin `json:"origin" required:"true"`
Version int `json:"version" required:"true"`
}
// GettableSpanMapperGroup is the HTTP response representation of a mapping group.
@@ -58,14 +73,42 @@ type GettableSpanMapperGroups struct {
Items []*GettableSpanMapperGroup `json:"items" required:"true" nullable:"false"`
}
// Validate requires at least one substring overall and rejects blank ones.
// All-off is allowed: a group with every substring disabled simply never runs.
func (c *SpanMapperGroupCondition) Validate() error {
if len(c.Attributes)+len(c.Resource) == 0 {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "condition must list at least one attribute or resource substring")
}
for _, k := range slices.Concat(c.Attributes, c.Resource) {
if strings.TrimSpace(k.Value) == "" {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "condition substrings must not be blank")
}
if !k.Origin.IsZero() && k.Origin != SpanMapperOriginUser && k.Origin != SpanMapperOriginSystem {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "condition origin must be one of %q or %q, got %q", SpanMapperOriginUser, SpanMapperOriginSystem, k.Origin.StringValue())
}
}
return nil
}
func (p *PostableSpanMapperGroup) Validate() error {
if strings.TrimSpace(p.Name) == "" {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "group name must not be blank")
}
return p.Condition.Validate()
}
func NewSpanMapperGroup(orgID valuer.UUID, createdBy string, p *PostableSpanMapperGroup) *SpanMapperGroup {
now := time.Now()
return &SpanMapperGroup{
ID: valuer.GenerateUUID(),
OrgID: orgID,
Name: p.Name,
Condition: p.Condition,
Enabled: p.Enabled,
ID: valuer.GenerateUUID(),
OrgID: orgID,
Name: p.Name,
Condition: SpanMapperGroupCondition{
Attributes: conditionKeysWithOrigin(p.Condition.Attributes, SpanMapperOriginUser),
Resource: conditionKeysWithOrigin(p.Condition.Resource, SpanMapperOriginUser),
},
Enabled: p.Enabled,
Origin: SpanMapperOriginUser,
TimeAuditable: types.TimeAuditable{
CreatedAt: now,
UpdatedAt: now,
@@ -77,18 +120,45 @@ func NewSpanMapperGroup(orgID valuer.UUID, createdBy string, p *PostableSpanMapp
}
}
func (g *SpanMapperGroup) Update(name *string, condition *SpanMapperGroupCondition, enabled *bool, updatedBy string) {
// Update applies a user edit. A system group keeps its name and its system
// substrings; see nextConditionKeys.
func (g *SpanMapperGroup) Update(name *string, condition *SpanMapperGroupCondition, enabled *bool, updatedBy string) error {
if name != nil {
if g.Origin == SpanMapperOriginSystem && *name != g.Name {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "system group %q cannot be renamed", g.Name)
}
if strings.TrimSpace(*name) == "" {
return errors.New(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "group name must not be blank")
}
g.Name = *name
}
if condition != nil {
g.Condition = *condition
attrs, err := nextConditionKeys(g.Condition.Attributes, condition.Attributes)
if err != nil {
return err
}
res, err := nextConditionKeys(g.Condition.Resource, condition.Resource)
if err != nil {
return err
}
g.Condition = SpanMapperGroupCondition{Attributes: attrs, Resource: res}
if err := g.Condition.Validate(); err != nil {
return err
}
}
if enabled != nil {
g.Enabled = *enabled
}
g.UpdatedAt = time.Now()
g.UpdatedBy = updatedBy
return nil
}
func (g *SpanMapperGroup) ErrIfNotDeletable() error {
if g.Origin == SpanMapperOriginSystem {
return errors.Newf(errors.TypeInvalidInput, ErrCodeMappingGroupNotDeletable, "system group %q cannot be deleted, disable it instead", g.Name)
}
return nil
}
func (g *SpanMapperGroup) ToStorable() *StorableSpanMapperGroup {
@@ -100,6 +170,8 @@ func (g *SpanMapperGroup) ToStorable() *StorableSpanMapperGroup {
Name: g.Name,
Condition: g.Condition,
Enabled: g.Enabled,
Origin: g.Origin,
Version: g.Version,
}
}
@@ -112,6 +184,8 @@ func (s *StorableSpanMapperGroup) ToSpanMapperGroup() *SpanMapperGroup {
Name: s.Name,
Condition: s.Condition,
Enabled: s.Enabled,
Origin: s.Origin,
Version: s.Version,
}
}
@@ -126,3 +200,39 @@ func NewSpanMapperGroupsFromStorable(ss []*StorableSpanMapperGroup) []*SpanMappe
func NewGettableSpanMapperGroups(g []*SpanMapperGroup) *GettableSpanMapperGroups {
return &GettableSpanMapperGroups{Items: g}
}
// nextConditionKeys builds a substring list from an edit: user substrings are
// taken from the edit as sent, system substrings stay as stored and the edit
// can only flip their enabled flag.
func nextConditionKeys(stored, edit []SpanMapperGroupConditionKey) ([]SpanMapperGroupConditionKey, error) {
var systemKeys, userKeys []SpanMapperGroupConditionKey
for _, k := range stored {
if k.Origin == SpanMapperOriginSystem {
systemKeys = append(systemKeys, k)
}
}
for _, k := range edit {
if k.Origin != SpanMapperOriginSystem {
k.Origin = SpanMapperOriginUser
userKeys = append(userKeys, k)
continue
}
idx := slices.IndexFunc(systemKeys, func(s SpanMapperGroupConditionKey) bool { return s.Value == k.Value })
if idx == -1 {
return nil, errors.Newf(errors.TypeInvalidInput, ErrCodeMappingInvalidInput, "system substring %q does not exist on this group; only its enabled flag can change", k.Value)
}
systemKeys[idx].Enabled = k.Enabled
}
return append(systemKeys, userKeys...), nil
}
func conditionKeysWithOrigin(keys []SpanMapperGroupConditionKey, origin SpanMapperOrigin) []SpanMapperGroupConditionKey {
out := make([]SpanMapperGroupConditionKey, len(keys))
for i, k := range keys {
k.Origin = origin
out[i] = k
}
return out
}

View File

@@ -0,0 +1,121 @@
package spantypes
import (
"bytes"
"encoding/json"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
)
var ErrCodeMappingDefinitionInvalid = errors.MustNewCode("span_attribute_mapping_definition_invalid")
// ProvisionerIdentity is stamped into created_by/updated_by by the reconciler.
const ProvisionerIdentity = "signoz"
// SpanMapperGroupDefinition is one shipped mapping group. Version is bumped on
// every content change and drives upgrades; the group name is the stable key
// and never changes. Once parsed, every substring and source carries the
// system origin and is enabled.
type SpanMapperGroupDefinition struct {
Version int `json:"version"`
Definition PostableSpanMapperTestGroup `json:"definition"`
}
func (d SpanMapperGroupDefinition) Name() string {
return d.Definition.Name
}
func NewSpanMapperGroupDefinition(raw []byte) (SpanMapperGroupDefinition, error) {
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.DisallowUnknownFields()
var d SpanMapperGroupDefinition
if err := decoder.Decode(&d); err != nil {
return SpanMapperGroupDefinition{}, errors.WrapInvalidInputf(err, ErrCodeMappingDefinitionInvalid, "%s", err.Error())
}
if err := d.validate(); err != nil {
return SpanMapperGroupDefinition{}, err
}
for _, keys := range [][]SpanMapperGroupConditionKey{d.Definition.Condition.Attributes, d.Definition.Condition.Resource} {
for i := range keys {
keys[i].Enabled = true
keys[i].Origin = SpanMapperOriginSystem
}
}
for i := range d.Definition.Mappers {
sources := d.Definition.Mappers[i].Config.Sources
for j := range sources {
sources[j].Enabled = true
sources[j].Origin = SpanMapperOriginSystem
}
}
return d, nil
}
// SpanMapperGroupRegistry holds every definition embedded in the binary, keyed by name.
type SpanMapperGroupRegistry struct {
definitions map[string]SpanMapperGroupDefinition
}
func NewSpanMapperGroupRegistry(definitions []SpanMapperGroupDefinition) (SpanMapperGroupRegistry, error) {
byName := make(map[string]SpanMapperGroupDefinition, len(definitions))
for _, d := range definitions {
if _, dup := byName[d.Name()]; dup {
return SpanMapperGroupRegistry{}, errors.NewInvalidInputf(ErrCodeMappingDefinitionInvalid, "duplicate span mapper group name %q", d.Name())
}
byName[d.Name()] = d
}
return SpanMapperGroupRegistry{definitions: byName}, nil
}
func (r SpanMapperGroupRegistry) IsReserved(name string) bool {
_, ok := r.definitions[name]
return ok
}
// List returns the definitions sorted by name so provisioning order is stable.
func (r SpanMapperGroupRegistry) List() []SpanMapperGroupDefinition {
out := make([]SpanMapperGroupDefinition, 0, len(r.definitions))
for _, d := range r.definitions {
out = append(out, d)
}
slices.SortFunc(out, func(a, b SpanMapperGroupDefinition) int { return strings.Compare(a.Name(), b.Name()) })
return out
}
func (d SpanMapperGroupDefinition) validate() error {
if d.Version < 1 {
return errors.NewInvalidInputf(ErrCodeMappingDefinitionInvalid, "version must be at least 1, got %d", d.Version)
}
if err := d.Definition.Validate(); err != nil {
return errors.Wrapf(err, errors.TypeInvalidInput, ErrCodeMappingDefinitionInvalid, "%s", d.Name())
}
if len(d.Definition.Mappers) == 0 {
return errors.NewInvalidInputf(ErrCodeMappingDefinitionInvalid, "%s: at least one mapper is required", d.Name())
}
names := make(map[string]struct{}, len(d.Definition.Mappers))
for i := range d.Definition.Mappers {
m := &d.Definition.Mappers[i]
if err := m.Validate(); err != nil {
return errors.Wrapf(err, errors.TypeInvalidInput, ErrCodeMappingDefinitionInvalid, "%s: mapper %q", d.Name(), m.Name)
}
if _, dup := names[m.Name]; dup {
return errors.NewInvalidInputf(ErrCodeMappingDefinitionInvalid, "%s: duplicate mapper %q", d.Name(), m.Name)
}
names[m.Name] = struct{}{}
for _, s := range m.Config.Sources {
if !s.Origin.IsZero() || s.Enabled {
return errors.NewInvalidInputf(ErrCodeMappingDefinitionInvalid, "%s: mapper %q: sources must not set origin or enabled", d.Name(), m.Name)
}
}
}
for _, k := range slices.Concat(d.Definition.Condition.Attributes, d.Definition.Condition.Resource) {
if !k.Origin.IsZero() || k.Enabled {
return errors.NewInvalidInputf(ErrCodeMappingDefinitionInvalid, "%s: condition substrings must not set origin or enabled", d.Name())
}
}
return nil
}

View File

@@ -86,17 +86,31 @@ func buildProcessorConfig(groups []*SpanMapperGroupWithMappers) *spanMapperProce
out := make([]spanMapperProcessorGroup, 0, len(groups))
for _, gm := range groups {
existsAny := spanMapperProcessorExistsAny{
Attributes: enabledConditionValues(gm.Group.Condition.Attributes),
Resource: enabledConditionValues(gm.Group.Condition.Resource),
}
// The collector rejects an empty exists_any and empty sources; with
// per-item toggles, all-off is valid stored state and means "never runs".
if len(existsAny.Attributes)+len(existsAny.Resource) == 0 {
continue
}
rules := make([]spanMapperProcessorAttribute, 0, len(gm.Mappers))
for _, m := range gm.Mappers {
rules = append(rules, buildAttributeRule(m))
rule := buildAttributeRule(m)
if len(rule.Sources) == 0 {
continue
}
rules = append(rules, rule)
}
if len(rules) == 0 {
continue
}
out = append(out, spanMapperProcessorGroup{
ID: gm.Group.Name,
ExistsAny: spanMapperProcessorExistsAny{
Attributes: gm.Group.Condition.Attributes,
Resource: gm.Group.Condition.Resource,
},
ID: gm.Group.Name,
ExistsAny: existsAny,
Attributes: rules,
})
}
@@ -104,14 +118,31 @@ func buildProcessorConfig(groups []*SpanMapperGroupWithMappers) *spanMapperProce
return &spanMapperProcessorConfig{Groups: out}
}
func enabledConditionValues(keys []SpanMapperGroupConditionKey) []string {
out := make([]string, 0, len(keys))
for _, k := range keys {
if k.Enabled {
out = append(out, k.Value)
}
}
if len(out) == 0 {
return nil
}
return out
}
// buildAttributeRule maps a single SpanMapper to a collector attribute rule.
// Sources are sorted by Priority DESC (highest-priority first); read-from-
// resource sources are encoded via the "resource." prefix on the key. Each
// source carries its own action — "copy" is omitted to keep the emitted YAML
// compact, and only "move" is set explicitly.
// Disabled sources are skipped and the rest are sorted by Priority DESC
// (highest-priority first); read-from-resource sources are encoded via the
// "resource." prefix on the key. Each source carries its own action — "copy"
// is omitted to keep the emitted YAML compact, and only "move" is set explicitly.
func buildAttributeRule(m *SpanMapper) spanMapperProcessorAttribute {
sources := make([]SpanMapperSource, len(m.Config.Sources))
copy(sources, m.Config.Sources)
sources := make([]SpanMapperSource, 0, len(m.Config.Sources))
for _, s := range m.Config.Sources {
if s.Enabled {
sources = append(sources, s)
}
}
sort.SliceStable(sources, func(i, j int) bool { return sources[i].Priority > sources[j].Priority })
out := make([]spanMapperProcessorSource, 0, len(sources))

View File

@@ -145,6 +145,22 @@ func TestBuildAttributeRule(t *testing.T) {
},
},
},
{
name: "disabled_sources_skipped",
mapper: newMapper("gen_ai.input.messages", FieldContextSpanAttribute,
systemSrc("gen_ai.prompt", SpanMapperOperationCopy, 30, false),
systemSrc("input.value", SpanMapperOperationCopy, 20, true),
attrSrc("gen_ai.prompt", SpanMapperOperationMove, 40),
),
want: spanMapperProcessorAttribute{
Target: "gen_ai.input.messages",
Context: FieldContextSpanAttribute.StringValue(),
Sources: []spanMapperProcessorSource{
{Key: "gen_ai.prompt", Action: SpanMapperOperationMove.StringValue()},
{Key: "input.value"},
},
},
},
}
for _, tc := range tests {
@@ -155,6 +171,33 @@ func TestBuildAttributeRule(t *testing.T) {
}
}
func TestBuildProcessorConfigDropsAllOffItems(t *testing.T) {
t.Parallel()
offGroup := newGroup("all-off", nil, nil)
offGroup.Condition.Attributes = []SpanMapperGroupConditionKey{{Value: "model", Enabled: false, Origin: SpanMapperOriginSystem}}
mixed := newGroup("llm", nil, nil)
mixed.Condition.Attributes = []SpanMapperGroupConditionKey{
{Value: "model", Enabled: false, Origin: SpanMapperOriginSystem},
{Value: "gen_ai.request.model", Enabled: true, Origin: SpanMapperOriginUser},
}
got := buildProcessorConfig([]*SpanMapperGroupWithMappers{
{Group: offGroup, Mappers: []*SpanMapper{newMapper("gen_ai.request.model", FieldContextSpanAttribute, attrSrc("llm.model", SpanMapperOperationCopy, 1))}},
{Group: mixed, Mappers: []*SpanMapper{
newMapper("gen_ai.request.model", FieldContextSpanAttribute, systemSrc("llm.model", SpanMapperOperationCopy, 10, false)),
newMapper("gen_ai.provider.name", FieldContextSpanAttribute, systemSrc("llm.vendor", SpanMapperOperationCopy, 10, true)),
}},
})
require.Len(t, got.Groups, 1)
assert.Equal(t, "llm", got.Groups[0].ID)
assert.Equal(t, []string{"gen_ai.request.model"}, got.Groups[0].ExistsAny.Attributes)
require.Len(t, got.Groups[0].Attributes, 1)
assert.Equal(t, "gen_ai.provider.name", got.Groups[0].Attributes[0].Target)
}
func loadFixture(t *testing.T, name string) []byte {
t.Helper()
b, err := os.ReadFile(filepath.Join("testdata", name))
@@ -174,12 +217,23 @@ func assertYAMLEqual(t *testing.T, want, got []byte) {
func newGroup(name string, attrs, res []string) *SpanMapperGroup {
return &SpanMapperGroup{
Name: name,
Condition: SpanMapperGroupCondition{Attributes: attrs, Resource: res},
Enabled: true,
Name: name,
Condition: SpanMapperGroupCondition{
Attributes: userConditionKeys(attrs),
Resource: userConditionKeys(res),
},
Enabled: true,
}
}
func userConditionKeys(values []string) []SpanMapperGroupConditionKey {
out := make([]SpanMapperGroupConditionKey, len(values))
for i, v := range values {
out[i] = SpanMapperGroupConditionKey{Value: v, Enabled: true, Origin: SpanMapperOriginUser}
}
return out
}
func newMapper(name string, target FieldContext, sources ...SpanMapperSource) *SpanMapper {
return &SpanMapper{
Name: name,
@@ -190,9 +244,13 @@ func newMapper(name string, target FieldContext, sources ...SpanMapperSource) *S
}
func attrSrc(key string, op SpanMapperOperation, priority int) SpanMapperSource {
return SpanMapperSource{Key: key, Context: FieldContextSpanAttribute, Operation: op, Priority: priority}
return SpanMapperSource{Key: key, Context: FieldContextSpanAttribute, Operation: op, Priority: priority, Enabled: true, Origin: SpanMapperOriginUser}
}
func resSrc(key string, op SpanMapperOperation, priority int) SpanMapperSource {
return SpanMapperSource{Key: key, Context: FieldContextResource, Operation: op, Priority: priority}
return SpanMapperSource{Key: key, Context: FieldContextResource, Operation: op, Priority: priority, Enabled: true, Origin: SpanMapperOriginUser}
}
func systemSrc(key string, op SpanMapperOperation, priority int, enabled bool) SpanMapperSource {
return SpanMapperSource{Key: key, Context: FieldContextSpanAttribute, Operation: op, Priority: priority, Enabled: enabled, Origin: SpanMapperOriginSystem}
}

View File

@@ -15,14 +15,14 @@ func TestSimulateSpanMappersProcessing_EndToEnd(t *testing.T) {
groups := []*SpanMapperGroupWithMappers{{
Group: &SpanMapperGroup{
Name: "llm",
Condition: SpanMapperGroupCondition{Attributes: []string{"model"}},
Condition: SpanMapperGroupCondition{Attributes: userConditionKeys([]string{"model"})},
Enabled: true,
},
Mappers: []*SpanMapper{{
Name: "gen_ai.request.model",
FieldContext: FieldContextSpanAttribute,
Config: SpanMapperConfig{Sources: []SpanMapperSource{
{Key: "llm.model", Context: FieldContextSpanAttribute, Operation: SpanMapperOperationCopy, Priority: 1},
{Key: "llm.model", Context: FieldContextSpanAttribute, Operation: SpanMapperOperationCopy, Priority: 1, Enabled: true, Origin: SpanMapperOriginUser},
}},
Enabled: true,
}},

View File

@@ -20,7 +20,9 @@ type StorableSpanMapperGroup struct {
OrgID valuer.UUID `bun:"org_id,type:text,notnull"`
Name string `bun:"name,type:text,notnull"`
Condition SpanMapperGroupCondition `bun:"condition,type:jsonb,notnull"`
Enabled bool `bun:"enabled,notnull,default:true"`
Enabled bool `bun:"enabled,notnull"`
Origin SpanMapperOrigin `bun:"origin,type:text,notnull"`
Version int `bun:"version,notnull"`
}
type StorableSpanMapper struct {
@@ -34,7 +36,8 @@ type StorableSpanMapper struct {
Name string `bun:"name,type:text,notnull"`
FieldContext FieldContext `bun:"field_context,type:text,notnull"`
Config SpanMapperConfig `bun:"config,type:jsonb,notnull"`
Enabled bool `bun:"enabled,notnull,default:true"`
Enabled bool `bun:"enabled,notnull"`
Origin SpanMapperOrigin `bun:"origin,type:text,notnull"`
}
func (c SpanMapperGroupCondition) Value() (driver.Value, error) {

View File

@@ -9,9 +9,14 @@ import (
)
type SpanMapperStore interface {
// RunInTx runs cb in one transaction; every store call made with the
// callback's ctx joins it.
RunInTx(ctx context.Context, cb func(ctx context.Context) error) error
// Group operations
ListGroups(ctx context.Context, orgID valuer.UUID, q *ListSpanMapperGroupsQuery) ([]*SpanMapperGroup, error)
GetGroup(ctx context.Context, orgID, id valuer.UUID) (*SpanMapperGroup, error)
GetGroupByName(ctx context.Context, orgID valuer.UUID, name string) (*SpanMapperGroup, error)
CreateGroup(ctx context.Context, group *SpanMapperGroup) error
UpdateGroup(ctx context.Context, group *SpanMapperGroup) error
DeleteGroup(ctx context.Context, orgID, id valuer.UUID) error
@@ -21,7 +26,7 @@ type SpanMapperStore interface {
GetMapper(ctx context.Context, orgID, groupID, id valuer.UUID) (*SpanMapper, error)
CreateMapper(ctx context.Context, mapper *SpanMapper) error
UpdateMapper(ctx context.Context, mapper *SpanMapper) error
DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID) error
DeleteMapper(ctx context.Context, orgID, groupID, id valuer.UUID, origin SpanMapperOrigin) error
}
// TraceStore defines the data access interface for trace detail queries.

View File

@@ -41,7 +41,7 @@ def test_create_groups_and_simulate_with_backfill(
},
json={
"name": "llm-backfill",
"condition": {"attributes": ["model"], "resource": []},
"condition": {"attributes": [{"value": "model", "enabled": True}], "resource": []},
"enabled": True,
},
)
@@ -69,6 +69,7 @@ def test_create_groups_and_simulate_with_backfill(
"context": "attribute",
"operation": "copy",
"priority": 1,
"enabled": True,
}
]
},
@@ -126,13 +127,13 @@ def test_create_groups_and_simulate_with_backfill(
# No "mappers" key: the server backfills them from the saved group.
{
"name": "llm-backfill",
"condition": {"attributes": ["model"], "resource": []},
"condition": {"attributes": [{"value": "model", "enabled": True}], "resource": []},
"enabled": True,
},
# Unsaved group; mappers provided inline.
{
"name": "db-inline",
"condition": {"attributes": ["db"], "resource": []},
"condition": {"attributes": [{"value": "db", "enabled": True}], "resource": []},
"enabled": True,
"mappers": [
{
@@ -145,6 +146,7 @@ def test_create_groups_and_simulate_with_backfill(
"context": "attribute",
"operation": "move",
"priority": 1,
"enabled": True,
}
]
},

View File

@@ -0,0 +1,152 @@
from collections.abc import Callable
from http import HTTPStatus
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
GROUPS_PATH = "/api/v1/span_mapper_groups"
def test_default_groups_are_seeded_and_shipped_items_are_toggle_only(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
) -> None:
"""
Setup:
A fresh org. The reconciler seeds the shipped mapping groups at startup
and on org creation, so nothing has to be created here.
Tests:
1. The list contains llm, agent and tool as system groups with shipped
substrings
2. Shipped mappers and their sources are system-owned and enabled
3. A shipped name cannot be taken by a user group, and system groups and
mappers cannot be deleted
4. A shipped source can be switched off and a user override added; both
round-trip through PATCH and the simulator honours them
5. A shipped substring can be switched off and a user one added
"""
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
headers = {"authorization": f"Bearer {token}", "content-type": "application/json"}
list_groups = requests.get(signoz.self.host_configs["8080"].get(GROUPS_PATH), timeout=10, headers=headers)
assert list_groups.status_code == HTTPStatus.OK
groups = {g["name"]: g for g in list_groups.json()["data"]["items"]}
assert {"gen_ai.llm", "gen_ai.agent", "gen_ai.tool"} <= set(groups)
for name in ("gen_ai.llm", "gen_ai.agent", "gen_ai.tool"):
assert groups[name]["origin"] == "system"
assert groups[name]["version"] >= 1
assert groups[name]["createdBy"] == "signoz"
llm = groups["gen_ai.llm"]
assert llm["condition"]["attributes"] == [{"value": "model", "enabled": True, "origin": "system"}]
list_mappers = requests.get(signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}/span_mappers"), timeout=10, headers=headers)
assert list_mappers.status_code == HTTPStatus.OK
mappers = {m["name"]: m for m in list_mappers.json()["data"]["items"]}
model = mappers["gen_ai.request.model"]
assert model["origin"] == "system"
assert model["enabled"] is True
assert all(s["origin"] == "system" and s["enabled"] is True for s in model["config"]["sources"])
assert "llm.model_name" in [s["key"] for s in model["config"]["sources"]]
reserved = requests.post(
signoz.self.host_configs["8080"].get(GROUPS_PATH),
timeout=10,
headers=headers,
json={"name": "gen_ai.tool", "condition": {"attributes": [{"value": "tool", "enabled": True}], "resource": []}, "enabled": True},
)
assert reserved.status_code == HTTPStatus.BAD_REQUEST
assert reserved.json()["error"]["code"] == "span_attribute_mapping_group_name_reserved"
delete_group = requests.delete(signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}"), timeout=10, headers=headers)
assert delete_group.status_code == HTTPStatus.BAD_REQUEST
assert delete_group.json()["error"]["code"] == "span_attribute_mapping_group_not_deletable"
delete_mapper = requests.delete(signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}/span_mappers/{model['id']}"), timeout=10, headers=headers)
assert delete_mapper.status_code == HTTPStatus.BAD_REQUEST
assert delete_mapper.json()["error"]["code"] == "span_attribute_mapper_not_deletable"
# Switch the shipped llm.model_name source off and re-add it as a user move.
sources = [{**s, "enabled": s["key"] != "llm.model_name"} for s in model["config"]["sources"]] + [{"key": "llm.model_name", "context": "attribute", "operation": "move", "priority": 1, "enabled": True}]
patch_mapper = requests.patch(
signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}/span_mappers/{model['id']}"),
timeout=10,
headers=headers,
json={"config": {"sources": sources}},
)
assert patch_mapper.status_code == HTTPStatus.NO_CONTENT
list_mappers = requests.get(signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}/span_mappers"), timeout=10, headers=headers)
updated = {m["name"]: m for m in list_mappers.json()["data"]["items"]}["gen_ai.request.model"]
by_origin = {(s["key"], s["origin"]): s for s in updated["config"]["sources"]}
assert by_origin[("llm.model_name", "system")]["enabled"] is False
assert by_origin[("llm.model_name", "user")]["operation"] == "move"
assert len(updated["config"]["sources"]) == len(model["config"]["sources"]) + 1
# The user override wins: the source is moved, not copied.
simulate = requests.post(
signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/test"),
timeout=10,
headers=headers,
json={
"spans": [{"attributes": {"llm.model_name": "gpt-4o"}, "resource": {}}],
"groups": [{"name": "gen_ai.llm", "condition": llm["condition"], "enabled": True}],
},
)
assert simulate.status_code == HTTPStatus.OK
attrs = simulate.json()["data"]["spans"][0]["attributes"]
assert attrs["gen_ai.request.model"] == "gpt-4o"
assert "llm.model_name" not in attrs
# A shipped substring that does not exist is rejected; toggling one is not.
bad_condition = requests.patch(
signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}"),
timeout=10,
headers=headers,
json={"condition": {"attributes": [{"value": "nope", "enabled": True, "origin": "system"}], "resource": []}},
)
assert bad_condition.status_code == HTTPStatus.BAD_REQUEST
patch_group = requests.patch(
signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}"),
timeout=10,
headers=headers,
json={
"condition": {
"attributes": [
{"value": "model", "enabled": False, "origin": "system"},
{"value": "gen_ai.request.model", "enabled": True},
],
"resource": [],
}
},
)
assert patch_group.status_code == HTTPStatus.NO_CONTENT
list_groups = requests.get(signoz.self.host_configs["8080"].get(GROUPS_PATH), timeout=10, headers=headers)
llm_after = {g["name"]: g for g in list_groups.json()["data"]["items"]}["gen_ai.llm"]
assert llm_after["condition"]["attributes"] == [
{"value": "model", "enabled": False, "origin": "system"},
{"value": "gen_ai.request.model", "enabled": True, "origin": "user"},
]
assert llm_after["origin"] == "system"
assert llm_after["name"] == "gen_ai.llm"
# Leave the shipped group as seeded for the other suites.
restore_group = requests.patch(
signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}"),
timeout=10,
headers=headers,
json={"condition": {"attributes": [{"value": "model", "enabled": True, "origin": "system"}], "resource": []}},
)
assert restore_group.status_code == HTTPStatus.NO_CONTENT
restore_mapper = requests.patch(
signoz.self.host_configs["8080"].get(f"{GROUPS_PATH}/{llm['id']}/span_mappers/{model['id']}"),
timeout=10,
headers=headers,
json={"config": {"sources": model["config"]["sources"]}},
)
assert restore_mapper.status_code == HTTPStatus.NO_CONTENT