Compare commits

...

14 Commits

Author SHA1 Message Date
Nikhil Soni
2177206b35 fix(querier): decode span attributes to flat dotted keys
The span attribute bag is flattened and merged with the legacy maps, so a
key stored as both a scalar and an object (e.g. `scope` and
`scope.attributes.name`) must stay two distinct dotted keys. NestedJSON
collapsed them and surfaced a mislabeled top-level key that no legacy value
overwrote. Decode the attributes column with FlattenJSON (ValuesByPath); the
log body keeps NestedJSON for its nested contract.

Assisted-by: Claude Opus 4.8
2026-09-29 19:09:03 +05:30
Nikhil Soni
ee1aacc32b fix(telemetrystore): decode JSON via MarshalJSON to keep nested arrays
The manual NestedMap/Variant walk collapsed arrays of objects to empty
maps (education: [{}]). Decode through the driver's MarshalJSON plus sonic
instead, which preserves arrays of objects and matches the prior
string-read semantics (float64 numbers, nested tree).

Assisted-by: Claude Opus 4.8
2026-09-29 18:32:55 +05:30
Nikhil Soni
7cd0135fa4 refactor(telemetrystore): drop WrapRows and JSONValue
The driver already reports chcol.JSON as the up-front scan type for JSON
columns under the flattened native serialization, so the WrapRows scan-type
override (and its test conn wrapper) is dead weight; verified across
JSON(mdp=0), JSON(mdp=0, message String) and JSON(mdp=100). Decode with a
plain map[string]any instead of the JSONValue alias.

Assisted-by: Claude Opus 4.8
2026-09-29 18:06:58 +05:30
Nikhil Soni
33db5c467c chore: go mod tidy
sonic is no longer imported directly after the JSON scan refactor.

Assisted-by: Claude Opus 4.8
2026-09-29 17:49:50 +05:30
Nikhil Soni
dcdf63f823 refactor(telemetrystore): scan JSON columns as chcol.JSON
Read JSON columns into chcol.JSON and decode to a nested document, so a
typed sub-path (body_v2.message) no longer errors under the flattened
native serialization. Kept nested, not flattened to dotted paths: the
query builder body contract is nested, so a key stored as both a scalar
and an object still collapses here.

Assisted-by: Claude Opus 4.8
2026-09-29 17:39:57 +05:30
Nikhil Soni
aecbb593bd revert: drop the chcol.JSON scan refactor, keep only the setting swap
Reading JSON in the flattened native serialization lets the driver decode it into the existing JSONValue map, so removing output_format_native_write_json_as_string needs no scan-side changes. The query builder still returns the JSON nested (its scalar-and-object collapse is handled separately for the waterfall).

Assisted-by: Claude Opus 4.8
2026-09-29 15:34:55 +05:30
Nikhil Soni
6a4e412a47 refactor(telemetrystore): drop JSONValue in favour of map[string]any
JSONValue was a thin named map used only as a type-switch discriminator; FlattenJSON now returns map[string]any and the consumers assert it directly.

Assisted-by: Claude Opus 4.8
2026-09-29 15:14:33 +05:30
Nikhil Soni
90be4a498d feat(telemetrystore): read JSON columns natively to keep scalar-and-object keys
Swap the read connection from output_format_native_write_json_as_string=1 to output_format_native_use_flattened_dynamic_and_json_serialization=1 and scan JSON columns as chcol.JSON flattened to dotted paths (via WrapRows + unwrapVariant), instead of a map the driver collapses. A key stored as both a scalar and an object is no longer dropped in logs or the trace list view. JSONValue is now a flat result type; its Scan is removed.

Assisted-by: Claude Opus 4.8
2026-09-29 14:49:20 +05:30
praneeth-signoz
47dd1fabf3 chore(channel-specs): Move channel specs to separate files (#12989)
Some checks failed
build-staging / prepare (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
<!--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

- Moved channel specs to separate files under alert manager types
- Channel receivers are also moved the same file

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR
Closes
https://github.com/orgs/SigNoz/projects/34/views/26?pane=issue&itemId=252814150&issue=SigNoz%7Cpulse-pod%7C374
2026-09-28 17:44:35 +00:00
Vikrant Gupta
270988fb48 fix(tokenizer): persist last_observed_at on postgres (#12982)
#### Description

- The flush CTE rendered `last_observed_at` as an untyped literal, which
postgres resolves to `text` and refuses to assign to the `timestamptz`
column. The column never populated, so the idle expiry never applied.
- Build the CTE from the token model with only `id`, `last_observed_at`
and `updated_at`, so bun casts per dialect and no token secrets land in
the statement.
- Flush now applies cached times through `Token.UpdateLastObservedAt`,
which also skips rows with a newer stored value.
- Integration test in `passwordauthn` runs with a short GC interval and
asserts the column populates on both sql stores.
2026-09-28 14:23:01 +00:00
Naman Verma
39badeb591 fix: remove rules types package import from migration #049 (#12998)
Some checks failed
Release Drafter / update_release_draft (push) Has been cancelled
build-staging / prepare (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
cacheci / tests (push) Has been cancelled
<!--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

Migration package ideally should have have things from types package
imported. Given that rules v1->v2 will (most probably) update/remove
some of the types, such as removing `PreferredChannels` from
`PostableRule`, better not to have this type imported in migrations
package.

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

Part of https://github.com/SigNoz/pulse-pod/issues/225
2026-09-28 11:39:07 +00:00
Nityananda Gohain
ec05bfe755 fix: empty patterns in pricing are rejected (#12995)
<!--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
Empty patterns were not rejected because of which corrupt config was
created, rejecting them at the handler layer.

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR
Part of https://github.com/SigNoz/nerve-pod/issues/282
2026-09-28 08:09:38 +00:00
Aditya Singh
ed1bf7ab89 fix(bottom-strip): size pages from the layout instead of the viewport (#12940)
Some checks failed
build-staging / prepare (push) Has been cancelled
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
#### Description

- pages that hardcoded `100vh` minus a guess at what sits above them
came out taller
than the pane they live in, which showed up as scroll that should not be
there. they
  now take what the layout gives them.
- most of the `100vh` in the app turned out to be harmless.. either
flex-shrink absorbs
it or the pane scrolls anyway. those are left alone, only the ones with
a real symptom
  are changed here.
- alert rules and triggered alerts also needed the AlertList tabs chain
to hand height
  down, that page uses antd `Tabs` directly instead of `RouteTab`.
- licenses, status and support pages had `max-height: 100vh` with
`overflow: hidden` and
no inner scroller, so anything past a viewport was clipped with no way
to reach it.
  removed the cap on all three.

#### Issues closed by this PR

Part of https://github.com/SigNoz/engineering-pod/issues/6074


#### Screenshots/ Screen recording

Home page


https://github.com/user-attachments/assets/8ba3e3c9-1959-4393-b503-ddf579e9d139

Without bottom strip


https://github.com/user-attachments/assets/2ee08181-710a-4c34-a4e7-893c4a320453



Status page

<img width="1728" height="1000" alt="status"
src="https://github.com/user-attachments/assets/f44ec7d7-4ee0-4866-9a90-fe25622fe25b"
/>

Without bottom strip

<img width="1728" height="997" alt="status2"
src="https://github.com/user-attachments/assets/be57dce6-65ca-4450-8f0d-0337625a1d64"
/>

Alert rules


https://github.com/user-attachments/assets/7cc25f56-54e4-4489-bdb4-453409151bac

Triggered alerts


https://github.com/user-attachments/assets/2d8e1885-7cad-4659-90f6-59960f4afce8

Without bottom strip



https://github.com/user-attachments/assets/5247f6c0-3675-4910-b1e4-4076bf93c16e



Support
<img width="1728" height="997" alt="support"
src="https://github.com/user-attachments/assets/c2b3c0b2-2abd-4538-bb52-2660e75817b3"
/>

Trace funnel



https://github.com/user-attachments/assets/c4872135-3cb5-4348-b6c6-2d3e4dabdc1b

Without bottom strip


https://github.com/user-attachments/assets/aef94d5d-02e5-4154-995b-76819d55c2d8


Trace details



https://github.com/user-attachments/assets/24d701b0-296f-4fb4-ba9d-615a47b33285

Without bottom strip


https://github.com/user-attachments/assets/23083f20-67bf-4b06-aaec-aba390a6f594



#### Additional Information

- every page here was checked on screen before and after. the ones left
untouched
(infra hosts/k8s, traces + llm explorer list views, llm settings tables,
the k8s logs
  drawer) were checked too and are fine.
2026-09-25 15:14:44 +00:00
Aditya Singh
6e979c8318 feat(bottom-strip): add the layout shell behind a feature flag (#12936)
#### Description

- adds the bottom strip to the app layout behind a localStorage flag.
shows the build
version on the left for now.. right side actions and the per page count
come in the
  next tickets.
- `.app-content` is a column flex now and `LayoutContent` takes the
height left over
instead of `height: 100%`, so the strip has a stable box to sit under.
this is the
  only bit not behind the flag.
- fixed bottom elements read `--bottom-strip-height`. the var only
exists while the
strip is mounted, so with the flag off everything falls back to where it
is today.
- hides nothing. each later ticket hides the piece it replaces.

#### Issues closed by this PR

Part of https://github.com/SigNoz/engineering-pod/issues/6074

<img width="3084" height="1566" alt="image"
src="https://github.com/user-attachments/assets/b1821fda-5c33-40e7-926a-5d91fedb797e"
/>


#### Additional Information

- pages that still hardcode `100vh` (infra hosts/k8s, trace details,
traces and llm
list views) push the strip off screen. that is the next PR on this
ticket.
- pylon chat window offset is not here.. needs a pylon enabled tenant to
verify so it
  goes with the right side actions ticket.
2026-09-25 12:04:21 +00:00
56 changed files with 1434 additions and 1058 deletions

View File

@@ -47,4 +47,5 @@ export enum LOCALSTORAGE {
DASHBOARDS_LIST_VIEWS = 'DASHBOARDS_LIST_VIEWS',
DASHBOARD_V2_PANEL_COLUMN_WIDTHS = 'DASHBOARD_V2_PANEL_COLUMN_WIDTHS',
LLM_ATTRIBUTE_MAPPING_TEST_SPAN = 'LLM_ATTRIBUTE_MAPPING_TEST_SPAN',
SAVED_VIEW_ENABLED = 'SAVED_VIEW_ENABLED',
}

View File

@@ -53,6 +53,10 @@
z-index: 0;
background: var(--l1-background);
// Column so the bottom strip sits under the scrolling content, not inside it.
display: flex;
flex-direction: column;
&.full-screen-content {
width: 100%;
}
@@ -70,7 +74,9 @@
.chat-support-gateway {
position: fixed;
bottom: 20px;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: calc(20px + var(--bottom-strip-height, 0px));
right: 20px;
z-index: 1000;

View File

@@ -43,6 +43,7 @@ import { USER_PREFERENCES } from 'constants/userPreferences';
import AIAssistantModal from 'container/AIAssistant/AIAssistantModal';
import AIAssistantPanel from 'container/AIAssistant/AIAssistantPanel';
import { useAIAssistantStore } from 'container/AIAssistant/store/useAIAssistantStore';
import BottomStrip from 'container/BottomStrip';
import SideNav from 'container/SideNav';
import TopNav from 'container/TopNav';
import dayjs from 'dayjs';
@@ -51,6 +52,7 @@ import { useIsDarkMode } from 'hooks/useDarkMode';
import { useGetTenantLicense } from 'hooks/useGetTenantLicense';
import { useIsAIAssistantEnabled } from 'hooks/useIsAIAssistantEnabled';
import { useNotifications } from 'hooks/useNotifications';
import { useSavedViewEnabled } from 'hooks/useSavedViewEnabled';
import useTabVisibility from 'hooks/useTabFocus';
import history from 'lib/history';
import { isNull } from 'lodash-es';
@@ -402,6 +404,7 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
}, [pathname]);
const isToDisplayLayout = isLoggedIn;
const isSavedViewEnabled = useSavedViewEnabled();
const routeKey = useMemo(() => getRouteKey(pathname), [pathname]);
const pageTitle = t(routeKey);
@@ -868,6 +871,10 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
</OverlayScrollbar>
</LayoutContent>
</Sentry.ErrorBoundary>
{isSavedViewEnabled && isToDisplayLayout && !renderFullScreen && (
<BottomStrip />
)}
</div>
{isLoggedIn && isAIAssistantEnabled && (

View File

@@ -12,8 +12,12 @@ export const Layout = styled(LayoutComponent)`
}
`;
// Takes the height left in `.app-content` after the bottom strip.
// `min-height: 0` is not needed right now, overlayscrollbars already sets
// `overflow: auto` here. Kept so this does not break if that goes away.
export const LayoutContent = styled(LayoutComponent.Content)`
height: 100%;
flex: 1;
min-height: 0;
&::-webkit-scrollbar {
width: 0.1rem;
}

View File

@@ -0,0 +1,36 @@
.strip {
display: flex;
align-items: center;
justify-content: space-between;
gap: var(--spacing-6);
flex-shrink: 0;
height: var(--bottom-strip-height);
padding: 0 var(--spacing-6);
background: var(--l2-background);
border-top: 1px solid var(--l2-border);
font-family: var(--font-family-sf-mono, monospace);
// Above page content, below the body-portalled overlays that are meant to
// cover the strip.
position: relative;
z-index: 1;
}
.left,
.right {
display: flex;
align-items: center;
gap: var(--spacing-6);
min-width: 0;
}
// Temporary placeholder for the left slot. Replaced later.
.version {
color: var(--l2-foreground);
white-space: nowrap;
overflow: hidden;
text-overflow: ellipsis;
}

View File

@@ -0,0 +1,49 @@
import { render } from 'tests/test-utils';
import BottomStrip, {
BOTTOM_STRIP_HEIGHT,
BOTTOM_STRIP_HEIGHT_VAR,
BOTTOM_STRIP_ON_CLASS,
} from '..';
describe('BottomStrip', () => {
it('publishes the body class and height property while mounted', () => {
const { unmount } = render(<BottomStrip />);
expect(document.body.classList.contains(BOTTOM_STRIP_ON_CLASS)).toBe(true);
expect(document.body.style.getPropertyValue(BOTTOM_STRIP_HEIGHT_VAR)).toBe(
`${BOTTOM_STRIP_HEIGHT}px`,
);
unmount();
expect(document.body.classList.contains(BOTTOM_STRIP_ON_CLASS)).toBe(false);
expect(document.body.style.getPropertyValue(BOTTOM_STRIP_HEIGHT_VAR)).toBe(
'',
);
});
// The string is whatever the Go build injected, so it is rendered untouched —
// same as SideNav. Release tags carry the "v", local builds do not.
it.each([['v0.134.67'], ['main-64f1c2a']])(
'renders the build version %p exactly as given',
(version) => {
const { getByTestId } = render(<BottomStrip />, undefined, {
appContextOverrides: {
versionData: { version, ee: 'Y', setupCompleted: true },
},
});
expect(getByTestId('bottom-strip-version')).toHaveTextContent(version);
},
);
it('renders the strip without a version when none is available', () => {
const { getByTestId, queryByTestId } = render(<BottomStrip />, undefined, {
appContextOverrides: { versionData: null },
});
expect(getByTestId('bottom-strip')).toBeInTheDocument();
expect(queryByTestId('bottom-strip-version')).not.toBeInTheDocument();
});
});

View File

@@ -0,0 +1,42 @@
import { useLayoutEffect } from 'react';
import { useAppContext } from 'providers/App/App';
import styles from './BottomStrip.module.scss';
export const BOTTOM_STRIP_HEIGHT = 24;
export const BOTTOM_STRIP_ON_CLASS = 'bottom-strip-on';
export const BOTTOM_STRIP_HEIGHT_VAR = '--bottom-strip-height';
function BottomStrip(): JSX.Element {
const { versionData } = useAppContext();
const version = versionData?.version?.trim();
useLayoutEffect(() => {
document.body.classList.add(BOTTOM_STRIP_ON_CLASS);
document.body.style.setProperty(
BOTTOM_STRIP_HEIGHT_VAR,
`${BOTTOM_STRIP_HEIGHT}px`,
);
return (): void => {
document.body.classList.remove(BOTTOM_STRIP_ON_CLASS);
document.body.style.removeProperty(BOTTOM_STRIP_HEIGHT_VAR);
};
}, []);
return (
<div className={styles.strip} data-testid="bottom-strip">
<div className={styles.left}>
{version && (
<span className={styles.version} data-testid="bottom-strip-version">
{version}
</span>
)}
</div>
<div className={styles.right} />
</div>
);
}
export default BottomStrip;

View File

@@ -1,6 +1,8 @@
.create-alert-v2-footer {
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
left: 63px;
right: 0;
background-color: var(--l1-background);

View File

@@ -1,6 +1,8 @@
.explorer-options-container {
position: fixed;
bottom: 0px;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
left: calc(50% + 240px);
transform: translate(calc(-50% - 120px), 0);
transition: left 0.2s linear;

View File

@@ -1,6 +1,8 @@
.explorer-option-droppable-container {
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
width: -webkit-fill-available;
height: 24px;
display: flex;

View File

@@ -1,7 +1,6 @@
.home-container {
display: flex;
flex-direction: column;
min-height: 100vh;
overflow-y: auto;
height: 100%;
width: 100%;

View File

@@ -1,7 +1,4 @@
.licenses-page {
max-height: 100vh;
overflow: hidden;
.licenses-page-header {
border-bottom: 1px solid var(--l1-border);
background: var(--l1-background);
@@ -32,7 +29,6 @@
.licenses-page-content {
flex: 1;
height: calc(100vh - 48px);
background: var(--l1-background);
padding: 10px 8px;
overflow-y: auto;

View File

@@ -2,7 +2,7 @@
display: flex;
flex-direction: column;
gap: 1rem;
height: calc(100vh - 62px);
flex: 1;
min-height: 400px;
}

View File

@@ -181,7 +181,9 @@
.ant-pagination {
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new
// fixed-bottom UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
width: calc(100% - 54px);
background: var(--l1-background);
padding: 16px;

View File

@@ -2,7 +2,7 @@
display: flex;
flex-direction: column;
gap: 1rem;
height: calc(100vh - 62px);
flex: 1;
min-height: 400px;
padding-top: var(--spacing-8);
}

View File

@@ -1,7 +1,4 @@
.version-container {
max-height: 100vh;
overflow: hidden;
.version-page-header {
border-bottom: 1px solid var(--l1-border);
background: var(--l1-background);

View File

@@ -0,0 +1,11 @@
import getLocalStorageKey from 'api/browser/localstorage/get';
import { LOCALSTORAGE } from 'constants/localStorage';
import { useState } from 'react';
export function useSavedViewEnabled(): boolean {
const [isEnabled] = useState(
() => getLocalStorageKey(LOCALSTORAGE.SAVED_VIEW_ENABLED) === 'true',
);
return isEnabled;
}

View File

@@ -1,4 +1,29 @@
.alerts-container {
// Hands the page height down to the active tab so its content can bound itself
// instead of guessing with 100vh. Child combinators only, nested Tabs
// (Configuration) must not be caught.
flex: 1;
min-height: 0;
> .ant-tabs-content-holder {
display: flex;
flex-direction: column;
> .ant-tabs-content {
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
> .ant-tabs-tabpane-active {
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}
}
}
.top-level-tab.periscope-tab {
padding: 2px 0;
}
@@ -40,5 +65,9 @@
.alert-rules-container {
margin-top: 10px;
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}
}

View File

@@ -2,7 +2,9 @@
display: flex;
flex-direction: column;
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
left: 0;
width: 100%;
z-index: 100;

View File

@@ -1,7 +1,4 @@
.support-page-container {
max-height: 100vh;
overflow: hidden;
.support-page-header {
border-bottom: 1px solid var(--l1-border);
background: var(--l1-background);

View File

@@ -1,5 +1,6 @@
.root {
height: calc(100vh);
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}

View File

@@ -1,13 +1,24 @@
.traces-funnel-details {
display: flex;
// 45px -> height of the tab bar
height: calc(100vh - 45px);
height: 100%;
&__steps-config {
flex-shrink: 0;
width: 600px;
border-right: 1px solid var(--l1-border);
// Positioning context for the absolute .steps-footer.
position: relative;
display: flex;
flex-direction: column;
// Scoped here so the modal usage of FunnelConfiguration on trace details
// stays in normal flow.
.funnel-configuration {
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}
}
&__steps-results {
width: 100%;

View File

@@ -4,14 +4,17 @@
flex-direction: column;
justify-content: flex-start;
&.funnel-details-page {
height: calc(
100vh - 170px
); // 64px bottom bar + 61px configuration header + 45px page navbar
flex: 1;
min-height: 0;
// .steps-footer is absolute against the config column, so its 64px is
// reserved rather than laid out.
margin-bottom: 64px;
overflow: auto;
}
}
&__header {
flex-shrink: 0;
display: flex;
align-items: center;
justify-content: space-between;

View File

@@ -40,11 +40,15 @@ func stripKeyAlias(name string) string {
return keyAliasRe.ReplaceAllString(name, "")
}
// unwrapVariant returns the concrete value inside the chcol.Variant envelope the driver scans a
// Dynamic column — a JSON path such as body_v2.level — into.
// unwrapVariant returns the concrete value inside the driver's scan envelopes: chcol.Variant for a
// Dynamic column (a JSON path such as body_v2.level), and chcol.JSON for a whole JSON column, decoded
// into a nested document.
func unwrapVariant(val any) any {
if v, ok := val.(chcol.Variant); ok {
switch v := val.(type) {
case chcol.Variant:
return v.Any()
case chcol.JSON:
return telemetrystoretypes.NestedJSON(v)
}
return val
}
@@ -58,7 +62,7 @@ func labelValue(val any) string {
if val == nil {
return ""
}
if v, ok := val.(telemetrystoretypes.JSONValue); ok {
if v, ok := val.(map[string]any); ok {
if raw, err := json.Marshal(v); err == nil {
return string(raw)
}
@@ -204,7 +208,7 @@ func readAsTimeSeries(rows driver.Rows, queryWindow *qbtypes.TimeRange, step qbt
Value: *val,
})
case *telemetrystoretypes.JSONValue, *chcol.Variant:
case *chcol.JSON, *chcol.Variant:
val := labelValue(derefValue(ptr))
lblVals = append(lblVals, val)
lblObjs = append(lblObjs, &qbtypes.Label{
@@ -536,7 +540,14 @@ func readAsRaw(rows driver.Rows, queryName string) (*qbtypes.RawData, error) {
name := stripKeyAlias(colNames[i])
// de-reference the typed pointer to any
val := unwrapVariant(reflect.ValueOf(cellPtr).Elem().Interface())
raw := reflect.ValueOf(cellPtr).Elem().Interface()
// the attributes bag is flattened to dotted keys downstream; decode it flat so a key stored as both a scalar and an object is not collapsed into a mislabeled key.
var val any
if j, ok := raw.(chcol.JSON); ok && name == "attributes" {
val = telemetrystoretypes.FlattenJSON(j)
} else {
val = unwrapVariant(raw)
}
// special-case: timestamp column
if name == "timestamp" || name == "timestamp_datetime" {
@@ -576,8 +587,6 @@ func flattenJSONPaths(prefix string, m map[string]any, out map[string]any) {
switch child := v.(type) {
case map[string]any:
flattenJSONPaths(key, child, out)
case telemetrystoretypes.JSONValue:
flattenJSONPaths(key, child, out)
default:
out[key] = v
}
@@ -593,7 +602,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
attrStr, hasStr := data["attributes_string"]
attrNum, hasNum := data["attributes_number"]
attrBool, hasBool := data["attributes_bool"]
attrJSON, _ := data["attributes"].(telemetrystoretypes.JSONValue)
attrJSON, _ := data["attributes"].(map[string]any)
// todo(nitya): move to resource json
resStr, hasRes := data["resources_string"]
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {

View File

@@ -3,14 +3,9 @@ package querier
import (
"reflect"
"testing"
"time"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
cmock "github.com/SigNoz/clickhouse-go-mock"
"github.com/SigNoz/signoz/pkg/telemetrystore"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -83,101 +78,31 @@ func TestMergeSpanAttributeColumns_ParsesEventsAndLinks(t *testing.T) {
}
}
// A ClickHouse query can put a JSON column in the result of any request type — e.g.
// `select * from signoz_logs.logs_v2` on a body_v2 stack, where `*` covers body_v2.
func TestConsume_JSONColumn(t *testing.T) {
ts := time.Date(2026, 8, 14, 10, 0, 0, 0, time.UTC)
body := `{"level":"error","attrs":{"code":500}}`
wantBody := telemetrystoretypes.JSONValue{
"level": "error",
"attrs": map[string]any{"code": float64(500)},
}
// the scalar reader reuses its scan slots across rows, so each row must still carry its own body
t.Run("scalar", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "body_v2", Type: "JSON"},
{Name: "__result_0", Type: "UInt64"},
}, [][]any{{body, uint64(3)}, {`{"level":"warn"}`, uint64(1)}}))
payload, err := consume(rows, qbtypes.RequestTypeScalar, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.ScalarData)
require.Len(t, data.Data, 2)
assert.Equal(t, wantBody, data.Data[0][0])
assert.Equal(t, uint64(3), data.Data[0][1])
assert.Equal(t, telemetrystoretypes.JSONValue{"level": "warn"}, data.Data[1][0])
assert.Equal(t, uint64(1), data.Data[1][1])
})
t.Run("time series", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "ts", Type: "DateTime"},
{Name: "body_v2", Type: "JSON"},
{Name: "__result_0", Type: "UInt64"},
}, [][]any{{ts, body, uint64(3)}}))
payload, err := consume(rows, qbtypes.RequestTypeTimeSeries, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.TimeSeriesData)
require.Len(t, data.Aggregations, 1)
require.Len(t, data.Aggregations[0].Series, 1)
require.Len(t, data.Aggregations[0].Series[0].Values, 1)
assert.Equal(t, float64(3), data.Aggregations[0].Series[0].Values[0].Value)
})
// grouping by a JSON column is legal in ClickHouse, so each document has to label its own
// series rather than being dropped, which would merge every group into one
t.Run("time series grouped by the JSON column", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "ts", Type: "DateTime"},
{Name: "body_v2", Type: "JSON"},
{Name: "__result_0", Type: "UInt64"},
}, [][]any{
{ts, `{"level":"error"}`, uint64(7)},
{ts, `{"level":"warn"}`, uint64(2)},
}))
payload, err := consume(rows, qbtypes.RequestTypeTimeSeries, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.TimeSeriesData)
require.Len(t, data.Aggregations, 1)
require.Len(t, data.Aggregations[0].Series, 2)
got := map[string]float64{}
for _, series := range data.Aggregations[0].Series {
require.Len(t, series.Labels, 1)
require.Len(t, series.Values, 1)
got[series.Labels[0].Value.(string)] = series.Values[0].Value
}
assert.Equal(t, map[string]float64{`{"level":"error"}`: 7, `{"level":"warn"}`: 2}, got)
})
t.Run("raw", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "timestamp", Type: "DateTime"},
{Name: "body_v2", Type: "JSON"},
}, [][]any{{ts, body}}))
payload, err := consume(rows, qbtypes.RequestTypeRaw, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.RawData)
require.Len(t, data.Rows, 1)
assert.Equal(t, ts, data.Rows[0].Timestamp.UTC())
assert.Equal(t, wantBody, data.Rows[0].Data["body_v2"])
})
}
// A JSON path (e.g. `body_v2.level`) comes back as a Dynamic column, which the driver scans
// into a chcol.Variant envelope rather than the value itself.
// A JSON path (e.g. `body_v2.level`) comes back as a Dynamic column, which the driver scans into a
// chcol.Variant envelope; a whole JSON column comes back as chcol.JSON, decoded into a nested document.
func TestUnwrapVariant(t *testing.T) {
assert.Equal(t, "error", unwrapVariant(chcol.NewDynamicWithType("error", "String")))
assert.Nil(t, unwrapVariant(chcol.Dynamic{}))
assert.Equal(t, uint64(3), unwrapVariant(uint64(3)))
j := chcol.NewJSON()
j.SetValueAtPath("level", "error")
j.SetValueAtPath("attrs.code", int64(500))
assert.Equal(t, map[string]any{
"level": "error",
"attrs": map[string]any{"code": float64(500)},
}, unwrapVariant(*j))
}
// labelValue renders a JSON group-by value as a stable, sorted-key string so structurally equal
// documents share a series.
func TestLabelValue(t *testing.T) {
assert.Equal(t, "", labelValue(nil))
assert.Equal(t, "error", labelValue(chcol.NewDynamicWithType("error", "String")))
assert.Equal(t, `{"attrs":{"code":500},"level":"error"}`, labelValue(map[string]any{
"level": "error",
"attrs": map[string]any{"code": 500},
}))
}
func TestMergeSpanAttributeColumns_EmptyEventsAndLinks(t *testing.T) {
@@ -207,7 +132,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
{
name: "JSONOnly_FlattensNestedPaths_PreservesTypes",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{
"attributes": map[string]any{
"http": map[string]any{"route": "/api/pay", "retry": map[string]any{"count": float64(3)}},
"cache.hit": true,
},
@@ -219,7 +144,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
data: map[string]any{
"attributes_string": map[string]string{"http.route": "/old", "only.map": "m"},
"attributes_number": map[string]float64{"http.status": 500},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"route": "/new"}, "only.json": "j"},
"attributes": map[string]any{"http": map[string]any{"route": "/new"}, "only.json": "j"},
},
want: map[string]any{"http.route": "/old", "only.map": "m", "http.status": float64(500), "only.json": "j"},
},
@@ -229,7 +154,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
"attributes_string": map[string]string{"http.route": "/map"},
"attributes_number": map[string]float64{"http.status": 200},
"attributes_bool": map[string]bool{"cache.hit": true},
"attributes": telemetrystoretypes.JSONValue{},
"attributes": map[string]any{},
},
want: map[string]any{"http.route": "/map", "http.status": float64(200), "cache.hit": true},
},
@@ -237,28 +162,28 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
name: "MapOnly_NilJSON_BehavesAsAbsent",
data: map[string]any{
"attributes_string": map[string]string{"http.route": "/map"},
"attributes": telemetrystoretypes.JSONValue(nil),
"attributes": map[string]any(nil),
},
want: map[string]any{"http.route": "/map"},
},
{
name: "Arrays_StayLeafValues",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"tags": []any{"a", "b"}, "codes": []any{float64(1), float64(2)}}},
"attributes": map[string]any{"http": map[string]any{"tags": []any{"a", "b"}, "codes": []any{float64(1), float64(2)}}},
},
want: map[string]any{"http.tags": []any{"a", "b"}, "http.codes": []any{float64(1), float64(2)}},
},
{
name: "TopLevelArrayOfMaps_StaysNativeLeaf",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
"attributes": map[string]any{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
},
want: map[string]any{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
},
{
name: "NestedArrayOfMaps_StaysNativeLeaf_NoIndexPaths",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
"attributes": map[string]any{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
},
want: map[string]any{"http.items": []any{map[string]any{"a": float64(1)}}},
},
@@ -266,28 +191,28 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
name: "DualWritten_NestedArray_IndexKeysAndJSONArrayCoexist",
data: map[string]any{
"attributes_number": map[string]float64{"http.items.0.a": 1},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
"attributes": map[string]any{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
},
want: map[string]any{"http.items.0.a": float64(1), "http.items": []any{map[string]any{"a": float64(1)}}},
},
{
name: "JSONNull_KeptAsNil",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{"k": nil},
"attributes": map[string]any{"k": nil},
},
want: map[string]any{"k": nil},
},
{
name: "KeyIsLeafValue_NotFlattened",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{"http": "plaintext"},
"attributes": map[string]any{"http": "plaintext"},
},
want: map[string]any{"http": "plaintext"},
},
{
name: "KeyIsParent_FlattensToDottedPath",
data: map[string]any{
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"route": "/a"}},
"attributes": map[string]any{"http": map[string]any{"route": "/a"}},
},
want: map[string]any{"http.route": "/a"},
},

View File

@@ -16,7 +16,6 @@ import (
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
@@ -1202,7 +1201,7 @@ func (q *querier) postProcessLogBody(ctx context.Context, orgID valuer.UUID, res
// carried one. Anything that is not a decoded document — the legacy string body, a NULL cell —
// is legal under these names and left alone.
func stripEmptyBodyMessage(val any) {
bodyMap, ok := val.(telemetrystoretypes.JSONValue)
bodyMap, ok := val.(map[string]any)
if !ok {
return
}

View File

@@ -13,7 +13,6 @@ import (
"github.com/SigNoz/signoz/pkg/sqlschema"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/types/ruletypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
@@ -50,6 +49,11 @@ type rule struct {
OrgID string `bun:"org_id,type:text"`
}
type routePolicyRuleData struct {
PreferredChannels []string `json:"preferredChannels"`
Labels map[string]string `json:"labels"`
}
type addRoutePolicies struct {
sqlstore sqlstore.SQLStore
sqlschema sqlschema.SQLSchema
@@ -187,20 +191,20 @@ func (migration *addRoutePolicies) migrateRulesToRoutePolicies(ctx context.Conte
func (migration *addRoutePolicies) convertRulesToRoutes(rules []*rule, channelsByOrg map[string][]string) ([]*expressionRoute, error) {
var routes []*expressionRoute
for _, r := range rules {
var gettableRule ruletypes.GettableRule
if err := json.Unmarshal([]byte(r.Data), &gettableRule); err != nil {
var ruleData routePolicyRuleData
if err := json.Unmarshal([]byte(r.Data), &ruleData); err != nil {
return nil, errors.NewInternalf(errors.CodeInternal, "failed to unmarshal rule data for rule ID %s: %v", r.ID, err)
}
if len(gettableRule.PreferredChannels) == 0 {
if len(ruleData.PreferredChannels) == 0 {
channels, exists := channelsByOrg[r.OrgID]
if !exists || len(channels) == 0 {
continue
}
gettableRule.PreferredChannels = channels
ruleData.PreferredChannels = channels
}
severity := "critical"
if v, ok := gettableRule.Labels["severity"]; ok {
if v, ok := ruleData.Labels["severity"]; ok {
severity = v
}
expression := fmt.Sprintf(`%s == "%s" && %s == "%s"`, "threshold.name", severity, "ruleId", r.ID.String())
@@ -218,7 +222,7 @@ func (migration *addRoutePolicies) convertRulesToRoutes(rules []*rule, channelsB
},
Expression: expression,
ExpressionKind: "rule",
Channels: gettableRule.PreferredChannels,
Channels: ruleData.PreferredChannels,
Name: r.ID.StringValue(),
Enabled: true,
OrgID: r.OrgID,

View File

@@ -97,8 +97,8 @@ func New(ctx context.Context, providerSettings factory.ProviderSettings, config
options.MaxIdleConns = config.Connection.MaxIdleConns
options.MaxOpenConns = config.Connection.MaxOpenConns
options.DialTimeout = config.Connection.DialTimeout
// This is to avoid the driver decoding issues with JSON columns
options.Settings["output_format_native_write_json_as_string"] = 1
// Decode JSON columns via the flattened native serialization (CH 25.6+); without it clickhouse-go mis-decodes the SharedData layout of JSON(max_dynamic_paths=0) columns and desyncs the native protocol.
options.Settings["output_format_native_use_flattened_dynamic_and_json_serialization"] = 1
chConn, err := clickhouse.Open(options)
if err != nil {
@@ -184,7 +184,7 @@ func (p *provider) Query(ctx context.Context, query string, args ...interface{})
}
return &rowsWithHooks{
Rows: telemetrystore.WrapRows(rows),
Rows: rows,
ctx: ctx,
event: event,
onClose: func() { telemetrystore.WrapAfterQuery(p.hooks, ctx, event) },

View File

@@ -1,39 +0,0 @@
package telemetrystore
import (
"reflect"
"strings"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
)
// WrapRows reports JSONValue as the scan type of every JSON column. Nested JSON — Array(JSON),
// Map(String, JSON) — is not covered.
func WrapRows(rows driver.Rows) driver.Rows {
return &rowsWithJSONScanType{Rows: rows}
}
type rowsWithJSONScanType struct {
driver.Rows
}
func (r *rowsWithJSONScanType) ColumnTypes() []driver.ColumnType {
colTypes := r.Rows.ColumnTypes()
wrapped := make([]driver.ColumnType, len(colTypes))
for i, colType := range colTypes {
wrapped[i] = colType
if strings.HasPrefix(strings.ToUpper(colType.DatabaseTypeName()), "JSON") {
wrapped[i] = jsonColumnType{ColumnType: colType}
}
}
return wrapped
}
type jsonColumnType struct {
driver.ColumnType
}
func (jsonColumnType) ScanType() reflect.Type {
return reflect.TypeFor[telemetrystoretypes.JSONValue]()
}

View File

@@ -1,23 +0,0 @@
package telemetrystoretest
import (
"context"
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/SigNoz/signoz/pkg/telemetrystore"
)
// conn wraps rows the way the clickhouse provider does, so mocked JSON columns report the scan
// type they do in production.
type conn struct {
clickhouse.Conn
}
func (c conn) Query(ctx context.Context, query string, args ...any) (driver.Rows, error) {
rows, err := c.Conn.Query(ctx, query, args...)
if err != nil {
return nil, err
}
return telemetrystore.WrapRows(rows), nil
}

View File

@@ -32,7 +32,7 @@ func New(_ telemetrystore.Config, matcher sqlmock.QueryMatcher) *Provider {
// ClickhouseDB returns the mock Clickhouse connection.
func (p *Provider) ClickhouseDB() clickhouse.Conn {
return conn{Conn: p.clickhouseDB.(clickhouse.Conn)}
return p.clickhouseDB.(clickhouse.Conn)
}
// Cluster returns the cluster name.

View File

@@ -364,12 +364,26 @@ func (provider *provider) gc(ctx context.Context, org *types.Organization) error
}
func (provider *provider) flushLastObservedAt(ctx context.Context, org *types.Organization) error {
accessTokenToLastObservedAt, err := provider.listLastObservedAtDesc(ctx, org.ID)
tokens, err := provider.tokenStore.ListByOrgID(ctx, org.ID)
if err != nil {
return err
}
if err := provider.tokenStore.UpdateLastObservedAtByAccessToken(ctx, accessTokenToLastObservedAt); err != nil {
observedTokens := make([]*authtypes.StorableToken, 0, len(tokens))
for _, token := range tokens {
cachedLastObservedAt, ok := provider.lastObservedAtCache.Get(lastObservedAtCacheKey(token.AccessToken, token.UserID))
if !ok {
continue
}
if err := token.UpdateLastObservedAt(cachedLastObservedAt); err != nil {
continue
}
observedTokens = append(observedTokens, token)
}
if err := provider.tokenStore.UpdateLastObservedAt(ctx, observedTokens); err != nil {
return err
}

View File

@@ -232,15 +232,16 @@ func (store *store) ListByUserID(ctx context.Context, userID valuer.UUID) ([]*au
return tokens, nil
}
func (store *store) UpdateLastObservedAtByAccessToken(ctx context.Context, accessTokenToLastObservedAt []map[string]any) error {
if len(accessTokenToLastObservedAt) == 0 {
func (store *store) UpdateLastObservedAt(ctx context.Context, tokens []*authtypes.StorableToken) error {
if len(tokens) == 0 {
return nil
}
values := store.
sqlstore.
BunDBCtx(ctx).
NewValues(&accessTokenToLastObservedAt)
NewValues(&tokens).
Column("id", "last_observed_at", "updated_at")
_, err := store.
sqlstore.
@@ -250,8 +251,8 @@ func (store *store) UpdateLastObservedAtByAccessToken(ctx context.Context, acces
Model((*authtypes.StorableToken)(nil)).
TableExpr("update_cte").
Set("last_observed_at = update_cte.last_observed_at").
Where("auth_token.access_token = update_cte.access_token").
Where("auth_token.user_id = update_cte.user_id").
Set("updated_at = update_cte.updated_at").
Where("auth_token.id = update_cte.id").
Exec(ctx)
if err != nil {
return err

View File

@@ -0,0 +1,62 @@
package alertmanagertypes
import (
"maps"
"net/textproto"
"slices"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
)
// ChannelEmailConfig carries no SMTP transport fields: the smarthost,
// credentials and TLS settings come from the deployment's global config, so a
// channel can only choose recipients and body.
type ChannelEmailConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
To string `json:"to" required:"true"`
HTML valuer.UnsetOrNonEmptyString `json:"html"`
Headers map[string]string `json:"headers,omitempty"`
}
func (c ChannelEmailConfig) Validate() error {
if c.To == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.to is required for an email channel")
}
// A read reports header names as textproto canonicalizes them, turning
// "subject" into "Subject", so a name that is not already in that form is
// rejected rather than answered with one the caller never sent.
for _, header := range slices.Sorted(maps.Keys(c.Headers)) {
if canonical := textproto.CanonicalMIMEHeaderKey(header); canonical != header {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.headers name %q must be written as %q", header, canonical)
}
}
return nil
}
func (c ChannelEmailConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
EmailConfigs: []*config.EmailConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultEmailConfig.VSendResolved)},
To: c.To,
HTML: c.HTML.StringValue(),
Headers: c.Headers,
}},
}}, nil
}
func newChannelEmailConfigFromReceiver(_ string, receiver *Receiver) (ChannelSpec, error) {
email := receiver.EmailConfigs[0]
sendResolved := email.VSendResolved
return &ChannelEmailConfig{
SendResolved: &sendResolved,
To: email.To,
HTML: valuer.UnsetIfEmpty(email.HTML),
Headers: email.Headers,
}, nil
}

View File

@@ -5,10 +5,59 @@ import (
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
type ChannelGoogleChatConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
}
func (c ChannelGoogleChatConfig) Validate() error {
if c.WebhookURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for a googlechat channel")
}
return nil
}
func (c ChannelGoogleChatConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
webhookURL, err := parseSecretURL(c.WebhookURL)
if err != nil {
return nil, err
}
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
GoogleChatConfigs: []*GoogleChatReceiverConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultGoogleChatReceiverConfig.VSendResolved)},
WebhookURL: webhookURL,
Title: c.Title.StringValue(),
Text: c.Text.StringValue(),
}},
}, nil
}
func newChannelGoogleChatConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
googlechat := receiver.GoogleChatConfigs[0]
sendResolved := googlechat.VSendResolved
if err := rejectAnyHTTPAuth(name, googlechat.HTTPConfig); err != nil {
return nil, err
}
return &ChannelGoogleChatConfig{
SendResolved: &sendResolved,
WebhookURL: formatSecretURL(googlechat.WebhookURL),
Title: valuer.UnsetIfEmpty(googlechat.Title),
Text: valuer.UnsetIfEmpty(googlechat.Text),
}, nil
}
type GoogleChatReceiverConfig struct {
config.NotifierConfig `yaml:",inline" json:",inline"`

View File

@@ -6,10 +6,64 @@ import (
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
type ChannelIncidentIOConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
URL string `json:"url" required:"true"`
Token string `json:"token" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Metadata map[string]string `json:"metadata,omitempty"`
}
func (c ChannelIncidentIOConfig) Validate() error {
if c.URL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for an incidentio channel")
}
if c.Token == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.token is required for an incidentio channel")
}
return nil
}
func (c ChannelIncidentIOConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
IncidentIOConfigs: []*IncidentIOReceiverConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultIncidentIOReceiverConfig.VSendResolved)},
URL: c.URL,
Token: config.Secret(c.Token),
Title: c.Title.StringValue(),
Description: c.Description.StringValue(),
Metadata: c.Metadata,
}},
}, nil
}
func newChannelIncidentIOConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
incidentio := receiver.IncidentIOConfigs[0]
sendResolved := incidentio.VSendResolved
if err := rejectAnyHTTPAuth(name, incidentio.HTTPConfig); err != nil {
return nil, err
}
return &ChannelIncidentIOConfig{
SendResolved: &sendResolved,
URL: incidentio.URL,
Token: string(incidentio.Token),
Title: valuer.UnsetIfEmpty(incidentio.Title),
Description: valuer.UnsetIfEmpty(incidentio.Description),
Metadata: incidentio.Metadata,
}, nil
}
// incidentIOEventsPathPrefix is the path of incident.io's HTTP alert source
// endpoint (Alert Events V2 API). The full URL is per-source:
// https://api.incident.io/v2/alert_events/http/<source_config_id>.

View File

@@ -7,11 +7,147 @@ import (
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
"github.com/prometheus/common/model"
)
type ChannelJiraConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
// Site is the Jira Cloud base URL, https://<site>.atlassian.net. Only Jira
// Cloud is supported; the REST base is derived from it.
Site string `json:"site" required:"true"`
Project string `json:"project" required:"true"`
IssueType string `json:"issueType" required:"true"`
Summary valuer.UnsetOrNonEmptyString `json:"summary"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Priority string `json:"priority"`
Labels []string `json:"labels,omitempty"`
ResolveTransition string `json:"resolveTransition"`
ReopenTransition string `json:"reopenTransition"`
ReopenDuration valuer.UnsetOrNonEmptyString `json:"reopenDuration"`
WontFixResolution string `json:"wontFixResolution"`
CustomFields map[string]any `json:"customFields,omitempty"`
Email string `json:"email" required:"true"`
APIToken string `json:"apiToken" required:"true" format:"password"`
}
func (c ChannelJiraConfig) Validate() error {
for _, required := range []struct {
value string
field string
}{
{c.Site, "site"},
{c.Project, "project"},
{c.IssueType, "issueType"},
{c.Email, "email"},
{c.APIToken, "apiToken"},
} {
if required.value == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.%s is required for a jira channel", required.field)
}
}
if !c.ReopenDuration.IsZero() {
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
if err != nil {
return errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q is not a valid duration", c.ReopenDuration)
}
// A read reports the duration as model.Duration formats it, collapsing
// "72h" into "3d", so a value that is not already in that form is rejected
// rather than answered with one the caller never sent.
if canonical := reopenDuration.String(); canonical != c.ReopenDuration.StringValue() {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q must be written as %q", c.ReopenDuration, canonical)
}
}
return nil
}
func (c ChannelJiraConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
// Seeded from upstream's default rather than a zero value: FollowRedirects
// and EnableHTTP2 marshal unconditionally, so a zero value would persist them
// as false and read back as a config ChannelJiraConfig cannot represent.
httpConfig := commoncfg.DefaultHTTPClientConfig
httpConfig.BasicAuth = &commoncfg.BasicAuth{
Username: c.Email,
Password: commoncfg.Secret(c.APIToken),
}
jira := &JiraReceiverConfig{
// JiraReceiverConfig seeds no send_resolved of its own, so unset means off.
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, false)},
Site: c.Site,
Project: c.Project,
IssueType: c.IssueType,
Summary: c.Summary.StringValue(),
Description: c.Description.StringValue(),
Priority: c.Priority,
Labels: c.Labels,
ResolveTransition: c.ResolveTransition,
ReopenTransition: c.ReopenTransition,
WontFixResolution: c.WontFixResolution,
CustomFields: c.CustomFields,
HTTPConfig: &httpConfig,
}
if !c.ReopenDuration.IsZero() {
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "parse reopenDuration %q", c.ReopenDuration)
}
jira.ReopenDuration = reopenDuration
}
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
JiraConfigs: []*JiraReceiverConfig{jira},
}, nil
}
func newChannelJiraConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
jira := receiver.JiraConfigs[0]
sendResolved := jira.VSendResolved
if err := rejectUnsupportedHTTPConfig(name, jira.HTTPConfig); err != nil {
return nil, err
}
if jira.HTTPConfig != nil && jira.HTTPConfig.Authorization != nil {
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization, which is not supported", name)
}
if err := rejectHTTPBasicAuthBeyondPassword(name, jira.HTTPConfig); err != nil {
return nil, err
}
spec := &ChannelJiraConfig{
SendResolved: &sendResolved,
Site: jira.Site,
Project: jira.Project,
IssueType: jira.IssueType,
Summary: valuer.UnsetIfEmpty(jira.Summary),
Description: valuer.UnsetIfEmpty(jira.Description),
Priority: jira.Priority,
Labels: jira.Labels,
ResolveTransition: jira.ResolveTransition,
ReopenTransition: jira.ReopenTransition,
ReopenDuration: valuer.UnsetIfEmpty(jira.ReopenDuration.String()),
WontFixResolution: jira.WontFixResolution,
CustomFields: jira.CustomFields,
}
if jira.HTTPConfig != nil && jira.HTTPConfig.BasicAuth != nil {
spec.Email = jira.HTTPConfig.BasicAuth.Username
spec.APIToken = string(jira.HTTPConfig.BasicAuth.Password)
}
return spec, nil
}
const defaultJiraReopenDuration = model.Duration(3 * 24 * time.Hour)
// Service accounts authenticate against the api.atlassian.com gateway (keyed by

View File

@@ -2,10 +2,63 @@ package alertmanagertypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
// ChannelJSMOpsConfig carries no API URL: JSM Ops is a single global gateway
// keyed by the integration API key, which the notifier pins itself.
type ChannelJSMOpsConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIKey string `json:"apiKey" required:"true" format:"password"`
Message valuer.UnsetOrNonEmptyString `json:"message"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Priority string `json:"priority"`
// Tags is the comma-separated list JSM Ops attaches to the alert.
Tags valuer.UnsetOrNonEmptyString `json:"tags"`
}
func (c ChannelJSMOpsConfig) Validate() error {
if c.APIKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for a jsmops channel")
}
return nil
}
func (c ChannelJSMOpsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
JSMOpsConfigs: []*JSMOpsReceiverConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultJSMOpsReceiverConfig.VSendResolved)},
APIKey: config.Secret(c.APIKey),
Message: c.Message.StringValue(),
Description: c.Description.StringValue(),
Priority: c.Priority,
Tags: c.Tags.StringValue(),
}},
}, nil
}
func newChannelJSMOpsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
jsmops := receiver.JSMOpsConfigs[0]
sendResolved := jsmops.VSendResolved
if err := rejectAnyHTTPAuth(name, jsmops.HTTPConfig); err != nil {
return nil, err
}
return &ChannelJSMOpsConfig{
SendResolved: &sendResolved,
APIKey: string(jsmops.APIKey),
Message: valuer.UnsetIfEmpty(jsmops.Message),
Description: valuer.UnsetIfEmpty(jsmops.Description),
Priority: jsmops.Priority,
Tags: valuer.UnsetIfEmpty(jsmops.Tags),
}, nil
}
// JSMOpsAPIBaseURL is the native JSM Ops integration-events gateway. It is a
// single global host keyed by the integration API key (no region/cloud id in
// the path). The trailing slash is required: the Opsgenie notifier appends

View File

@@ -0,0 +1,55 @@
package alertmanagertypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
)
type ChannelMSTeamsConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
}
func (c ChannelMSTeamsConfig) Validate() error {
if c.WebhookURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for an msteams channel")
}
return nil
}
func (c ChannelMSTeamsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
webhookURL, err := parseSecretURL(c.WebhookURL)
if err != nil {
return nil, err
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
MSTeamsV2Configs: []*config.MSTeamsV2Config{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultMSTeamsV2Config.VSendResolved)},
WebhookURL: webhookURL,
Title: c.Title.StringValue(),
Text: c.Text.StringValue(),
}},
}}, nil
}
func newChannelMSTeamsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
msteams := receiver.MSTeamsV2Configs[0]
sendResolved := msteams.VSendResolved
if err := rejectAnyHTTPAuth(name, msteams.HTTPConfig); err != nil {
return nil, err
}
return &ChannelMSTeamsConfig{
SendResolved: &sendResolved,
WebhookURL: formatSecretURL(msteams.WebhookURL),
Title: valuer.UnsetIfEmpty(msteams.Title),
Text: valuer.UnsetIfEmpty(msteams.Text),
}, nil
}

View File

@@ -0,0 +1,71 @@
package alertmanagertypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
)
type ChannelOpsgenieConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIKey string `json:"apiKey" required:"true" format:"password"`
APIURL string `json:"apiUrl"`
Message valuer.UnsetOrNonEmptyString `json:"message"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Source valuer.UnsetOrNonEmptyString `json:"source"`
Details map[string]string `json:"details,omitempty"`
Priority string `json:"priority"`
}
func (c ChannelOpsgenieConfig) Validate() error {
if c.APIKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for an opsgenie channel")
}
return nil
}
func (c ChannelOpsgenieConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
var apiURL *config.URL
if c.APIURL != "" {
parsed, err := parseUpstreamURL(c.APIURL)
if err != nil {
return nil, err
}
apiURL = parsed
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
OpsGenieConfigs: []*config.OpsGenieConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultOpsGenieConfig.VSendResolved)},
APIKey: config.Secret(c.APIKey),
APIURL: apiURL,
Message: c.Message.StringValue(),
Description: c.Description.StringValue(),
Source: c.Source.StringValue(),
Priority: c.Priority,
Details: c.Details,
}},
}}, nil
}
func newChannelOpsgenieConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
opsgenie := receiver.OpsGenieConfigs[0]
sendResolved := opsgenie.VSendResolved
if err := rejectAnyHTTPAuth(name, opsgenie.HTTPConfig); err != nil {
return nil, err
}
return &ChannelOpsgenieConfig{
SendResolved: &sendResolved,
APIKey: string(opsgenie.APIKey),
APIURL: formatUpstreamURL(opsgenie.APIURL),
Message: valuer.UnsetIfEmpty(opsgenie.Message),
Description: valuer.UnsetIfEmpty(opsgenie.Description),
Source: valuer.UnsetIfEmpty(opsgenie.Source),
Priority: opsgenie.Priority,
Details: opsgenie.Details,
}, nil
}

View File

@@ -0,0 +1,92 @@
package alertmanagertypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
)
type ChannelPagerdutyConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
RoutingKey string `json:"routingKey" required:"true" format:"password"`
URL string `json:"url"`
Source valuer.UnsetOrNonEmptyString `json:"source"`
Client valuer.UnsetOrNonEmptyString `json:"client"`
ClientURL valuer.UnsetOrNonEmptyString `json:"clientUrl"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Severity string `json:"severity"`
Component string `json:"component"`
Group string `json:"group"`
Class string `json:"class"`
Details map[string]string `json:"details,omitempty"`
}
func (c ChannelPagerdutyConfig) Validate() error {
if c.RoutingKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.routingKey is required for a pagerduty channel")
}
return nil
}
func (c ChannelPagerdutyConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
var eventsURL *config.URL
if c.URL != "" {
parsed, err := parseUpstreamURL(c.URL)
if err != nil {
return nil, err
}
eventsURL = parsed
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
PagerdutyConfigs: []*config.PagerdutyConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultPagerdutyConfig.VSendResolved)},
RoutingKey: config.Secret(c.RoutingKey),
URL: eventsURL,
Source: c.Source.StringValue(),
Client: c.Client.StringValue(),
ClientURL: c.ClientURL.StringValue(),
Description: c.Description.StringValue(),
Severity: c.Severity,
Component: c.Component,
Group: c.Group,
Class: c.Class,
Details: newUpstreamDetails(c.Details),
}},
}}, nil
}
func newChannelPagerdutyConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
pagerduty := receiver.PagerdutyConfigs[0]
sendResolved := pagerduty.VSendResolved
if err := rejectAnyHTTPAuth(name, pagerduty.HTTPConfig); err != nil {
return nil, err
}
var details map[string]string
if len(pagerduty.Details) > 0 {
extracted, err := extractStringDetails(name, pagerduty.Details)
if err != nil {
return nil, err
}
details = extracted
}
return &ChannelPagerdutyConfig{
SendResolved: &sendResolved,
RoutingKey: string(pagerduty.RoutingKey),
URL: formatUpstreamURL(pagerduty.URL),
Source: valuer.UnsetIfEmpty(pagerduty.Source),
Client: valuer.UnsetIfEmpty(pagerduty.Client),
ClientURL: valuer.UnsetIfEmpty(pagerduty.ClientURL),
Description: valuer.UnsetIfEmpty(pagerduty.Description),
Severity: pagerduty.Severity,
Component: pagerduty.Component,
Group: pagerduty.Group,
Class: pagerduty.Class,
Details: details,
}, nil
}

View File

@@ -0,0 +1,182 @@
package alertmanagertypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
)
type ChannelSlackConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIURL string `json:"apiUrl" required:"true" format:"password"`
Channel string `json:"channel"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
Color valuer.UnsetOrNonEmptyString `json:"color"`
TitleLink valuer.UnsetOrNonEmptyString `json:"titleLink"`
Pretext valuer.UnsetOrNonEmptyString `json:"pretext"`
Fallback valuer.UnsetOrNonEmptyString `json:"fallback"`
Footer valuer.UnsetOrNonEmptyString `json:"footer"`
Fields []ChannelSlackField `json:"fields,omitempty"`
Actions []ChannelSlackAction `json:"actions,omitempty"`
}
type ChannelSlackField struct {
Title string `json:"title" required:"true"`
Value string `json:"value" required:"true"`
Short *bool `json:"short,omitempty"`
}
// ChannelSlackAction is a link button when URL is set, otherwise a message
// button that needs Name. Upstream clears whichever side is not in use.
type ChannelSlackAction struct {
Type string `json:"type" required:"true"`
Text string `json:"text" required:"true"`
URL string `json:"url"`
Style string `json:"style"`
Name string `json:"name"`
Value string `json:"value"`
Confirm *ChannelSlackConfirmation `json:"confirm,omitempty"`
}
type ChannelSlackConfirmation struct {
Text string `json:"text" required:"true"`
Title string `json:"title"`
OkText string `json:"okText"`
DismissText string `json:"dismissText"`
}
func (c ChannelSlackConfig) Validate() error {
if c.APIURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiUrl is required for a slack channel")
}
for i, field := range c.Fields {
if field.Title == "" || field.Value == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.fields[%d] requires title and value", i)
}
}
for i, action := range c.Actions {
if action.Type == "" || action.Text == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires type and text", i)
}
if action.URL == "" && action.Name == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires url or name", i)
}
if action.Confirm != nil && action.Confirm.Text == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d].confirm requires text", i)
}
}
return nil
}
func (c ChannelSlackConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
apiURL, err := parseSecretURL(c.APIURL)
if err != nil {
return nil, err
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
SlackConfigs: []*config.SlackConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultSlackConfig.VSendResolved)},
APIURL: apiURL,
Channel: c.Channel,
Title: c.Title.StringValue(),
Text: c.Text.StringValue(),
Color: c.Color.StringValue(),
TitleLink: c.TitleLink.StringValue(),
Pretext: c.Pretext.StringValue(),
Fallback: c.Fallback.StringValue(),
Footer: c.Footer.StringValue(),
Fields: newUpstreamSlackFields(c.Fields),
Actions: newUpstreamSlackActions(c.Actions),
}},
}}, nil
}
func newChannelSlackConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
slack := receiver.SlackConfigs[0]
sendResolved := slack.VSendResolved
if err := rejectAnyHTTPAuth(name, slack.HTTPConfig); err != nil {
return nil, err
}
return &ChannelSlackConfig{
SendResolved: &sendResolved,
APIURL: formatSecretURL(slack.APIURL),
Channel: slack.Channel,
Title: valuer.UnsetIfEmpty(slack.Title),
Text: valuer.UnsetIfEmpty(slack.Text),
Color: valuer.UnsetIfEmpty(slack.Color),
TitleLink: valuer.UnsetIfEmpty(slack.TitleLink),
Pretext: valuer.UnsetIfEmpty(slack.Pretext),
Fallback: valuer.UnsetIfEmpty(slack.Fallback),
Footer: valuer.UnsetIfEmpty(slack.Footer),
Fields: newChannelSlackFields(slack.Fields),
Actions: newChannelSlackActions(slack.Actions),
}, nil
}
func newUpstreamSlackFields(fields []ChannelSlackField) []*config.SlackField {
if len(fields) == 0 {
return nil
}
upstream := make([]*config.SlackField, 0, len(fields))
for _, field := range fields {
upstream = append(upstream, &config.SlackField{Title: field.Title, Value: field.Value, Short: field.Short})
}
return upstream
}
func newChannelSlackFields(upstream []*config.SlackField) []ChannelSlackField {
if len(upstream) == 0 {
return nil
}
fields := make([]ChannelSlackField, 0, len(upstream))
for _, field := range upstream {
fields = append(fields, ChannelSlackField{Title: field.Title, Value: field.Value, Short: field.Short})
}
return fields
}
func newUpstreamSlackActions(actions []ChannelSlackAction) []*config.SlackAction {
if len(actions) == 0 {
return nil
}
upstream := make([]*config.SlackAction, 0, len(actions))
for _, action := range actions {
upstreamAction := &config.SlackAction{Type: action.Type, Text: action.Text, URL: action.URL, Style: action.Style, Name: action.Name, Value: action.Value}
if action.Confirm != nil {
upstreamAction.ConfirmField = &config.SlackConfirmationField{Text: action.Confirm.Text, Title: action.Confirm.Title, OkText: action.Confirm.OkText, DismissText: action.Confirm.DismissText}
}
upstream = append(upstream, upstreamAction)
}
return upstream
}
func newChannelSlackActions(upstream []*config.SlackAction) []ChannelSlackAction {
if len(upstream) == 0 {
return nil
}
actions := make([]ChannelSlackAction, 0, len(upstream))
for _, upstreamAction := range upstream {
action := ChannelSlackAction{Type: upstreamAction.Type, Text: upstreamAction.Text, URL: upstreamAction.URL, Style: upstreamAction.Style, Name: upstreamAction.Name, Value: upstreamAction.Value}
if upstreamAction.ConfirmField != nil {
action.Confirm = &ChannelSlackConfirmation{Text: upstreamAction.ConfirmField.Text, Title: upstreamAction.ConfirmField.Title, OkText: upstreamAction.ConfirmField.OkText, DismissText: upstreamAction.ConfirmField.DismissText}
}
actions = append(actions, action)
}
return actions
}

View File

@@ -0,0 +1,98 @@
package alertmanagertypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
)
// ChannelWebhookConfig splits apart the two authentication modes the legacy API
// overloaded onto one password field, where an empty username meant the password
// was really a bearer token. Username or Password may be set without the other,
// as upstream allows, but not together with BearerToken.
type ChannelWebhookConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
URL string `json:"url" required:"true" format:"password"`
Username string `json:"username"`
Password string `json:"password" format:"password"`
BearerToken string `json:"bearerToken" format:"password"`
}
func (c ChannelWebhookConfig) Validate() error {
if c.URL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for a webhook channel")
}
usesBasicAuth := c.Username != "" || c.Password != ""
if usesBasicAuth && c.BearerToken != "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.bearerToken cannot be combined with config.spec.username or config.spec.password")
}
return nil
}
func (c ChannelWebhookConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
webhook := &config.WebhookConfig{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultWebhookConfig.VSendResolved)},
URL: config.SecretTemplateURL(c.URL),
}
// Seeded from upstream's default rather than a zero value: FollowRedirects
// and EnableHTTP2 marshal unconditionally, so a zero value would persist
// them as false and read back as a config ChannelWebhookConfig cannot represent.
switch {
case c.Username != "" || c.Password != "":
httpConfig := commoncfg.DefaultHTTPClientConfig
httpConfig.BasicAuth = &commoncfg.BasicAuth{
Username: c.Username,
Password: commoncfg.Secret(c.Password),
}
webhook.HTTPConfig = &httpConfig
case c.BearerToken != "":
httpConfig := commoncfg.DefaultHTTPClientConfig
httpConfig.Authorization = &commoncfg.Authorization{
Type: bearerAuthorizationType,
Credentials: commoncfg.Secret(c.BearerToken),
}
webhook.HTTPConfig = &httpConfig
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
WebhookConfigs: []*config.WebhookConfig{webhook},
}}, nil
}
func newChannelWebhookConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
upstream := receiver.WebhookConfigs[0]
sendResolved := upstream.VSendResolved
if err := rejectUnsupportedHTTPConfig(name, upstream.HTTPConfig); err != nil {
return nil, err
}
if err := rejectHTTPBasicAuthBeyondPassword(name, upstream.HTTPConfig); err != nil {
return nil, err
}
if err := rejectHTTPAuthorizationBeyondBearer(name, upstream.HTTPConfig); err != nil {
return nil, err
}
webhook := &ChannelWebhookConfig{
SendResolved: &sendResolved,
URL: string(upstream.URL),
}
if upstream.HTTPConfig != nil {
if basicAuth := upstream.HTTPConfig.BasicAuth; basicAuth != nil {
webhook.Username = basicAuth.Username
webhook.Password = string(basicAuth.Password)
}
if authorization := upstream.HTTPConfig.Authorization; authorization != nil {
webhook.BearerToken = string(authorization.Credentials)
}
}
return webhook, nil
}

View File

@@ -3,8 +3,6 @@ package alertmanagertypes
import (
"bytes"
"encoding/json"
"maps"
"net/textproto"
"net/url"
"reflect"
"slices"
@@ -14,7 +12,6 @@ import (
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/prometheus/alertmanager/config"
commoncfg "github.com/prometheus/common/config"
"github.com/prometheus/common/model"
"github.com/swaggest/jsonschema-go"
)
@@ -205,808 +202,6 @@ type ChannelSpec interface {
toUndefaultedReceiver(displayName string) (*Receiver, error)
}
type ChannelSlackConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIURL string `json:"apiUrl" required:"true" format:"password"`
Channel string `json:"channel"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
Color valuer.UnsetOrNonEmptyString `json:"color"`
TitleLink valuer.UnsetOrNonEmptyString `json:"titleLink"`
Pretext valuer.UnsetOrNonEmptyString `json:"pretext"`
Fallback valuer.UnsetOrNonEmptyString `json:"fallback"`
Footer valuer.UnsetOrNonEmptyString `json:"footer"`
Fields []ChannelSlackField `json:"fields,omitempty"`
Actions []ChannelSlackAction `json:"actions,omitempty"`
}
type ChannelSlackField struct {
Title string `json:"title" required:"true"`
Value string `json:"value" required:"true"`
Short *bool `json:"short,omitempty"`
}
// ChannelSlackAction is a link button when URL is set, otherwise a message
// button that needs Name. Upstream clears whichever side is not in use.
type ChannelSlackAction struct {
Type string `json:"type" required:"true"`
Text string `json:"text" required:"true"`
URL string `json:"url"`
Style string `json:"style"`
Name string `json:"name"`
Value string `json:"value"`
Confirm *ChannelSlackConfirmation `json:"confirm,omitempty"`
}
type ChannelSlackConfirmation struct {
Text string `json:"text" required:"true"`
Title string `json:"title"`
OkText string `json:"okText"`
DismissText string `json:"dismissText"`
}
func (c ChannelSlackConfig) Validate() error {
if c.APIURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiUrl is required for a slack channel")
}
for i, field := range c.Fields {
if field.Title == "" || field.Value == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.fields[%d] requires title and value", i)
}
}
for i, action := range c.Actions {
if action.Type == "" || action.Text == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires type and text", i)
}
if action.URL == "" && action.Name == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires url or name", i)
}
if action.Confirm != nil && action.Confirm.Text == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d].confirm requires text", i)
}
}
return nil
}
func (c ChannelSlackConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
apiURL, err := parseSecretURL(c.APIURL)
if err != nil {
return nil, err
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
SlackConfigs: []*config.SlackConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultSlackConfig.VSendResolved)},
APIURL: apiURL,
Channel: c.Channel,
Title: c.Title.StringValue(),
Text: c.Text.StringValue(),
Color: c.Color.StringValue(),
TitleLink: c.TitleLink.StringValue(),
Pretext: c.Pretext.StringValue(),
Fallback: c.Fallback.StringValue(),
Footer: c.Footer.StringValue(),
Fields: newUpstreamSlackFields(c.Fields),
Actions: newUpstreamSlackActions(c.Actions),
}},
}}, nil
}
func newChannelSlackConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
slack := receiver.SlackConfigs[0]
sendResolved := slack.VSendResolved
if err := rejectAnyHTTPAuth(name, slack.HTTPConfig); err != nil {
return nil, err
}
return &ChannelSlackConfig{
SendResolved: &sendResolved,
APIURL: formatSecretURL(slack.APIURL),
Channel: slack.Channel,
Title: valuer.UnsetIfEmpty(slack.Title),
Text: valuer.UnsetIfEmpty(slack.Text),
Color: valuer.UnsetIfEmpty(slack.Color),
TitleLink: valuer.UnsetIfEmpty(slack.TitleLink),
Pretext: valuer.UnsetIfEmpty(slack.Pretext),
Fallback: valuer.UnsetIfEmpty(slack.Fallback),
Footer: valuer.UnsetIfEmpty(slack.Footer),
Fields: newChannelSlackFields(slack.Fields),
Actions: newChannelSlackActions(slack.Actions),
}, nil
}
func newUpstreamSlackFields(fields []ChannelSlackField) []*config.SlackField {
if len(fields) == 0 {
return nil
}
upstream := make([]*config.SlackField, 0, len(fields))
for _, field := range fields {
upstream = append(upstream, &config.SlackField{Title: field.Title, Value: field.Value, Short: field.Short})
}
return upstream
}
func newChannelSlackFields(upstream []*config.SlackField) []ChannelSlackField {
if len(upstream) == 0 {
return nil
}
fields := make([]ChannelSlackField, 0, len(upstream))
for _, field := range upstream {
fields = append(fields, ChannelSlackField{Title: field.Title, Value: field.Value, Short: field.Short})
}
return fields
}
func newUpstreamSlackActions(actions []ChannelSlackAction) []*config.SlackAction {
if len(actions) == 0 {
return nil
}
upstream := make([]*config.SlackAction, 0, len(actions))
for _, action := range actions {
upstreamAction := &config.SlackAction{Type: action.Type, Text: action.Text, URL: action.URL, Style: action.Style, Name: action.Name, Value: action.Value}
if action.Confirm != nil {
upstreamAction.ConfirmField = &config.SlackConfirmationField{Text: action.Confirm.Text, Title: action.Confirm.Title, OkText: action.Confirm.OkText, DismissText: action.Confirm.DismissText}
}
upstream = append(upstream, upstreamAction)
}
return upstream
}
func newChannelSlackActions(upstream []*config.SlackAction) []ChannelSlackAction {
if len(upstream) == 0 {
return nil
}
actions := make([]ChannelSlackAction, 0, len(upstream))
for _, upstreamAction := range upstream {
action := ChannelSlackAction{Type: upstreamAction.Type, Text: upstreamAction.Text, URL: upstreamAction.URL, Style: upstreamAction.Style, Name: upstreamAction.Name, Value: upstreamAction.Value}
if upstreamAction.ConfirmField != nil {
action.Confirm = &ChannelSlackConfirmation{Text: upstreamAction.ConfirmField.Text, Title: upstreamAction.ConfirmField.Title, OkText: upstreamAction.ConfirmField.OkText, DismissText: upstreamAction.ConfirmField.DismissText}
}
actions = append(actions, action)
}
return actions
}
// ChannelEmailConfig carries no SMTP transport fields: the smarthost,
// credentials and TLS settings come from the deployment's global config, so a
// channel can only choose recipients and body.
type ChannelEmailConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
To string `json:"to" required:"true"`
HTML valuer.UnsetOrNonEmptyString `json:"html"`
Headers map[string]string `json:"headers,omitempty"`
}
func (c ChannelEmailConfig) Validate() error {
if c.To == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.to is required for an email channel")
}
// A read reports header names as textproto canonicalizes them, turning
// "subject" into "Subject", so a name that is not already in that form is
// rejected rather than answered with one the caller never sent.
for _, header := range slices.Sorted(maps.Keys(c.Headers)) {
if canonical := textproto.CanonicalMIMEHeaderKey(header); canonical != header {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.headers name %q must be written as %q", header, canonical)
}
}
return nil
}
func (c ChannelEmailConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
EmailConfigs: []*config.EmailConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultEmailConfig.VSendResolved)},
To: c.To,
HTML: c.HTML.StringValue(),
Headers: c.Headers,
}},
}}, nil
}
func newChannelEmailConfigFromReceiver(_ string, receiver *Receiver) (ChannelSpec, error) {
email := receiver.EmailConfigs[0]
sendResolved := email.VSendResolved
return &ChannelEmailConfig{
SendResolved: &sendResolved,
To: email.To,
HTML: valuer.UnsetIfEmpty(email.HTML),
Headers: email.Headers,
}, nil
}
// ChannelWebhookConfig splits apart the two authentication modes the legacy API
// overloaded onto one password field, where an empty username meant the password
// was really a bearer token. Username or Password may be set without the other,
// as upstream allows, but not together with BearerToken.
type ChannelWebhookConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
URL string `json:"url" required:"true" format:"password"`
Username string `json:"username"`
Password string `json:"password" format:"password"`
BearerToken string `json:"bearerToken" format:"password"`
}
func (c ChannelWebhookConfig) Validate() error {
if c.URL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for a webhook channel")
}
usesBasicAuth := c.Username != "" || c.Password != ""
if usesBasicAuth && c.BearerToken != "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.bearerToken cannot be combined with config.spec.username or config.spec.password")
}
return nil
}
func (c ChannelWebhookConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
webhook := &config.WebhookConfig{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultWebhookConfig.VSendResolved)},
URL: config.SecretTemplateURL(c.URL),
}
// Seeded from upstream's default rather than a zero value: FollowRedirects
// and EnableHTTP2 marshal unconditionally, so a zero value would persist
// them as false and read back as a config ChannelWebhookConfig cannot represent.
switch {
case c.Username != "" || c.Password != "":
httpConfig := commoncfg.DefaultHTTPClientConfig
httpConfig.BasicAuth = &commoncfg.BasicAuth{
Username: c.Username,
Password: commoncfg.Secret(c.Password),
}
webhook.HTTPConfig = &httpConfig
case c.BearerToken != "":
httpConfig := commoncfg.DefaultHTTPClientConfig
httpConfig.Authorization = &commoncfg.Authorization{
Type: bearerAuthorizationType,
Credentials: commoncfg.Secret(c.BearerToken),
}
webhook.HTTPConfig = &httpConfig
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
WebhookConfigs: []*config.WebhookConfig{webhook},
}}, nil
}
func newChannelWebhookConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
upstream := receiver.WebhookConfigs[0]
sendResolved := upstream.VSendResolved
if err := rejectUnsupportedHTTPConfig(name, upstream.HTTPConfig); err != nil {
return nil, err
}
if err := rejectHTTPBasicAuthBeyondPassword(name, upstream.HTTPConfig); err != nil {
return nil, err
}
if err := rejectHTTPAuthorizationBeyondBearer(name, upstream.HTTPConfig); err != nil {
return nil, err
}
webhook := &ChannelWebhookConfig{
SendResolved: &sendResolved,
URL: string(upstream.URL),
}
if upstream.HTTPConfig != nil {
if basicAuth := upstream.HTTPConfig.BasicAuth; basicAuth != nil {
webhook.Username = basicAuth.Username
webhook.Password = string(basicAuth.Password)
}
if authorization := upstream.HTTPConfig.Authorization; authorization != nil {
webhook.BearerToken = string(authorization.Credentials)
}
}
return webhook, nil
}
type ChannelPagerdutyConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
RoutingKey string `json:"routingKey" required:"true" format:"password"`
URL string `json:"url"`
Source valuer.UnsetOrNonEmptyString `json:"source"`
Client valuer.UnsetOrNonEmptyString `json:"client"`
ClientURL valuer.UnsetOrNonEmptyString `json:"clientUrl"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Severity string `json:"severity"`
Component string `json:"component"`
Group string `json:"group"`
Class string `json:"class"`
Details map[string]string `json:"details,omitempty"`
}
func (c ChannelPagerdutyConfig) Validate() error {
if c.RoutingKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.routingKey is required for a pagerduty channel")
}
return nil
}
func (c ChannelPagerdutyConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
var eventsURL *config.URL
if c.URL != "" {
parsed, err := parseUpstreamURL(c.URL)
if err != nil {
return nil, err
}
eventsURL = parsed
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
PagerdutyConfigs: []*config.PagerdutyConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultPagerdutyConfig.VSendResolved)},
RoutingKey: config.Secret(c.RoutingKey),
URL: eventsURL,
Source: c.Source.StringValue(),
Client: c.Client.StringValue(),
ClientURL: c.ClientURL.StringValue(),
Description: c.Description.StringValue(),
Severity: c.Severity,
Component: c.Component,
Group: c.Group,
Class: c.Class,
Details: newUpstreamDetails(c.Details),
}},
}}, nil
}
func newChannelPagerdutyConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
pagerduty := receiver.PagerdutyConfigs[0]
sendResolved := pagerduty.VSendResolved
if err := rejectAnyHTTPAuth(name, pagerduty.HTTPConfig); err != nil {
return nil, err
}
var details map[string]string
if len(pagerduty.Details) > 0 {
extracted, err := extractStringDetails(name, pagerduty.Details)
if err != nil {
return nil, err
}
details = extracted
}
return &ChannelPagerdutyConfig{
SendResolved: &sendResolved,
RoutingKey: string(pagerduty.RoutingKey),
URL: formatUpstreamURL(pagerduty.URL),
Source: valuer.UnsetIfEmpty(pagerduty.Source),
Client: valuer.UnsetIfEmpty(pagerduty.Client),
ClientURL: valuer.UnsetIfEmpty(pagerduty.ClientURL),
Description: valuer.UnsetIfEmpty(pagerduty.Description),
Severity: pagerduty.Severity,
Component: pagerduty.Component,
Group: pagerduty.Group,
Class: pagerduty.Class,
Details: details,
}, nil
}
type ChannelOpsgenieConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIKey string `json:"apiKey" required:"true" format:"password"`
APIURL string `json:"apiUrl"`
Message valuer.UnsetOrNonEmptyString `json:"message"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Source valuer.UnsetOrNonEmptyString `json:"source"`
Details map[string]string `json:"details,omitempty"`
Priority string `json:"priority"`
}
func (c ChannelOpsgenieConfig) Validate() error {
if c.APIKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for an opsgenie channel")
}
return nil
}
func (c ChannelOpsgenieConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
var apiURL *config.URL
if c.APIURL != "" {
parsed, err := parseUpstreamURL(c.APIURL)
if err != nil {
return nil, err
}
apiURL = parsed
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
OpsGenieConfigs: []*config.OpsGenieConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultOpsGenieConfig.VSendResolved)},
APIKey: config.Secret(c.APIKey),
APIURL: apiURL,
Message: c.Message.StringValue(),
Description: c.Description.StringValue(),
Source: c.Source.StringValue(),
Priority: c.Priority,
Details: c.Details,
}},
}}, nil
}
func newChannelOpsgenieConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
opsgenie := receiver.OpsGenieConfigs[0]
sendResolved := opsgenie.VSendResolved
if err := rejectAnyHTTPAuth(name, opsgenie.HTTPConfig); err != nil {
return nil, err
}
return &ChannelOpsgenieConfig{
SendResolved: &sendResolved,
APIKey: string(opsgenie.APIKey),
APIURL: formatUpstreamURL(opsgenie.APIURL),
Message: valuer.UnsetIfEmpty(opsgenie.Message),
Description: valuer.UnsetIfEmpty(opsgenie.Description),
Source: valuer.UnsetIfEmpty(opsgenie.Source),
Priority: opsgenie.Priority,
Details: opsgenie.Details,
}, nil
}
type ChannelMSTeamsConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
}
func (c ChannelMSTeamsConfig) Validate() error {
if c.WebhookURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for an msteams channel")
}
return nil
}
func (c ChannelMSTeamsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
webhookURL, err := parseSecretURL(c.WebhookURL)
if err != nil {
return nil, err
}
return &Receiver{Receiver: &config.Receiver{
Name: displayName,
MSTeamsV2Configs: []*config.MSTeamsV2Config{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultMSTeamsV2Config.VSendResolved)},
WebhookURL: webhookURL,
Title: c.Title.StringValue(),
Text: c.Text.StringValue(),
}},
}}, nil
}
func newChannelMSTeamsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
msteams := receiver.MSTeamsV2Configs[0]
sendResolved := msteams.VSendResolved
if err := rejectAnyHTTPAuth(name, msteams.HTTPConfig); err != nil {
return nil, err
}
return &ChannelMSTeamsConfig{
SendResolved: &sendResolved,
WebhookURL: formatSecretURL(msteams.WebhookURL),
Title: valuer.UnsetIfEmpty(msteams.Title),
Text: valuer.UnsetIfEmpty(msteams.Text),
}, nil
}
type ChannelGoogleChatConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Text valuer.UnsetOrNonEmptyString `json:"text"`
}
func (c ChannelGoogleChatConfig) Validate() error {
if c.WebhookURL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for a googlechat channel")
}
return nil
}
func (c ChannelGoogleChatConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
webhookURL, err := parseSecretURL(c.WebhookURL)
if err != nil {
return nil, err
}
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
GoogleChatConfigs: []*GoogleChatReceiverConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultGoogleChatReceiverConfig.VSendResolved)},
WebhookURL: webhookURL,
Title: c.Title.StringValue(),
Text: c.Text.StringValue(),
}},
}, nil
}
func newChannelGoogleChatConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
googlechat := receiver.GoogleChatConfigs[0]
sendResolved := googlechat.VSendResolved
if err := rejectAnyHTTPAuth(name, googlechat.HTTPConfig); err != nil {
return nil, err
}
return &ChannelGoogleChatConfig{
SendResolved: &sendResolved,
WebhookURL: formatSecretURL(googlechat.WebhookURL),
Title: valuer.UnsetIfEmpty(googlechat.Title),
Text: valuer.UnsetIfEmpty(googlechat.Text),
}, nil
}
type ChannelJiraConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
// Site is the Jira Cloud base URL, https://<site>.atlassian.net. Only Jira
// Cloud is supported; the REST base is derived from it.
Site string `json:"site" required:"true"`
Project string `json:"project" required:"true"`
IssueType string `json:"issueType" required:"true"`
Summary valuer.UnsetOrNonEmptyString `json:"summary"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Priority string `json:"priority"`
Labels []string `json:"labels,omitempty"`
ResolveTransition string `json:"resolveTransition"`
ReopenTransition string `json:"reopenTransition"`
ReopenDuration valuer.UnsetOrNonEmptyString `json:"reopenDuration"`
WontFixResolution string `json:"wontFixResolution"`
CustomFields map[string]any `json:"customFields,omitempty"`
Email string `json:"email" required:"true"`
APIToken string `json:"apiToken" required:"true" format:"password"`
}
func (c ChannelJiraConfig) Validate() error {
for _, required := range []struct {
value string
field string
}{
{c.Site, "site"},
{c.Project, "project"},
{c.IssueType, "issueType"},
{c.Email, "email"},
{c.APIToken, "apiToken"},
} {
if required.value == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.%s is required for a jira channel", required.field)
}
}
if !c.ReopenDuration.IsZero() {
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
if err != nil {
return errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q is not a valid duration", c.ReopenDuration)
}
// A read reports the duration as model.Duration formats it, collapsing
// "72h" into "3d", so a value that is not already in that form is rejected
// rather than answered with one the caller never sent.
if canonical := reopenDuration.String(); canonical != c.ReopenDuration.StringValue() {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q must be written as %q", c.ReopenDuration, canonical)
}
}
return nil
}
func (c ChannelJiraConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
// Seeded from upstream's default rather than a zero value: FollowRedirects
// and EnableHTTP2 marshal unconditionally, so a zero value would persist them
// as false and read back as a config ChannelJiraConfig cannot represent.
httpConfig := commoncfg.DefaultHTTPClientConfig
httpConfig.BasicAuth = &commoncfg.BasicAuth{
Username: c.Email,
Password: commoncfg.Secret(c.APIToken),
}
jira := &JiraReceiverConfig{
// JiraReceiverConfig seeds no send_resolved of its own, so unset means off.
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, false)},
Site: c.Site,
Project: c.Project,
IssueType: c.IssueType,
Summary: c.Summary.StringValue(),
Description: c.Description.StringValue(),
Priority: c.Priority,
Labels: c.Labels,
ResolveTransition: c.ResolveTransition,
ReopenTransition: c.ReopenTransition,
WontFixResolution: c.WontFixResolution,
CustomFields: c.CustomFields,
HTTPConfig: &httpConfig,
}
if !c.ReopenDuration.IsZero() {
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "parse reopenDuration %q", c.ReopenDuration)
}
jira.ReopenDuration = reopenDuration
}
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
JiraConfigs: []*JiraReceiverConfig{jira},
}, nil
}
func newChannelJiraConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
jira := receiver.JiraConfigs[0]
sendResolved := jira.VSendResolved
if err := rejectUnsupportedHTTPConfig(name, jira.HTTPConfig); err != nil {
return nil, err
}
if jira.HTTPConfig != nil && jira.HTTPConfig.Authorization != nil {
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization, which is not supported", name)
}
if err := rejectHTTPBasicAuthBeyondPassword(name, jira.HTTPConfig); err != nil {
return nil, err
}
spec := &ChannelJiraConfig{
SendResolved: &sendResolved,
Site: jira.Site,
Project: jira.Project,
IssueType: jira.IssueType,
Summary: valuer.UnsetIfEmpty(jira.Summary),
Description: valuer.UnsetIfEmpty(jira.Description),
Priority: jira.Priority,
Labels: jira.Labels,
ResolveTransition: jira.ResolveTransition,
ReopenTransition: jira.ReopenTransition,
ReopenDuration: valuer.UnsetIfEmpty(jira.ReopenDuration.String()),
WontFixResolution: jira.WontFixResolution,
CustomFields: jira.CustomFields,
}
if jira.HTTPConfig != nil && jira.HTTPConfig.BasicAuth != nil {
spec.Email = jira.HTTPConfig.BasicAuth.Username
spec.APIToken = string(jira.HTTPConfig.BasicAuth.Password)
}
return spec, nil
}
// ChannelJSMOpsConfig carries no API URL: JSM Ops is a single global gateway
// keyed by the integration API key, which the notifier pins itself.
type ChannelJSMOpsConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
APIKey string `json:"apiKey" required:"true" format:"password"`
Message valuer.UnsetOrNonEmptyString `json:"message"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Priority string `json:"priority"`
// Tags is the comma-separated list JSM Ops attaches to the alert.
Tags valuer.UnsetOrNonEmptyString `json:"tags"`
}
func (c ChannelJSMOpsConfig) Validate() error {
if c.APIKey == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for a jsmops channel")
}
return nil
}
func (c ChannelJSMOpsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
JSMOpsConfigs: []*JSMOpsReceiverConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultJSMOpsReceiverConfig.VSendResolved)},
APIKey: config.Secret(c.APIKey),
Message: c.Message.StringValue(),
Description: c.Description.StringValue(),
Priority: c.Priority,
Tags: c.Tags.StringValue(),
}},
}, nil
}
func newChannelJSMOpsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
jsmops := receiver.JSMOpsConfigs[0]
sendResolved := jsmops.VSendResolved
if err := rejectAnyHTTPAuth(name, jsmops.HTTPConfig); err != nil {
return nil, err
}
return &ChannelJSMOpsConfig{
SendResolved: &sendResolved,
APIKey: string(jsmops.APIKey),
Message: valuer.UnsetIfEmpty(jsmops.Message),
Description: valuer.UnsetIfEmpty(jsmops.Description),
Priority: jsmops.Priority,
Tags: valuer.UnsetIfEmpty(jsmops.Tags),
}, nil
}
type ChannelIncidentIOConfig struct {
SendResolved *bool `json:"sendResolved,omitempty"`
URL string `json:"url" required:"true"`
Token string `json:"token" required:"true" format:"password"`
Title valuer.UnsetOrNonEmptyString `json:"title"`
Description valuer.UnsetOrNonEmptyString `json:"description"`
Metadata map[string]string `json:"metadata,omitempty"`
}
func (c ChannelIncidentIOConfig) Validate() error {
if c.URL == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for an incidentio channel")
}
if c.Token == "" {
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.token is required for an incidentio channel")
}
return nil
}
func (c ChannelIncidentIOConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
return &Receiver{
Receiver: &config.Receiver{Name: displayName},
IncidentIOConfigs: []*IncidentIOReceiverConfig{{
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultIncidentIOReceiverConfig.VSendResolved)},
URL: c.URL,
Token: config.Secret(c.Token),
Title: c.Title.StringValue(),
Description: c.Description.StringValue(),
Metadata: c.Metadata,
}},
}, nil
}
func newChannelIncidentIOConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
incidentio := receiver.IncidentIOConfigs[0]
sendResolved := incidentio.VSendResolved
if err := rejectAnyHTTPAuth(name, incidentio.HTTPConfig); err != nil {
return nil, err
}
return &ChannelIncidentIOConfig{
SendResolved: &sendResolved,
URL: incidentio.URL,
Token: string(incidentio.Token),
Title: valuer.UnsetIfEmpty(incidentio.Title),
Description: valuer.UnsetIfEmpty(incidentio.Description),
Metadata: incidentio.Metadata,
}, nil
}
// ════════════════════════════════════════════════════════════════════════
// Helpers
// ════════════════════════════════════════════════════════════════════════

View File

@@ -258,6 +258,6 @@ type TokenStore interface {
// Delete a token by userID.
DeleteByUserID(context.Context, valuer.UUID) error
// Update last observed at by access token.
UpdateLastObservedAtByAccessToken(context.Context, []map[string]any) error
// Update last observed at of the given tokens.
UpdateLastObservedAt(context.Context, []*StorableToken) error
}

View File

@@ -208,6 +208,35 @@ func NewGettableUnmappedModels(items []*UnmappedModel) *GettableUnmappedModels {
}
}
func (u *UpdatableLLMPricingRule) UnmarshalJSON(data []byte) error {
type Alias UpdatableLLMPricingRule
var temp Alias
if err := json.Unmarshal(data, &temp); err != nil {
return err
}
*u = UpdatableLLMPricingRule(temp)
return u.Validate()
}
// Validate mirrors the collector's pattern check: at least one pattern, none
// empty, all valid path.Match globs.
func (u *UpdatableLLMPricingRule) Validate() error {
if len(u.ModelPattern) == 0 {
return errors.Newf(errors.TypeInvalidInput, ErrCodePricingRuleInvalidInput, "model %q: modelPattern must contain at least one pattern", u.Model)
}
for _, p := range u.ModelPattern {
if p == "" {
return errors.Newf(errors.TypeInvalidInput, ErrCodePricingRuleInvalidInput, "model %q: modelPattern must not contain an empty pattern", u.Model)
}
if _, err := path.Match(p, ""); err != nil {
return errors.Newf(errors.TypeInvalidInput, ErrCodePricingRuleInvalidInput, "model %q: modelPattern %q is not a valid glob", u.Model, p)
}
}
return nil
}
func NewLLMPricingRuleFromUpdatable(u *UpdatableLLMPricingRule, orgID valuer.UUID, userEmail string, now time.Time) *LLMPricingRule {
id := valuer.GenerateUUID()
if u.ID != nil {

View File

@@ -1,6 +1,7 @@
package llmpricingruletypes
import (
"encoding/json"
"os"
"path/filepath"
"testing"
@@ -126,3 +127,34 @@ func TestGenerateCollectorConfig_EmptyInputPassthrough(t *testing.T) {
assert.Equal(t, in, out)
}
}
func TestUpdatableLLMPricingRuleUnmarshalJSON(t *testing.T) {
tests := []struct {
name string
pattern string
wantErr bool
}{
{name: "valid", pattern: `["gpt-4o*", "gpt-4o"]`},
{name: "missing", pattern: ``, wantErr: true},
{name: "null", pattern: `null`, wantErr: true},
{name: "empty_list", pattern: `[]`, wantErr: true},
{name: "empty_entry", pattern: `["gpt-4o*", ""]`, wantErr: true},
{name: "bad_glob", pattern: `["gpt-["]`, wantErr: true},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
body := `{"modelName": "gpt-4o"}`
if tc.pattern != "" {
body = `{"modelName": "gpt-4o", "modelPattern": ` + tc.pattern + `}`
}
var req UpdatableLLMPricingRules
err := json.Unmarshal([]byte(`{"rules": [`+body+`]}`), &req)
if tc.wantErr {
assert.Error(t, err)
} else {
assert.NoError(t, err)
}
})
}
}

View File

@@ -1,37 +1,51 @@
package telemetrystoretypes
import (
"github.com/SigNoz/signoz/pkg/errors"
"encoding/json"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
"github.com/bytedance/sonic"
)
var ErrCodeUnmarshalJSONColumn = errors.MustNewCode("fail_unmarshal_json_column")
// JSONValue is the scan target for a ClickHouse JSON column: the connection sets
// output_format_native_write_json_as_string, so the column arrives as a raw document rather than
// the chcol.JSON the driver reports as its scan type.
type JSONValue map[string]any
// Scan decodes into a fresh map every time: a scan target is reused across rows, and unmarshalling
// into the map already there would both keep its keys and hand every row the same map.
func (v *JSONValue) Scan(src any) error {
var raw []byte
switch value := src.(type) {
case nil:
*v = nil
// NestedJSON decodes a native JSON column into a nested document via the driver's own marshaler, so
// arrays of objects and typed sub-paths survive and Dynamic values arrive unwrapped. A key stored as
// both a scalar and an object collapses, as the nested form cannot hold both.
func NestedJSON(j chcol.JSON) map[string]any {
raw, err := j.MarshalJSON()
if err != nil {
return nil
case string:
raw = []byte(value)
case []byte:
raw = value
default:
return errors.NewInternalf(ErrCodeUnmarshalJSONColumn, "cannot decode %T as a JSON column", src)
}
decoded := JSONValue{}
if err := sonic.Unmarshal(raw, &decoded); err != nil {
return errors.WrapInternalf(err, ErrCodeUnmarshalJSONColumn, "failed to unmarshal JSON column")
var out map[string]any
if err := sonic.Unmarshal(raw, &out); err != nil {
return nil
}
*v = decoded
return nil
return out
}
// FlattenJSON decodes a native JSON column into its leaf paths as dotted keys, so a key stored as
// both a scalar and an object survives as two distinct keys — unlike the nested form, which cannot
// hold both. Dynamic values arrive unwrapped, arrays of objects intact.
func FlattenJSON(j chcol.JSON) map[string]any {
paths := j.ValuesByPath()
out := make(map[string]any, len(paths))
for path, value := range paths {
out[path] = decodePathValue(value)
}
return out
}
func decodePathValue(value any) any {
variant, ok := value.(chcol.Variant)
if !ok {
return value
}
raw, err := json.Marshal(variant)
if err != nil {
return nil
}
var out any
if err := sonic.Unmarshal(raw, &out); err != nil {
return nil
}
return out
}

View File

@@ -0,0 +1,82 @@
package telemetrystoretypes
import (
"testing"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
"github.com/stretchr/testify/assert"
)
func TestNestedJSON(t *testing.T) {
testCases := []struct {
name string
paths map[string]any
want map[string]any
}{
{
name: "Empty",
paths: nil,
want: map[string]any{},
},
{
name: "FlatScalars",
paths: map[string]any{"level": "error", "status": int64(500)},
want: map[string]any{"level": "error", "status": float64(500)},
},
{
name: "DottedPathsBecomeNested",
paths: map[string]any{"attrs.code": int64(500), "attrs.path": "/checkout"},
want: map[string]any{"attrs": map[string]any{"code": float64(500), "path": "/checkout"}},
},
{
name: "ArrayOfObjectsPreserved",
paths: map[string]any{"education": []any{map[string]any{"name": "IIT"}}},
want: map[string]any{"education": []any{map[string]any{"name": "IIT"}}},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
j := chcol.NewJSON()
for path, value := range testCase.paths {
j.SetValueAtPath(path, value)
}
assert.Equal(t, testCase.want, NestedJSON(*j))
})
}
}
func TestFlattenJSON(t *testing.T) {
testCases := []struct {
name string
paths map[string]any
want map[string]any
}{
{
name: "Empty",
paths: nil,
want: map[string]any{},
},
{
name: "DottedPathsStayFlat",
paths: map[string]any{"http.method": "GET", "level": "error"},
want: map[string]any{"http.method": "GET", "level": "error"},
},
{
// A scalar and an object under the same prefix survive as two distinct dotted keys.
name: "ScalarAndObjectKey_BothSurvive",
paths: map[string]any{"scope": "x", "scope.attributes.name": "y"},
want: map[string]any{"scope": "x", "scope.attributes.name": "y"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
j := chcol.NewJSON()
for path, value := range testCase.paths {
j.SetValueAtPath(path, value)
}
assert.Equal(t, testCase.want, FlattenJSON(*j))
})
}
}

View File

@@ -130,3 +130,17 @@ def test_bulk_sync(
assert all(r["pricing"]["input"] == 5 for r in stored)
delete_all_llm_pricing_rules(signoz, token)
def test_rejects_rule_without_pattern(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
):
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
delete_all_llm_pricing_rules(signoz, token)
rules = zeus_rules(10)
rules[1]["modelPattern"] = []
assert upsert_llm_pricing_rules(signoz, token, rules).status_code == HTTPStatus.BAD_REQUEST
assert list_llm_pricing_rules(signoz, token) == []

View File

@@ -0,0 +1,36 @@
import time
from collections.abc import Callable
from http import HTTPStatus
import requests
from sqlalchemy import sql
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
def test_last_observed_at_is_flushed(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
"""Verify the tokenizer GC persists the cached last observed at of a used token to the sql store."""
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
headers={"Authorization": f"Bearer {token}"},
timeout=5,
)
assert response.status_code == HTTPStatus.OK
deadline = time.time() + 30
while time.time() < deadline:
with signoz.sqlstore.conn.connect() as conn:
row = conn.execute(
sql.text("SELECT last_observed_at FROM auth_token WHERE access_token = :access_token"),
{"access_token": token},
).fetchone()
if row is not None and row[0] is not None:
return
time.sleep(1)
raise AssertionError("last_observed_at was not flushed to the sql store within 30s")

View File

@@ -0,0 +1,33 @@
import pytest
from testcontainers.core.container import Network
from fixtures import types
from fixtures.signoz import create_signoz
@pytest.fixture(name="signoz", scope="package")
def signoz_passwordauthn(
network: Network,
zeus: types.TestContainerDocker,
gateway: types.TestContainerDocker,
sqlstore: types.TestContainerSQL,
clickhouse: types.TestContainerClickhouse,
request: pytest.FixtureRequest,
pytestconfig: pytest.Config,
) -> types.SigNoz:
"""
Package-scoped fixture for SigNoz with a short tokenizer GC interval so the last observed at flush runs within a test.
"""
return create_signoz(
network=network,
zeus=zeus,
gateway=gateway,
sqlstore=sqlstore,
clickhouse=clickhouse,
request=request,
pytestconfig=pytestconfig,
cache_key="signoz-passwordauthn",
env_overrides={
"SIGNOZ_TOKENIZER_OPAQUE_GC_INTERVAL": "5s",
},
)