Compare commits

...

19 Commits

Author SHA1 Message Date
Naman Verma
1c3b981179 chore: add api to retry migration for a dashboard 2026-08-04 00:09:40 +05:30
Naman Verma
cfaa7de165 fix: enforce the required tag on dashboard spec fields (#12381)
Some checks are pending
build-staging / prepare (push) Waiting to run
build-staging / js-build (push) Blocked by required conditions
build-staging / go-build (push) Blocked by required conditions
build-staging / staging (push) Blocked by required conditions
cacheci / tests (push) Waiting to run
Release Drafter / update_release_draft (push) Waiting to run
* fix: enforace the required tag on dashboard spec fields

* test: add empty list and objects for required fields in integration tests
2026-08-03 18:22:43 +00:00
Naman Verma
5b3cc2400f chore: delta temporality metrics should always be considered as a sum metric (#12313)
* chore: delta temporality metrics should always be considered as a sum metric

* test: add integration tests

* test: parametrise the non-reduced test for rate and increase both

* test: add comment explaining last samples in test

* chore: remove unneeded comments

---------

Co-authored-by: Srikanth Chekuri <srikanth.chekuri92@gmail.com>
2026-08-03 16:08:45 +00:00
Tushar Vats
8eb3f6bc1b refactor(querier): squash statement-builder config under querier (#12385)
- Embed statementbuilder.Config into querier.Config with mapstructure ",squash";
  keys move to querier.skip_resource_fingerprint.* (env SIGNOZ_QUERIER_*).
- Drop the standalone statementbuilder section and its config factory.
- Pass statementbuilder.Config wholesale into NewLogQueryStatementBuilder.
2026-08-03 15:40:07 +00:00
Srikanth Chekuri
1ac6103685 feat(querier): wire clickhouseprometheusv2 for shadow comparison and pinned serving (#12324)
* feat(querier): wire clickhouseprometheusv2 for shadow comparison and pinned serving

Stand up the v2 provider next to the default one and give the querier two
flag-gated ways to exercise it, neither affecting default serving:

- shadow: with use_prometheus_clickhouse_v2 on, every PromQL query re-runs
  on v2 after the response is sent; result diffs are logged. Bounded by a
  small per-process admission cap (skip, not queue, at the cap).
- pin: the X-SigNoz-PromQL-Provider header serves the response from v2
  directly, for side-by-side comparison by integration tests and support.
  Pinned requests bypass the cache in both directions.

Typed storage errors (series/sample budgets) survive the engine's wrapping
and surface as 4xx instead of internal 500s. PromQL results now carry
ClickHouse scan stats.

prometheus::provider: clickhousev2 also makes v2 the serving provider
outright (no shadow in that mode - nothing to compare against).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* test(promql): replay the conformance corpus on both providers

Every corpus case now runs twice — default provider, and pinned to
clickhousev2 via the flag-gated X-SigNoz-PromQL-Provider header — each leg
asserted against the same frozen expectations, each leg with its own
known-divergences ledger enforced in both directions. The new
known_divergences_v2.json is the rollout scorecard: the provider swap is
measured by burning it to empty.

The legs are never asserted against each other: both can sit within one
rounding quantum of the expected value yet differ by up to two quanta at a
rounding boundary, so a leg-vs-leg check would reintroduce the boundary
noise the tolerance absorbs. Anchoring both to the same oracle over the
same ingested bytes already localizes any disagreement to a provider.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2026-08-03 15:00:52 +00:00
Tushar Vats
302a40e8df chore(lint): satisfy golangci-lint v2.12.2 and de-flake http middleware tests (#12383)
primus bumped golangci-lint to v2.12.2, whose govet now runs the inline
analyzer and whose sloglint is stricter. CI resolves primus.workflows@main,
so every PR started failing lint the moment that landed.

- reflect.Ptr is a deprecated alias carrying //go:fix inline, so it is now
  reflect.Pointer at all five call sites
- metricsstatementbuilder imported golang.org/x/exp/slices, which carries
  //go:fix inline pointing at the stdlib; the analyzer cannot inline
  generics, so switch the import to stdlib slices as the directive intends
- pkg/instrumentation/loghandler emits OpenTelemetry semantic-convention
  attributes (code.filepath, exception.type, ...), which are dotted rather
  than snake_case by definition. Renaming them would break every log
  consumer, so the keys move to constants in instrumentationtypes, which
  already held this kind of key -- and already defined code.function, so
  source.go was duplicating it. Six of the seven alias the semconv
  constants that define them; exception.code has no OTel equivalent.
  sloglint resolves a same-package constant back to its literal but skips
  a qualified one, so this needs no exclusion.

CI reported 8 issues but capped at max-same-issues=3, hiding 2 more
reflect.Ptr sites and 4 more sloglint ones.

Separately, TestTimeout/WaitTillNoTimeoutForExcludedPath failed with
"transport connection broken: http: CloseIdleConnections called".
TestTimeout and TestCache issue requests through http.DefaultClient while
a parallel subtest in response_test.go closes an httptest.Server, and
httptest.Server.Close calls http.DefaultTransport.CloseIdleConnections.
Both tests now use their own client and transport, and are closed via
t.Cleanup. That makes Serve return ErrServerClosed on every run, so the
require.NoError wrapping it is dropped -- it could never have held, and
require runs t.FailNow off the test goroutine anyway. Bare Serve in a
goroutine matches routerweb and render tests.
2026-08-03 14:41:16 +00:00
Pandey
87560eaeb9 feat(querybuilder): allow row-generator table functions (#12376)
Some checks failed
build-staging / staging (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
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
* feat(querybuilder): allow row-generator table functions

The table-function rule refuses everything, which costs a false positive
on queries that use numbers() or generateSeries() to build a dense
interval axis to join a sparse series against. Neither reads through
anything: they compute their rows from their arguments, open no file or
socket, reach no other host, and name no table, database or dictionary.

Everything else stays refused. Most table functions read through
something, and merge('system', '.*'), remote() and cluster() reach the
internal databases without producing a TableIdentifier for the database
rule to catch, so this rule is all that sees them. generateRandom is pure
but streams rows its arguments do not bound, so it stays out too.

Arguments are visited before the table function itself, which is what
keeps an allowed generator from being usable as a wrapper to smuggle a
read through.

* test(querybuilder): address review on the generator allow list

Names the allowed table functions in the rejection error, collapses the
added comments, and covers the generators in JOIN, CTE, subquery and
UNION position alongside the reads they must not be usable to smuggle.

* refactor(querybuilder): derive the allowed table function message from the map

Values the map by the spelling to name back to the caller, so the message
comes off the same data the lookup uses and cannot drift from it. Drops
the test that existed only to catch that drift.
2026-08-03 10:28:52 +00:00
Abhi kumar
98030c29ef fix(dashboard): substitute per-row groupBy values in context-link URLs (#12328)
Per-row field variables are registered `_`-prefixed in useContextVariables so
a clicked row's field can't shadow a dashboard variable of the same name, but
resolveText/resolveTextWithTruncation only ever looked up the bare placeholder
name. A link authored as `/trace/{{trace_id}}` therefore never matched the
registered `_trace_id`, was left literal, and got percent-encoded into
`/trace/%7B%7Btrace_id%7D%7D` on navigation.

Look up the exact name first — so dashboard/global variables keep precedence —
then fall back to the `_`-prefixed field value. Fixes every placeholder syntax
({{var}}, {{.var}}, [[var]], $var) and both V1 and V2 dashboards, which share
this resolver and the same `_` convention.

Closes #11325
2026-08-03 08:31:21 +00:00
Pandey
b1b668925e chore(deps): bump clickhouse-sql-parser to v0.5.4 (#12375)
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
* chore(deps): bump clickhouse-sql-parser to v0.5.4

v0.5.4 carries the four SigNoz-reported grammar fixes: unquoted `interval`
as a column name, GLOBAL before a join type and GLOBAL NOT IN, arithmetic
inside table function arguments, and the SQL-standard keyword-argument
forms of trim/substring/overlay.

Replaying the saved production corpus takes the v5 clickhouse_sql rejection
rate from 11/198 to 1/198, and the one remaining rejection is a real rule
hit rather than a parser gap. TestErrIfStatementIsNotValid_ShouldPassButFails
therefore has nothing left to hold; its queries move to _Pass as regression
canaries.

* test(querybuilder): bound how long a valid statement may take to parse

Telling an INTERVAL operator from a column named interval needs
backtracking, and v0.5.4 keeps that from going exponential by remembering
the offsets it has already failed at. Losing that memoisation would hang
the parser rather than fail it, so a 30-term case covers the shape and the
loop bounds how long any valid statement may take.

* test(querybuilder): address review on the parser bump

Collapse the added comments to fewer lines, and fail the timeout branch
with assert rather than t.Fatal so a case that does hang reports alongside
the rest.

* test(querybuilder): drop the stale signed-literal comment

* test(querybuilder): pin the numbers table function false positive

Restores TestErrIfStatementIsNotValid_ShouldPassButFails around the one
rejection the saved production corpus still produces on v0.5.4. numbers()
generates rows rather than reading through anything, so this is the
blanket table-function rule being stricter than the threat it exists for
rather than a parser gap, and the table now asserts the error code.

* test(querybuilder): trim the numbers case comment and query

Collapses the query to a single line to match the other cases in the file
and swaps the tenant's metric name for a placeholder.
2026-08-02 11:35:08 +00:00
dependabot[bot]
2d90a9f5eb chore(deps): bump golang.org/x/net in /scripts/promqltestcorpus (#12373)
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
Bumps [golang.org/x/net](https://github.com/golang/net) from 0.54.0 to 0.55.0.
- [Commits](https://github.com/golang/net/compare/v0.54.0...v0.55.0)

---
updated-dependencies:
- dependency-name: golang.org/x/net
  dependency-version: 0.55.0
  dependency-type: indirect
...

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
2026-08-02 09:28:15 +00:00
dependabot[bot]
f75cd0854f chore(deps): bump golang.org/x/crypto in /scripts/promqltestcorpus (#12320)
Bumps [golang.org/x/crypto](https://github.com/golang/crypto) from 0.50.0 to 0.52.0.
- [Commits](https://github.com/golang/crypto/compare/v0.50.0...v0.52.0)

---
updated-dependencies:
- dependency-name: golang.org/x/crypto
  dependency-version: 0.52.0
  dependency-type: indirect
...

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
2026-08-02 08:54:18 +00:00
Pandey
ab91995ee5 ci: restore shared go and pnpm caches in integration and e2e jobs (#12369)
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
* feat(tests): add --signoz-image option to run a prebuilt image

* ci: collapse integrationci and e2eci into testsci with a shared image build

* ci: rename flag to --zeus-network-aliases

* ci: use singular --zeus-network-alias derived from a shared env var

* revert: drop the testsci workflow and the --signoz-image option

* ci: restore shared go and pnpm caches in integration and e2e jobs

* ci: inject cache tarballs instead of directory trees
2026-08-01 18:58:17 +00:00
Pandey
fd2bd85e90 ci(cacheci): extract cache mounts via docker cp instead of the local exporter (#12370)
* ci(cacheci): extract cache mounts via docker cp instead of the local exporter

* ci(cacheci): make extracted cache writable for unprivileged untar

* ci(cacheci): move cache contents as tarballs

* ci(cacheci): drop temporary pr trigger and cancel superseded runs
2026-08-01 17:45:15 +00:00
Pandey
c110e14505 ci: add cacheci workflow to maintain shared go and pnpm caches (#12368)
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
2026-08-01 16:11:09 +00:00
Pandey
5917f9fe31 perf(tests): cache go and pnpm stores across integration image builds (#12366)
* perf(tests): cache go and pnpm stores across integration image builds

Add BuildKit cache mounts for GOCACHE/GOMODCACHE and the pnpm store to the
integration Dockerfiles, and build the image via the docker CLI (docker-py,
used by testcontainers' DockerImage, does not support BuildKit). Embed the
go build command directly so Makefile changes do not invalidate the build
layer, and pin HOME/GOCACHE/GOMODCACHE/PNPM_HOME explicitly so cache-mount
targets match tool defaults by contract. The with-web node stage fetches
dependencies from the lockfile before the source copy, so frontend edits
only re-run the offline install and build.

* feat(tests): add --clean flag to prune buildkit cache mounts

The go and pnpm caches introduced for the integration image build survive
--teardown since they belong to the docker builder, not to any container.
--clean runs docker builder prune with a type=exec.cachemount filter at
session start, forcing the next image build to start cold. Documented in
the integration testing guide.

* feat(tests): add --rebuild flag to refresh the signoz container under --reuse

--reuse keeps the running signoz container, so backend source changes are
never picked up without tearing down the whole stack. --rebuild deletes the
cached signoz container and recreates it from the current sources (an
incremental image build), while databases, mocks and migrations stay reused.
Requires --reuse; combining with --teardown or --clean is a usage error.

* chore(tests): prune comments to non-obvious constraints

* docs(tests): make py-test-setup rebuild signoz and audit the integration guide

py-test-setup now passes --rebuild so re-running it after backend changes
transparently swaps in a signoz container built from the current sources.
The integration guide documents the iteration loop and fixes stale content:
option defaults (clickhouse 25.12.5, schema migrator v0.144.6), the
nonexistent --zookeeper-version option, Zookeeper vs ClickHouse Keeper, the
e2e doc path, and the lint toolchain (ruff).

* docs(tests): wire --rebuild into the e2e setup flow

The e2e bootstrap shares the signoz fixture, so --rebuild already applies;
with --with-web it also picks up frontend changes since the image bakes the
built frontend in. The setup command now passes --rebuild, and the guide
documents the iteration loop, the --rebuild/--clean flags, ClickHouse Keeper
instead of Zookeeper, and the corrected integration doc path.

* docs(tests): qualify --rebuild workflow for suites with custom signoz variants

make py-test-setup only rebuilds the default signoz instance; suites that
create their own via create_signoz(cache_key=...) keep a separately cached
container. Passing --rebuild on the suite run itself rebuilds every variant
that run instantiates.

* docs(tests): describe --clean behaviour instead of its exact command

Keeps the docs from drifting if the prune invocation behind --clean changes.
2026-08-01 13:52:19 +00:00
Srikanth Chekuri
052255abf5 feat(prometheus): add clickhouseprometheusv2 native read path (#12323)
Second-generation ClickHouse-backed Prometheus provider: the stock engine
evaluates over a native storage.Querier instead of the v1 remote-read
adapter. Per-selector fetch windows, last-sample-per-step reduction for
subquery-free instant selectors (gated on prometheus.QueryTraits),
identical-labelset merge, per-type __name__ matchers and anchored regexes,
inclusive series-lookup bounds for the exporter's hour-floored registration
rows.

Not wired: no factory registration, no config selection, nothing serves
from this package yet. Fetch budgets (series/sample ceilings) are
deliberately left out for now and will come separately.

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 13:15:07 +00:00
Pandey
afc90f532f chore: remove doc.go files (#12365) 2026-08-01 12:35:54 +00:00
Pandey
24390d7192 fix(ruletypes): expose above_or_equal and below_or_equal in CompareOperator enum (#12360)
* fix(ruletypes): expose above_or_equal and below_or_equal in CompareOperator enum

The operators are accepted by Validate(), normalized, evaluated and
returned by the rules API, but were commented out of Enum(), so the
generated OpenAPI spec (and clients generated from it, e.g.
terraform-provider-signoz) rejected rules the server itself creates.

* fix(alerts): support above_or_equal and below_or_equal operators in CreateAlertV2

Adds the two inclusive operators to the v2 alert form: selectable in the
threshold operator dropdown, normalized from all backend aliases
(5/6, above_or_eq/below_or_eq, >=/<=), rendered with their symbols in
threshold rows and match-type tooltips, and prefilled losslessly from
dashboard panel thresholds instead of collapsing onto the strict
variants. The v1 form is left untouched.
2026-08-01 11:07:51 +00:00
Nikhil Mantri
16c6fdc600 chore: combined function for goroutines (#12362)
Some checks failed
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
build-staging / prepare (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
2026-08-01 07:06:09 +00:00
129 changed files with 3872 additions and 596 deletions

92
.github/workflows/cacheci.yml vendored Normal file
View File

@@ -0,0 +1,92 @@
name: cacheci
on:
push:
branches:
- main
workflow_dispatch:
permissions:
contents: read
actions: write
# Cancelling mid-rotation is safe: the sequential delete-then-save order
# leaves at most one key missing at any moment.
concurrency:
group: cacheci
cancel-in-progress: true
jobs:
tests:
runs-on: ubuntu-latest
steps:
- name: checkout
uses: actions/checkout@v4
- name: restore
id: restore
uses: actions/cache/restore@v4
with:
path: ${{ runner.temp }}/cacheci
key: tests-primary
restore-keys: |
tests-secondary
- name: inject
if: steps.restore.outputs.cache-matched-key != ''
run: |
cat > "$RUNNER_TEMP/inject.Dockerfile" <<'EOF'
FROM busybox:1.37
RUN --mount=type=cache,target=/root/.cache/go-build \
--mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/pnpm/store \
--mount=type=bind,target=/restored \
tar -xf /restored/go-build.tar -C /root/.cache/go-build && \
tar -xf /restored/go-mod.tar -C /go/pkg/mod && \
tar -xf /restored/pnpm-store.tar -C /pnpm/store
EOF
docker build -f "$RUNNER_TEMP/inject.Dockerfile" "$RUNNER_TEMP/cacheci"
- name: build
run: |
docker build -f cmd/enterprise/Dockerfile.integration --build-arg TARGETARCH=amd64 --build-arg ZEUSURL=http://zeus:8080 .
docker build -f cmd/enterprise/Dockerfile.with-web.integration --build-arg TARGETARCH=amd64 --build-arg ZEUSURL=http://zeus:8080 .
# docker cp instead of --output type=local (the local exporter stalls on
# multi-GB outputs); tarballs instead of raw trees so the host never hits
# the permission and symlink semantics that broke docker cp.
- name: extract
run: |
rm -rf "$RUNNER_TEMP/cacheci"
mkdir -p "$RUNNER_TEMP/cacheci" "$RUNNER_TEMP/extract-context"
cat > "$RUNNER_TEMP/extract.Dockerfile" <<'EOF'
FROM busybox:1.37
RUN --mount=type=cache,target=/root/.cache/go-build \
--mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/pnpm/store \
mkdir -p /out && \
tar -cf /out/go-build.tar -C /root/.cache/go-build . && \
tar -cf /out/go-mod.tar -C /go/pkg/mod . && \
tar -cf /out/pnpm-store.tar -C /pnpm/store .
EOF
docker build -f "$RUNNER_TEMP/extract.Dockerfile" -t cacheci-extract "$RUNNER_TEMP/extract-context"
id=$(docker create cacheci-extract)
docker cp "$id":/out/. "$RUNNER_TEMP/cacheci/"
docker rm "$id"
# Fixed cache keys are immutable, so each key must be deleted before it
# can be saved again. Rotating primary and secondary one after the other
# keeps at least one key restorable for concurrent test runs.
- name: delete-primary
env:
GH_TOKEN: ${{ github.token }}
run: gh cache delete tests-primary --repo "$GITHUB_REPOSITORY" || true
- name: save-primary
uses: actions/cache/save@v4
with:
path: ${{ runner.temp }}/cacheci
key: tests-primary
- name: delete-secondary
env:
GH_TOKEN: ${{ github.token }}
run: gh cache delete tests-secondary --repo "$GITHUB_REPOSITORY" || true
- name: save-secondary
uses: actions/cache/save@v4
with:
path: ${{ runner.temp }}/cacheci
key: tests-secondary

View File

@@ -75,6 +75,30 @@ jobs:
docker rm pw
echo "PLAYWRIGHT_BROWSERS_PATH=$RUNNER_TEMP/ms-playwright" >> "$GITHUB_ENV"
cd tests/e2e && pnpm playwright install-deps ${{ matrix.project }}
# Restore-only: the cacheci workflow owns cache saves. Seeds the
# BuildKit cache mounts so the in-test image build is incremental.
- name: restore
id: restore
uses: actions/cache/restore@v4
with:
path: ${{ runner.temp }}/cacheci
key: tests-primary
restore-keys: |
tests-secondary
- name: inject
if: steps.restore.outputs.cache-matched-key != ''
run: |
cat > "$RUNNER_TEMP/inject.Dockerfile" <<'EOF'
FROM busybox:1.37
RUN --mount=type=cache,target=/root/.cache/go-build \
--mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/pnpm/store \
--mount=type=bind,target=/restored \
tar -xf /restored/go-build.tar -C /root/.cache/go-build && \
tar -xf /restored/go-mod.tar -C /go/pkg/mod && \
tar -xf /restored/pnpm-store.tar -C /pnpm/store
EOF
docker build -f "$RUNNER_TEMP/inject.Dockerfile" "$RUNNER_TEMP/cacheci"
- name: bring-up-stack
run: |
cd tests && \

View File

@@ -110,6 +110,30 @@ jobs:
sudo mv chromedriver-linux64/chromedriver /usr/local/bin/chromedriver
chromedriver -version
google-chrome-stable --version
# Restore-only: the cacheci workflow owns cache saves. Seeds the
# BuildKit cache mounts so the in-test image build is incremental.
- name: restore
id: restore
uses: actions/cache/restore@v4
with:
path: ${{ runner.temp }}/cacheci
key: tests-primary
restore-keys: |
tests-secondary
- name: inject
if: steps.restore.outputs.cache-matched-key != ''
run: |
cat > "$RUNNER_TEMP/inject.Dockerfile" <<'EOF'
FROM busybox:1.37
RUN --mount=type=cache,target=/root/.cache/go-build \
--mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/pnpm/store \
--mount=type=bind,target=/restored \
tar -xf /restored/go-build.tar -C /root/.cache/go-build && \
tar -xf /restored/go-mod.tar -C /go/pkg/mod && \
tar -xf /restored/pnpm-store.tar -C /pnpm/store
EOF
docker build -f "$RUNNER_TEMP/inject.Dockerfile" "$RUNNER_TEMP/cacheci"
- name: run
run: |
cd tests && \

View File

@@ -209,8 +209,8 @@ py-lint: ## Run ruff check across the shared tests project
@cd tests && uv run ruff check --fix .
.PHONY: py-test-setup
py-test-setup: ## Bring up the shared SigNoz backend used by integration and e2e tests
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --capture=no integration/bootstrap/setup.py::test_setup
py-test-setup: ## Bring up the shared SigNoz backend used by integration and e2e tests, rebuilding signoz from the current sources
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --rebuild --capture=no integration/bootstrap/setup.py::test_setup
.PHONY: py-test-teardown
py-test-teardown: ## Tear down the shared SigNoz backend

View File

@@ -4,9 +4,13 @@ ARG OS="linux"
ARG TARGETARCH
ARG ZEUSURL
# HOME comes from the build user, not the image config; declare it so the
# /root paths below trace to it.
ENV HOME=/root
# This path is important for stacktraces
WORKDIR $GOPATH/src/github.com/signoz/signoz
WORKDIR /root
WORKDIR $HOME
RUN set -eux; \
apt-get update; \
@@ -14,23 +18,36 @@ RUN set -eux; \
g++ \
gcc \
libc6-dev \
make \
pkg-config \
; \
rm -rf /var/lib/apt/lists/*
# Keep the literal cache-mount targets below in sync with these. The caches
# are shared with Dockerfile.with-web.integration (same target paths).
ENV GOCACHE=$HOME/.cache/go-build
ENV GOMODCACHE=$GOPATH/pkg/mod
COPY go.mod go.sum ./
RUN go mod download
RUN --mount=type=cache,target=/go/pkg/mod \
go mod download
COPY ./cmd/ ./cmd/
COPY ./ee/ ./ee/
COPY ./pkg/ ./pkg/
COPY ./templates /root/templates
COPY Makefile Makefile
RUN TARGET_DIR=/root ARCHS=${TARGETARCH} ZEUS_URL=${ZEUSURL} LICENSE_URL=${ZEUSURL}/api/v1 make go-build-enterprise-race
RUN mv /root/linux-${TARGETARCH}/signoz /root/signoz
# Invoked directly instead of via make so Makefile changes don't invalidate
# this layer; the Makefile's git-derived ldflags resolve to empty in here
# anyway (.git is dockerignored).
RUN --mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/root/.cache/go-build \
GOARCH=${TARGETARCH} GOOS=${OS} go build -C ./cmd/enterprise -race -tags timetzdata -o /root/signoz \
-ldflags "-s -w \
-X github.com/SigNoz/signoz/pkg/version.version=integration \
-X github.com/SigNoz/signoz/pkg/version.variant=enterprise \
-X github.com/SigNoz/signoz/ee/zeus.url=${ZEUSURL} \
-X github.com/SigNoz/signoz/ee/zeus.deprecatedURL=${ZEUSURL}/api/v1"
RUN chmod 755 /root /root/signoz

View File

@@ -1,10 +1,23 @@
FROM node:22-bookworm AS build
WORKDIR /opt/
COPY ./frontend/ ./
# HOME comes from the build user, not the image config.
ENV HOME=/root
# pnpm's store lives at $PNPM_HOME/store — a dedicated directory pnpm
# manages. Keep the literal cache-mount targets below in sync.
ENV PNPM_HOME=/pnpm
ENV NODE_OPTIONS=--max-old-space-size=8192
RUN CI=1 npm i -g pnpm@10
RUN CI=1 pnpm install
# pnpm fetch resolves from the lockfile alone and runs no lifecycle scripts;
# the repo's postinstall needs source files that are not copied yet.
COPY ./frontend/package.json ./frontend/pnpm-lock.yaml ./frontend/pnpm-workspace.yaml ./
RUN --mount=type=cache,target=/pnpm/store CI=1 pnpm fetch
COPY ./frontend/ ./
RUN --mount=type=cache,target=/pnpm/store CI=1 pnpm install --offline
RUN CI=1 pnpm build
FROM golang:1.25-bookworm
@@ -13,9 +26,13 @@ ARG OS="linux"
ARG TARGETARCH
ARG ZEUSURL
# HOME comes from the build user, not the image config; declare it so the
# /root paths below trace to it.
ENV HOME=/root
# This path is important for stacktraces
WORKDIR $GOPATH/src/github.com/signoz/signoz
WORKDIR /root
WORKDIR $HOME
RUN set -eux; \
apt-get update; \
@@ -23,23 +40,36 @@ RUN set -eux; \
g++ \
gcc \
libc6-dev \
make \
pkg-config \
; \
rm -rf /var/lib/apt/lists/*
# Keep the literal cache-mount targets below in sync with these. The caches
# are shared with Dockerfile.integration (same target paths).
ENV GOCACHE=$HOME/.cache/go-build
ENV GOMODCACHE=$GOPATH/pkg/mod
COPY go.mod go.sum ./
RUN go mod download
RUN --mount=type=cache,target=/go/pkg/mod \
go mod download
COPY ./cmd/ ./cmd/
COPY ./ee/ ./ee/
COPY ./pkg/ ./pkg/
COPY ./templates /root/templates
COPY Makefile Makefile
RUN TARGET_DIR=/root ARCHS=${TARGETARCH} ZEUS_URL=${ZEUSURL} LICENSE_URL=${ZEUSURL}/api/v1 make go-build-enterprise-race
RUN mv /root/linux-${TARGETARCH}/signoz /root/signoz
# Invoked directly instead of via make so Makefile changes don't invalidate
# this layer; the Makefile's git-derived ldflags resolve to empty in here
# anyway (.git is dockerignored).
RUN --mount=type=cache,target=/go/pkg/mod \
--mount=type=cache,target=/root/.cache/go-build \
GOARCH=${TARGETARCH} GOOS=${OS} go build -C ./cmd/enterprise -race -tags timetzdata -o /root/signoz \
-ldflags "-s -w \
-X github.com/SigNoz/signoz/pkg/version.version=integration \
-X github.com/SigNoz/signoz/pkg/version.variant=enterprise \
-X github.com/SigNoz/signoz/ee/zeus.url=${ZEUSURL} \
-X github.com/SigNoz/signoz/ee/zeus.deprecatedURL=${ZEUSURL}/api/v1"
COPY --from=build /opt/build ./web/

View File

@@ -7430,6 +7430,8 @@ components:
- below
- equal
- not_equal
- above_or_equal
- below_or_equal
- outside_bounds
type: string
RuletypesCumulativeSchedule:
@@ -15475,6 +15477,72 @@ paths:
summary: Lock dashboard (v2)
tags:
- dashboard
/api/v2/dashboards/{id}/migrate:
post:
deprecated: false
description: 'This endpoint retries the v1→v2 (Perses) migration on a dashboard
still stored in the v1 schema and returns the v2-shape result. It is idempotent:
a dashboard already in the v2 schema is returned unchanged.'
operationId: MigrateDashboardV2
parameters:
- in: path
name: id
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/DashboardtypesGettableDashboardV2'
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- EDITOR
- tokenizer:
- EDITOR
summary: Migrate dashboard to v2
tags:
- dashboard
/api/v2/factor_password/forgot:
post:
deprecated: false

View File

@@ -0,0 +1,123 @@
# PromQL Serving — clickhouseprometheusv2
This document is the subsystem context for `pkg/prometheus/clickhouseprometheusv2`,
the second-generation ClickHouse-backed Prometheus provider. It explains why the
package exists, the correctness constraints that shaped it, and how each fetch
reduction is proven not to change results. Any change to the provider must keep
these invariants; if a change would violate one, it must be flagged and
discussed.
---
## Why a second provider
The v1 provider (`pkg/prometheus/clickhouseprometheus`) serves the promql engine
through the remote-read protobuf adapter: every raw sample of a query's union
window is fetched, serialized, and handed to the engine. The cost is a function
of ingested data, not of the question asked — which is how a dashboard of PromQL
panels can take an instance down.
In v2 the stock promql engine evaluates over a native `storage.Querier`: no
translation layer, per-selector fetch windows, and fetch reductions that are
provably invisible to the engine.
**The core constraint: every reduction either preserves engine semantics exactly
or is not performed.** A PromQL result that differs from upstream Prometheus is
a lost user. The conformance suite
(`tests/integration/tests/promqlconformance/`) replays Prometheus' own test
corpus against both providers and is the arbiter.
---
## Series lookup
Matchers resolve to series once per selector (`selectSeries`) against the series
tables, which hold one row per (fingerprint, bucket) at 1h/6h/1d/1w
granularities. Table selection and window rounding delegate to the shared
metrics schema package (`pkg/telemetryschema/metricstelemetryschema`); the
window start rounds down to the bucket boundary so a window beginning mid-bucket
still matches the bucket's row.
How matchers become SQL is documented at `applySeriesConditions`. The rules that
carry semantics:
- `__name__` matchers (all four types) translate to the `metric_name` column.
- Every other matcher becomes a `JSONExtractString` condition on the labels
column. An equality matcher against `""` matches series *without* the label,
mirroring PromQL, because `JSONExtractString` returns `""` for missing keys.
- Regexes are anchored (`^(?:...)$`) before they reach `match()`: PromQL
matchers match the whole value, ClickHouse `match()` searches for a
substring.
- The series-lookup upper bound is inclusive (`unix_milli <= end`) because the
exporter floors registration rows to the bucket start: a series first
registered in the bucket beginning exactly at `end` would otherwise be
invisible while its samples are in range.
Empty-valued labels come off at this boundary: an empty value means "label
absent" in Prometheus, but stored attribute JSON can carry them.
---
## Sample fetch
Samples are fetched per selector using the engine's per-selector hints, not the
query-wide union window — `foo / foo offset 1d` reads two narrow windows
instead of the widest one twice.
**Last-sample-per-step reduction.** Instant selectors of subquery-free queries
fetch only the last sample per step bucket. The engine resolves an instant
selector at each grid timestamp `t` to the latest sample in the left-open
lookback window `(t lookback, t]`. Buckets anchor at the selector's first
evaluation timestamp — recovered from the hints as
`hints.Start + lookback 1ms`, the inverse of how the engine derives
`hints.Start` — so bucket boundaries coincide with evaluation timestamps, and a
non-final sample of a bucket can never be the latest sample in
`(t lookback, t]` for any grid `t`. Real timestamps are preserved, so the
engine's own lookback and staleness handling stay exact.
Range selectors always fetch raw — every sample feeds the range function. The
subquery-free proof travels in the context as `prometheus.QueryTraits`, because
subquery selectors evaluate at the subquery's step while the hints carry the
top-level step; call sites that do not attach traits get the conservative raw
fetch.
**Row assembly** maps stale flags to the engine's `StaleNaN` and merges series
with identical label sets (`sortAndMerge`) — the engine assumes storages never
emit duplicates. Duplicate timestamps pass through as stored: uniqueness is
ingest's job, and v1 feeds them to the engine as-is over the same data.
**The fingerprint filter is a shard-local semi-join.** The samples query
restricts to the matched series by re-running the series predicates as an
`IN (SELECT fingerprint FROM <local series table> ...)` subquery, not a GLOBAL
broadcast of the matched set. ClickHouse materializes the subquery's set per
shard before the scan, so it still engages the fingerprint primary-key column.
Because the subquery re-executes the predicates after the lookup ran, it can
match series registered in between; sample rows whose fingerprint the lookup
never saw are skipped — the lookup is the read snapshot.
---
## Sharding
`samples_v4` and `time_series_v4` (and all their rollups) shard on the same key
`cityHash64(env, temporality, metric_name, fingerprint)` — so a series'
samples and catalog rows live on the same shard. The semi-join above exploits
that: each shard filters by its own series rows, which are exactly the series
of that shard's samples.
The temporality filter on every samples statement
(`temporality IN ['Cumulative', 'Unspecified']`) is a semantic no-op — the
matched fingerprints already come from those temporalities — that engages the
leading samples primary-key column.
Delta-temporality series stay invisible to PromQL exactly as they are in v1:
the rollout gate is parity with v1, and a Delta stream fed to `rate()`
as-if-cumulative would be wrong, not just new.
---
## Observability
Every statement carries a `log_comment` with
`code.namespace=clickhouse-prometheus-v2` and `code.function.name` naming the
call site, so this provider's work is attributable in `system.query_log`.

View File

@@ -30,11 +30,11 @@ yarn install:browsers # one-time Playwright browser install
### Starting the Test Environment
To spin up the backend stack (SigNoz, ClickHouse, Postgres, Zookeeper, Zeus mock, gateway mock, seeder, migrator-with-web) and keep it running:
To spin up the backend stack (SigNoz, ClickHouse, Postgres, ClickHouse Keeper, Zeus mock, gateway mock, seeder, migrator-with-web) and keep it running:
```bash
cd tests
uv run pytest --basetemp=./tmp/ -vv --reuse --with-web \
uv run pytest --basetemp=./tmp/ -vv --reuse --rebuild --with-web \
e2e/bootstrap/setup.py::test_setup
```
@@ -45,8 +45,13 @@ This command will:
- Start the HTTP seeder container (`tests/seeder/` — exposing `/telemetry/{traces,logs,metrics}` POST + DELETE)
- Write backend coordinates to `tests/e2e/.env.local` (loaded by `playwright.config.ts` via dotenv)
- Keep containers running via the `--reuse` flag
- Rebuild the SigNoz container from the current sources via the `--rebuild` flag
The `--with-web` flag builds the frontend into the SigNoz container — required for E2E. The build takes ~4 mins on a cold start.
The `--with-web` flag builds the frontend into the SigNoz container — required for E2E. The build takes ~4 mins on a cold start; later builds are incremental.
### Rebuilding After Source Changes
The `--with-web` image bakes the built frontend in, so neither backend nor frontend changes are picked up while `--reuse` keeps the container running. `--rebuild` fixes that for both: it kills the SigNoz container, rebuilds the image incrementally (go build cache + pnpm store — a frontend-only change rebuilds in about a minute), and starts a fresh one while databases, mocks, migrations, and the seeder stay reused. The setup command above passes it, so the iteration loop is: change code → re-run the setup command → re-run your specs. `--rebuild` requires `--reuse` and cannot be combined with `--teardown` or `--clean`.
### Stopping the Test Environment
@@ -281,13 +286,16 @@ The full `playwright.config.ts` is the source of truth. Common things to tweak:
The same pytest flags integration tests expose work here, since E2E reuses the shared fixture graph:
- `--reuse` — keep containers warm between runs (required for all iteration).
- `--rebuild` — recreate the SigNoz container from the current sources (backend and, with `--with-web`, frontend) while the rest of the stack stays up. Requires `--reuse`.
- `--teardown` — tear everything down.
- `--clean` — prune the docker build caches, forcing the next image build to start cold.
- `--with-web` — build the frontend into the SigNoz container. **Required for E2E**; integration tests don't need it.
- `--sqlstore-provider`, `--postgres-version`, `--clickhouse-version`, etc. — see `docs/contributing/integration.md`.
- `--sqlstore-provider`, `--postgres-version`, `--clickhouse-version`, etc. — see `docs/contributing/tests/integration.md`.
## What should I remember?
- **Always use the `--reuse` flag** when setting up the E2E stack. `--with-web` adds a ~4 min frontend build; you only want to pay that once.
- **Always use the `--reuse` flag** when setting up the E2E stack. `--with-web` adds a ~4 min frontend build on a cold start; later builds are incremental.
- **Changed backend or frontend code? Re-run the setup command** — it passes `--rebuild`, swapping the SigNoz container for one built from your current sources while the rest of the stack stays up.
- **Don't teardown before setup.** `--reuse` correctly handles partially-set-up state, so chaining teardown → setup wastes time.
- **Prefer UI-driven flows.** Playwright captures BE requests in the trace; a parallel `fetch` probe is almost always redundant. Drop to `page.request.*` only when the UI can't reach what you need.
- **Use `page.waitForResponse` on UI clicks** to assert BE contracts — it still exercises the UI trigger path.

View File

@@ -37,13 +37,34 @@ make py-test-setup
Under the hood this runs, from `tests/`:
```bash
uv run pytest --basetemp=./tmp/ -vv --reuse integration/bootstrap/setup.py::test_setup
uv run pytest --basetemp=./tmp/ -vv --reuse --rebuild --capture=no integration/bootstrap/setup.py::test_setup
```
This command will:
- Start all required services (ClickHouse, PostgreSQL, Zookeeper, SigNoz, Zeus mock, gateway mock)
- Start all required services (ClickHouse, PostgreSQL, ClickHouse Keeper, SigNoz, Zeus mock, gateway mock)
- Register an admin user
- Keep containers running via the `--reuse` flag
- Rebuild the SigNoz container from the current sources via the `--rebuild` flag
### Rebuilding After Source Changes
`--reuse` keeps the running SigNoz container, which means backend source changes are not picked up. `--rebuild` fixes exactly that: it kills the existing SigNoz container, rebuilds the image (incremental — only changed packages recompile thanks to the build cache), and starts a fresh one, while everything else (databases, mocks, migrations) stays reused. `make py-test-setup` passes it by default, so the iteration loop is simply:
```bash
make py-test-setup # (re)build signoz from your current sources
uv run pytest --basetemp=./tmp/ -vv --reuse integration/tests/<suite>/
# ... edit backend code or tests ...
make py-test-setup # pick up the backend changes
uv run pytest --basetemp=./tmp/ -vv --reuse integration/tests/<suite>/
```
The same applies to the e2e stack. `--rebuild` requires `--reuse` and cannot be combined with `--teardown` or `--clean`.
Some suites define their own SigNoz variant in a suite-local `conftest.py` (`create_signoz(..., cache_key=...)` — e.g. `basepath`, `metricreduction`, `querier_json_body`). Those containers are not touched by `make py-test-setup`, which only rebuilds the default instance. For such suites, pass `--rebuild` on the suite run itself — it rebuilds every SigNoz variant the run instantiates:
```bash
uv run pytest --basetemp=./tmp/ -vv --reuse --rebuild integration/tests/<suite>/
```
### Stopping the Test Environment
@@ -56,11 +77,21 @@ make py-test-teardown
Which runs:
```bash
uv run pytest --basetemp=./tmp/ -vv --teardown integration/bootstrap/setup.py::test_teardown
uv run pytest --basetemp=./tmp/ -vv --teardown --capture=no integration/bootstrap/setup.py::test_teardown
```
This destroys the running integration test setup and cleans up resources.
### Cleaning the Image Build Cache
The `signoz:integration` image build keeps its Go build and module caches in BuildKit cache mounts, so rebuilds only recompile what changed. These caches survive `--teardown` (they belong to the Docker builder, not to any container). If a cache ever needs to be nuked — suspected corruption, disk pressure, or to force a genuinely cold build — pass the `--clean` flag:
```bash
uv run pytest --basetemp=./tmp/ -vv --teardown --clean integration/bootstrap/setup.py::test_teardown
```
`--clean` prunes the docker build artifacts backing the incremental image build at session start, so the next build starts from a clean slate. Images and regular layer cache stay intact, but note the pruning is host-wide — it clears build caches for other projects too, not just SigNoz's. The flag composes with any invocation — passing it on a normal `--reuse` run simply makes the next image build start cold (~34 minutes instead of seconds).
## Understanding the Integration Test Framework
Python and pytest form the foundation of the integration testing framework. Testcontainers are used to spin up disposable integration environments. WireMock is used to spin up **test doubles** of external services (Zeus cloud API, gateway, etc.).
@@ -99,7 +130,7 @@ tests/
│ ├── passwordauthn/
│ ├── querier/
│ └── ...
└── e2e/ # Playwright suite (see docs/contributing/e2e.md)
└── e2e/ # Playwright suite (see docs/contributing/tests/e2e.md)
```
Each test suite follows these principles:
@@ -224,9 +255,9 @@ Tests can be configured using pytest options:
- `--sqlstore-provider` — Choose the SQL store provider (default: `postgres`)
- `--sqlite-mode` — SQLite journal mode: `delete` or `wal` (default: `delete`). Only relevant when `--sqlstore-provider=sqlite`.
- `--postgres-version` — PostgreSQL version (default: `15`)
- `--clickhouse-version` — ClickHouse version (default: `25.5.6`)
- `--zookeeper-version` — Zookeeper version (default: `3.7.1`)
- `--schema-migrator-version` — SigNoz schema migrator version (default: `v0.144.2`)
- `--clickhouse-version` — ClickHouse version, also used for ClickHouse Keeper (default: `25.12.5`)
- `--schema-migrator-version` — SigNoz schema migrator version (default: `v0.144.6`)
- `--with-web` — Build the frontend into the SigNoz image (required for e2e)
Example:
@@ -239,6 +270,7 @@ uv run pytest --basetemp=./tmp/ -vv --reuse \
## What should I remember?
- **Always use the `--reuse` flag** when setting up the environment or running tests to keep containers warm. Without it every run rebuilds the stack (~4 mins).
- **Changed backend code? Re-run `make py-test-setup`** — it passes `--rebuild`, swapping the SigNoz container for one built from your current sources while the rest of the stack stays up.
- **Use the `--teardown` flag** only when cleaning up — mixing `--teardown` with `--reuse` is a contradiction.
- **Do not pre-emptively teardown before setup.** If the stack is partially up, `--reuse` picks up from wherever it is. `make py-test-teardown` then `make py-test-setup` wastes minutes.
- **Follow the naming convention** with two-digit numeric prefixes (`01_`, `02_`) for ordered test execution within a suite.
@@ -247,5 +279,5 @@ uv run pytest --basetemp=./tmp/ -vv --reuse \
- **Use descriptive test names** that clearly indicate what is being tested.
- **Leverage fixtures** for common setup. The shared fixture package is at `tests/fixtures/` — reuse before adding new ones.
- **Test both success and failure scenarios** (4xx / 5xx paths) to ensure robust functionality.
- **Run `make py-fmt` and `make py-lint` before committing** Python changes — black + isort + autoflake + pylint.
- **Run `make py-fmt` and `make py-lint` before committing** Python changes — ruff format + ruff check.
- **`--sqlite-mode=wal` does not work on macOS.** The integration test environment runs SigNoz inside a Linux container with the SQLite database file mounted from the macOS host. WAL mode requires shared memory between connections, and connections crossing the VM boundary (macOS host ↔ Linux container) cannot share the WAL index, resulting in `SQLITE_IOERR_SHORT_READ`. WAL mode is tested in CI on Linux only.

View File

@@ -276,6 +276,10 @@ func (module *module) GetV2(ctx context.Context, orgID valuer.UUID, id valuer.UU
return module.pkgDashboardModule.GetV2(ctx, orgID, id)
}
func (module *module) MigrateV2(ctx context.Context, orgID valuer.UUID, id valuer.UUID) (*dashboardtypes.DashboardV2, error) {
return module.pkgDashboardModule.MigrateV2(ctx, orgID, id)
}
func (module *module) UpdateV2(ctx context.Context, orgID valuer.UUID, id valuer.UUID, updatedBy string, updatable dashboardtypes.UpdatableDashboardV2) (*dashboardtypes.DashboardV2, error) {
return module.pkgDashboardModule.UpdateV2(ctx, orgID, id, updatedBy, updatable)
}

View File

@@ -52,6 +52,8 @@ import type {
ListDashboardsV2200,
ListDashboardsV2Params,
LockDashboardV2PathParameters,
MigrateDashboardV2200,
MigrateDashboardV2PathParameters,
PatchDashboardV2200,
PatchDashboardV2PathParameters,
PinDashboardV2PathParameters,
@@ -1804,6 +1806,85 @@ export const useLockDashboardV2 = <
> => {
return useMutation(getLockDashboardV2MutationOptions(options));
};
/**
* This endpoint retries the v1→v2 (Perses) migration on a dashboard still stored in the v1 schema and returns the v2-shape result. It is idempotent: a dashboard already in the v2 schema is returned unchanged.
* @summary Migrate dashboard to v2
*/
export const migrateDashboardV2 = (
{ id }: MigrateDashboardV2PathParameters,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<MigrateDashboardV2200>({
url: `/api/v2/dashboards/${id}/migrate`,
method: 'POST',
signal,
});
};
export const getMigrateDashboardV2MutationOptions = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof migrateDashboardV2>>,
TError,
{ pathParams: MigrateDashboardV2PathParameters },
TContext
>;
}): UseMutationOptions<
Awaited<ReturnType<typeof migrateDashboardV2>>,
TError,
{ pathParams: MigrateDashboardV2PathParameters },
TContext
> => {
const mutationKey = ['migrateDashboardV2'];
const { mutation: mutationOptions } = options
? options.mutation &&
'mutationKey' in options.mutation &&
options.mutation.mutationKey
? options
: { ...options, mutation: { ...options.mutation, mutationKey } }
: { mutation: { mutationKey } };
const mutationFn: MutationFunction<
Awaited<ReturnType<typeof migrateDashboardV2>>,
{ pathParams: MigrateDashboardV2PathParameters }
> = (props) => {
const { pathParams } = props ?? {};
return migrateDashboardV2(pathParams);
};
return { mutationFn, ...mutationOptions };
};
export type MigrateDashboardV2MutationResult = NonNullable<
Awaited<ReturnType<typeof migrateDashboardV2>>
>;
export type MigrateDashboardV2MutationError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Migrate dashboard to v2
*/
export const useMigrateDashboardV2 = <
TError = ErrorType<RenderErrorResponseDTO>,
TContext = unknown,
>(options?: {
mutation?: UseMutationOptions<
Awaited<ReturnType<typeof migrateDashboardV2>>,
TError,
{ pathParams: MigrateDashboardV2PathParameters },
TContext
>;
}): UseMutationResult<
Awaited<ReturnType<typeof migrateDashboardV2>>,
TError,
{ pathParams: MigrateDashboardV2PathParameters },
TContext
> => {
return useMutation(getMigrateDashboardV2MutationOptions(options));
};
/**
* This endpoint returns the sanitized v2-shape dashboard data for public access. Each panel query is reduced to a safe field subset, so filters and raw query strings are not exposed.
* @summary Get public dashboard data (v2)

View File

@@ -8479,6 +8479,8 @@ export enum RuletypesCompareOperatorDTO {
below = 'below',
equal = 'equal',
not_equal = 'not_equal',
above_or_equal = 'above_or_equal',
below_or_equal = 'below_or_equal',
outside_bounds = 'outside_bounds',
}
export interface RuletypesBasicRuleThresholdDTO {
@@ -11162,6 +11164,17 @@ export type UnlockDashboardV2PathParameters = {
export type LockDashboardV2PathParameters = {
id: string;
};
export type MigrateDashboardV2PathParameters = {
id: string;
};
export type MigrateDashboardV2200 = {
data: DashboardtypesGettableDashboardV2DTO;
/**
* @type string
*/
status: string;
};
export type GetFeatures200 = {
/**
* @type array

View File

@@ -1,5 +1,8 @@
// ** Helpers
import { MetrictypesTypeDTO } from 'api/generated/services/sigNoz.schemas';
import {
MetrictypesTemporalityDTO,
MetrictypesTypeDTO,
} from 'api/generated/services/sigNoz.schemas';
import { defaultTraceSelectedColumns } from 'container/OptionsMenu/constants';
import { createIdFromObjectFields } from 'lib/createIdFromObjectFields';
import { createNewBuilderItemName } from 'lib/newQueryBuilder/createNewBuilderItemName';
@@ -389,11 +392,17 @@ const METRIC_TYPE_TO_ATTRIBUTE_TYPE: Record<
export function toAttributeType(
metricType: MetrictypesTypeDTO | undefined,
isMonotonic?: boolean,
temporality?: MetrictypesTemporalityDTO,
): ATTRIBUTE_TYPES | '' {
if (!metricType) {
return '';
}
if (metricType === MetrictypesTypeDTO.sum && isMonotonic === false) {
// Only non-monotonic cumulative sums are treated as gauges; delta sums stay Sum
if (
metricType === MetrictypesTypeDTO.sum &&
isMonotonic === false &&
temporality === MetrictypesTemporalityDTO.cumulative
) {
return ATTRIBUTE_TYPES.GAUGE;
}
return METRIC_TYPE_TO_ATTRIBUTE_TYPE[metricType] || '';

View File

@@ -66,6 +66,10 @@ function ThresholdItem({
return '=';
case AlertThresholdOperator.IS_NOT_EQUAL_TO:
return '!=';
case AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO:
return '>=';
case AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO:
return '<=';
default:
return '';
}

View File

@@ -83,6 +83,10 @@ const getOperatorWord = (op: AlertThresholdOperator): string => {
return 'equal';
case AlertThresholdOperator.IS_NOT_EQUAL_TO:
return 'not equal';
case AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO:
return 'equal or exceed';
case AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO:
return 'equal or fall below';
default:
return 'exceed';
}
@@ -98,6 +102,10 @@ const getThresholdValue = (op: AlertThresholdOperator): number => {
return 100;
case AlertThresholdOperator.IS_NOT_EQUAL_TO:
return 0;
case AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO:
return 80;
case AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO:
return 50;
default:
return 80;
}
@@ -116,6 +124,8 @@ const getDataPoints = (
[AlertThresholdOperator.IS_EQUAL_TO]: [95, 100, 105, 90, 100],
[AlertThresholdOperator.IS_NOT_EQUAL_TO]: [5, 0, 10, 15, 0],
[AlertThresholdOperator.IS_ABOVE]: [75, 85, 90, 78, 95],
[AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO]: [75, 80, 90, 78, 95],
[AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO]: [60, 50, 40, 55, 35],
[AlertThresholdOperator.ABOVE_BELOW]: [75, 85, 90, 78, 95],
},
[AlertThresholdMatchType.ALL_THE_TIME]: {
@@ -123,6 +133,8 @@ const getDataPoints = (
[AlertThresholdOperator.IS_EQUAL_TO]: [100, 100, 100, 100, 100],
[AlertThresholdOperator.IS_NOT_EQUAL_TO]: [5, 10, 15, 8, 12],
[AlertThresholdOperator.IS_ABOVE]: [85, 87, 90, 88, 95],
[AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO]: [80, 87, 90, 88, 95],
[AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO]: [50, 40, 35, 42, 38],
[AlertThresholdOperator.ABOVE_BELOW]: [85, 87, 90, 88, 95],
},
[AlertThresholdMatchType.ON_AVERAGE]: {
@@ -130,6 +142,8 @@ const getDataPoints = (
[AlertThresholdOperator.IS_EQUAL_TO]: [95, 105, 100, 95, 105],
[AlertThresholdOperator.IS_NOT_EQUAL_TO]: [5, 10, 15, 8, 12],
[AlertThresholdOperator.IS_ABOVE]: [75, 85, 90, 78, 95],
[AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO]: [70, 85, 90, 75, 80],
[AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO]: [60, 40, 55, 45, 50],
[AlertThresholdOperator.ABOVE_BELOW]: [75, 85, 90, 78, 95],
},
[AlertThresholdMatchType.IN_TOTAL]: {
@@ -137,6 +151,8 @@ const getDataPoints = (
[AlertThresholdOperator.IS_EQUAL_TO]: [20, 20, 20, 20, 20],
[AlertThresholdOperator.IS_NOT_EQUAL_TO]: [10, 15, 25, 5, 30],
[AlertThresholdOperator.IS_ABOVE]: [10, 15, 25, 5, 30],
[AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO]: [10, 15, 25, 5, 25],
[AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO]: [8, 5, 10, 12, 15],
[AlertThresholdOperator.ABOVE_BELOW]: [10, 15, 25, 5, 30],
},
[AlertThresholdMatchType.LAST]: {
@@ -144,6 +160,8 @@ const getDataPoints = (
[AlertThresholdOperator.IS_EQUAL_TO]: [75, 85, 90, 78, 100],
[AlertThresholdOperator.IS_NOT_EQUAL_TO]: [75, 85, 90, 78, 25],
[AlertThresholdOperator.IS_ABOVE]: [75, 85, 90, 78, 95],
[AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO]: [75, 85, 90, 78, 80],
[AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO]: [75, 85, 90, 78, 50],
[AlertThresholdOperator.ABOVE_BELOW]: [75, 85, 90, 78, 95],
},
};
@@ -157,6 +175,8 @@ const getTooltipOperatorSymbol = (op: AlertThresholdOperator): string => {
[AlertThresholdOperator.IS_BELOW]: '<',
[AlertThresholdOperator.IS_EQUAL_TO]: '=',
[AlertThresholdOperator.IS_NOT_EQUAL_TO]: '!=',
[AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO]: '>=',
[AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO]: '<=',
[AlertThresholdOperator.ABOVE_BELOW]: '>',
};
return symbolMap[op] || '>';
@@ -252,6 +272,10 @@ export const getMatchTypeTooltip = (
return p === thresholdValue;
case AlertThresholdOperator.IS_NOT_EQUAL_TO:
return p !== thresholdValue;
case AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO:
return p >= thresholdValue;
case AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO:
return p <= thresholdValue;
default:
return p > thresholdValue;
}
@@ -294,7 +318,8 @@ export const getMatchTypeTooltip = (
matchType={matchType}
>
Alert triggers (all points {operatorWord} {thresholdValue})<br />
If any point was {thresholdValue}, no alert would fire
If any point didn&apos;t {operatorWord} {thresholdValue}, no alert would
fire
</TooltipExample>
<TooltipLink />
</TooltipContent>

View File

@@ -532,7 +532,7 @@ describe('Footer utils', () => {
['symbol', '>', 'at_least_once'],
['literal', 'above', 'at_least_once'],
['short', 'eq', 'avg'],
['UI-unexposed', 'above_or_equal', 'at_least_once'],
['inclusive', 'above_or_equal', 'at_least_once'],
])(
'round-trips %s op/matchType unchanged through the submit payload (%s / %s)',
(_desc, op, matchType) => {

View File

@@ -332,25 +332,20 @@ describe('CreateAlertV2 utils', () => {
['not_equal', AlertThresholdOperator.IS_NOT_EQUAL_TO],
['not_eq', AlertThresholdOperator.IS_NOT_EQUAL_TO],
['!=', AlertThresholdOperator.IS_NOT_EQUAL_TO],
['5', AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO],
['above_or_equal', AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO],
['above_or_eq', AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO],
['>=', AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO],
['6', AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO],
['below_or_equal', AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO],
['below_or_eq', AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO],
['<=', AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO],
['7', AlertThresholdOperator.ABOVE_BELOW],
['outside_bounds', AlertThresholdOperator.ABOVE_BELOW],
])('maps backend alias %s to canonical enum', (alias, expected) => {
expect(normalizeOperator(alias)).toBe(expected);
});
it.each([
['5', 'above_or_equal'],
['above_or_equal', 'above_or_equal'],
['above_or_eq', 'above_or_equal'],
['>=', 'above_or_equal'],
['6', 'below_or_equal'],
['below_or_equal', 'below_or_equal'],
['below_or_eq', 'below_or_equal'],
['<=', 'below_or_equal'],
])('returns undefined for UI-unexposed alias %s (%s family)', (alias) => {
expect(normalizeOperator(alias)).toBeUndefined();
});
it('returns undefined for unknown values', () => {
expect(normalizeOperator('gibberish')).toBeUndefined();
expect(normalizeOperator(undefined)).toBeUndefined();
@@ -413,8 +408,8 @@ describe('CreateAlertV2 utils', () => {
['symbol', '>', 'at_least_once'],
['short form', 'eq', 'avg'],
['mixed numeric and literal', '7', 'last'],
['UI-unexposed operator', 'above_or_equal', 'at_least_once'],
['UI-unexposed numeric operator', '5', 'at_least_once'],
['inclusive literal operator', 'above_or_equal', 'at_least_once'],
['inclusive numeric operator', '5', 'at_least_once'],
])('preserves %s op/matchType verbatim (%s / %s)', (_desc, op, matchType) => {
const state = getThresholdStateFromAlertDef(buildDef(op, matchType));
expect(state.operator).toBe(op);

View File

@@ -2,9 +2,8 @@ import { AlertThresholdMatchType, AlertThresholdOperator } from './types';
// Mirrors the backend's CompareOperator.Normalize() in
// pkg/types/ruletypes/compare.go. Maps any accepted alias to the enum value
// the dropdown understands. Returns undefined for aliases the UI does not
// expose (e.g. above_or_equal, below_or_equal) so callers can keep the raw
// value on screen instead of silently rewriting it.
// the dropdown understands. Returns undefined for unknown values so callers
// can keep the raw value on screen instead of silently rewriting it.
export function normalizeOperator(
raw: string | undefined,
): AlertThresholdOperator | undefined {
@@ -27,6 +26,16 @@ export function normalizeOperator(
case 'not_eq':
case '!=':
return AlertThresholdOperator.IS_NOT_EQUAL_TO;
case '5':
case 'above_or_equal':
case 'above_or_eq':
case '>=':
return AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO;
case '6':
case 'below_or_equal':
case 'below_or_eq':
case '<=':
return AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO;
case '7':
case 'outside_bounds':
return AlertThresholdOperator.ABOVE_BELOW;

View File

@@ -125,6 +125,14 @@ export const THRESHOLD_OPERATOR_OPTIONS = [
{ value: AlertThresholdOperator.IS_BELOW, label: 'BELOW' },
{ value: AlertThresholdOperator.IS_EQUAL_TO, label: 'EQUAL TO' },
{ value: AlertThresholdOperator.IS_NOT_EQUAL_TO, label: 'NOT EQUAL TO' },
{
value: AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO,
label: 'ABOVE OR EQUAL TO',
},
{
value: AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO,
label: 'BELOW OR EQUAL TO',
},
];
export const ANOMALY_THRESHOLD_OPERATOR_OPTIONS = [

View File

@@ -99,6 +99,8 @@ export enum AlertThresholdOperator {
IS_BELOW = 'below',
IS_EQUAL_TO = 'equal',
IS_NOT_EQUAL_TO = 'not_equal',
IS_ABOVE_OR_EQUAL_TO = 'above_or_equal',
IS_BELOW_OR_EQUAL_TO = 'below_or_equal',
ABOVE_BELOW = 'outside_bounds',
}

View File

@@ -33,6 +33,7 @@ function AllAttributes({
metricName,
metricType,
isMonotonic,
temporality,
minTime,
maxTime,
}: AllAttributesProps): JSX.Element {
@@ -71,6 +72,7 @@ function AllAttributes({
groupBy,
limit,
isMonotonic,
temporality,
);
handleExplorerTabChange(
PANEL_TYPES.TIME_SERIES,
@@ -89,7 +91,7 @@ function AllAttributes({
[MetricsExplorerEventKeys.AttributeKey]: groupBy,
});
},
[metricName, metricType, isMonotonic, handleExplorerTabChange],
[metricName, metricType, isMonotonic, temporality, handleExplorerTabChange],
);
const goToMetricsExploreWithAppliedAttribute = useCallback(
@@ -101,6 +103,7 @@ function AllAttributes({
undefined,
undefined,
isMonotonic,
temporality,
);
handleExplorerTabChange(
PANEL_TYPES.TIME_SERIES,
@@ -120,7 +123,7 @@ function AllAttributes({
[MetricsExplorerEventKeys.AttributeValue]: value,
});
},
[metricName, metricType, isMonotonic, handleExplorerTabChange],
[metricName, metricType, isMonotonic, temporality, handleExplorerTabChange],
);
const handleKeyMenuItemClick = useCallback(

View File

@@ -86,6 +86,7 @@ function MetricDetails({
undefined,
undefined,
metadata?.isMonotonic,
metadata?.temporality,
);
handleExplorerTabChange(
PANEL_TYPES.TIME_SERIES,
@@ -108,6 +109,7 @@ function MetricDetails({
handleExplorerTabChange,
metadata?.type,
metadata?.isMonotonic,
metadata?.temporality,
]);
useEffect(() => {
@@ -196,6 +198,7 @@ function MetricDetails({
metricName={metricName}
metricType={metadata?.type}
isMonotonic={metadata?.isMonotonic}
temporality={metadata?.temporality}
minTime={minTime}
maxTime={maxTime}
/>

View File

@@ -147,6 +147,44 @@ describe('MetricDetails utils', () => {
expect(query.builder.queryData[0]?.spaceAggregation).toBe('sum');
});
it('treats a cumulative non-monotonic Sum as a Gauge', () => {
const query = getMetricDetailsQuery(
TEST_METRIC_NAME,
MetrictypesTypeDTO.sum,
undefined,
undefined,
undefined,
false,
MetrictypesTemporalityDTO.cumulative,
);
expect(query.builder.queryData[0]?.aggregateAttribute?.type).toBe(
ATTRIBUTE_TYPES.GAUGE,
);
expect(query.builder.queryData[0]?.aggregateOperator).toBe('avg');
expect(query.builder.queryData[0]?.timeAggregation).toBe('avg');
expect(query.builder.queryData[0]?.spaceAggregation).toBe('avg');
});
it('treats a delta non-monotonic Sum as a Sum', () => {
const query = getMetricDetailsQuery(
TEST_METRIC_NAME,
MetrictypesTypeDTO.sum,
undefined,
undefined,
undefined,
false,
MetrictypesTemporalityDTO.delta,
);
expect(query.builder.queryData[0]?.aggregateAttribute?.type).toBe(
ATTRIBUTE_TYPES.SUM,
);
expect(query.builder.queryData[0]?.aggregateOperator).toBe('rate');
expect(query.builder.queryData[0]?.timeAggregation).toBe('rate');
expect(query.builder.queryData[0]?.spaceAggregation).toBe('sum');
});
it('should create correct query for GAUGE metric type', () => {
const query = getMetricDetailsQuery(
TEST_METRIC_NAME,

View File

@@ -35,6 +35,7 @@ export interface AllAttributesProps {
metricName: string;
metricType: MetrictypesTypeDTO | undefined;
isMonotonic?: boolean;
temporality?: MetrictypesTemporalityDTO;
minTime?: number;
maxTime?: number;
}

View File

@@ -89,12 +89,16 @@ export function getMetricDetailsQuery(
groupBy?: string,
limit?: number,
isMonotonic?: boolean,
temporality?: MetrictypesTemporalityDTO,
): Query {
let timeAggregation;
let spaceAggregation;
let aggregateOperator;
// Only non-monotonic cumulative sums are treated as gauges; delta sums stay Sum
const isNonMonotonicSum =
metricType === MetrictypesTypeDTO.sum && isMonotonic === false;
metricType === MetrictypesTypeDTO.sum &&
isMonotonic === false &&
temporality === MetrictypesTemporalityDTO.cumulative;
switch (metricType) {
case MetrictypesTypeDTO.sum:
@@ -131,7 +135,7 @@ export function getMetricDetailsQuery(
break;
}
const attributeType = toAttributeType(metricType, isMonotonic);
const attributeType = toAttributeType(metricType, isMonotonic, temporality);
return {
...initialQueriesMap[DataSource.METRICS],

View File

@@ -0,0 +1,40 @@
import { processContextLinks } from '../utils';
describe('processContextLinks', () => {
// Regression for #11325.
it('substitutes per-row groupBy values in the path and query params', () => {
const [link] = processContextLinks(
[
{
id: '1',
label: 'Open trace {{trace_id}}',
url: '/trace/{{trace_id}}?spanId={{span_id}}&levelUp=0&levelDown=0',
},
],
{ _trace_id: 'abc123', _span_id: 'def456' },
);
expect(link.url).toBe('/trace/abc123?spanId=def456&levelUp=0&levelDown=0');
expect(link.label).toBe('Open trace abc123');
});
it('resolves dashboard and global variables alongside row fields', () => {
const [link] = processContextLinks(
[
{ id: '1', label: 'Logs', url: '/logs/{{service}}?ts={{timestamp_start}}' },
],
{ service: 'frontend', timestamp_start: '1720512000000', _service: 'redis' },
);
expect(link.url).toBe('/logs/frontend?ts=1720512000000');
});
it('leaves an unresolvable token in place', () => {
const [link] = processContextLinks(
[{ id: '1', label: 'Open', url: '/trace/{{trace_id}}' }],
{},
);
expect(link.url).toBe('/trace/{{trace_id}}');
});
});

View File

@@ -393,12 +393,13 @@ describe('selecting a metric type updates the aggregation options', () => {
]);
});
it('non-monotonic Sum metric is treated as Gauge', () => {
it('cumulative non-monotonic Sum metric is treated as Gauge', () => {
returnMetrics([
makeMetric({
metricName: 'active_connections',
type: MetrictypesTypeDTO.sum,
isMonotonic: false,
temporality: 'cumulative' as never,
}),
]);
@@ -427,6 +428,36 @@ describe('selecting a metric type updates the aggregation options', () => {
]);
});
it('delta non-monotonic Sum metric is treated as Sum', () => {
returnMetrics([
makeMetric({
metricName: 'queue_depth_delta',
type: MetrictypesTypeDTO.sum,
isMonotonic: false,
temporality: 'delta' as never,
}),
]);
render(<MetricQueryHarness query={makeQuery()} />);
const input = screen.getByRole('combobox');
fireEvent.change(input, {
target: { value: 'queue_depth_delta' },
});
fireEvent.blur(input);
expect(getOptionLabels('time-agg-options')).toStrictEqual([
'Rate',
'Increase',
]);
expect(getOptionLabels('space-agg-options')).toStrictEqual([
'Sum',
'Avg',
'Min',
'Max',
]);
});
it('Histogram metric shows no time options and P50P99 space options', () => {
returnMetrics([
makeMetric({

View File

@@ -34,7 +34,7 @@ export type MetricNameSelectorProps = {
function getAttributeType(
metric: MetricsexplorertypesListMetricDTO,
): ATTRIBUTE_TYPES | '' {
return toAttributeType(metric.type, metric.isMonotonic);
return toAttributeType(metric.type, metric.isMonotonic, metric.temporality);
}
function createAutocompleteData(

View File

@@ -0,0 +1,61 @@
import { resolveTexts } from '../useContextVariables';
const ROW_VARIABLES = { _trace_id: 'abc123', _span_id: 'def456' };
const resolveOne = (
text: string,
processedVariables: Record<string, string>,
): string => resolveTexts({ texts: [text], processedVariables }).fullTexts[0];
describe('resolveTexts', () => {
it('resolves bare {{field}} placeholders from per-row field variables', () => {
expect(
resolveOne('/trace/{{trace_id}}?spanId={{span_id}}', ROW_VARIABLES),
).toBe('/trace/abc123?spanId=def456');
});
it('still resolves the explicitly prefixed {{_field}} form', () => {
expect(resolveOne('/trace/{{_trace_id}}', ROW_VARIABLES)).toBe(
'/trace/abc123',
);
});
it.each([
['{{.trace_id}}', '/trace/{{.trace_id}}'],
['[[trace_id]]', '/trace/[[trace_id]]'],
['$trace_id', '/trace/$trace_id'],
])('resolves the %s placeholder syntax', (_syntax, text) => {
expect(resolveOne(text, ROW_VARIABLES)).toBe('/trace/abc123');
});
it('resolves a field name containing dots', () => {
expect(
resolveOne('/svc/{{service.name}}', { '_service.name': 'redis' }),
).toBe('/svc/redis');
});
it('gives a dashboard variable precedence over a same-named row field', () => {
expect(
resolveOne('/svc/{{service}}', {
service: 'from-dashboard',
_service: 'from-row',
}),
).toBe('/svc/from-dashboard');
});
it('leaves an unknown placeholder untouched', () => {
expect(resolveOne('/trace/{{unknown}}', ROW_VARIABLES)).toBe(
'/trace/{{unknown}}',
);
});
it('picks the truncated side of a multi-value field for truncatedTexts', () => {
const { fullTexts, truncatedTexts } = resolveTexts({
texts: ['services: {{service}}'],
processedVariables: { _service: 'a, b +1-|-a, b, c' },
});
expect(fullTexts[0]).toBe('services: a, b, c');
expect(truncatedTexts[0]).toBe('services: a, b +1');
});
});

View File

@@ -223,6 +223,14 @@ const extractVarName = (
return match;
};
// Per-row fields are registered `_`-prefixed, but templates use the bare name (`{{trace_id}}`).
// Exact match first so dashboard/global variables keep precedence over a same-named row field.
const lookupVariableValue = (
varName: string,
processedVariables: Record<string, string>,
): string | undefined =>
processedVariables[varName] ?? processedVariables[`_${varName}`];
// Utility function to resolve text with processed variables
const resolveText = (
text: string,
@@ -233,7 +241,7 @@ const resolveText = (
return text.replace(combinedPattern, (match) => {
const varName = extractVarName(match, matcher, processedVariables);
const value = processedVariables[varName];
const value = lookupVariableValue(varName, processedVariables);
if (value != null) {
const parts = value.split('-|-');
@@ -254,7 +262,7 @@ const resolveTextWithTruncation = (
const result = text.replace(combinedPattern, (match) => {
const varName = extractVarName(match, matcher, processedVariables);
const value = processedVariables[varName];
const value = lookupVariableValue(varName, processedVariables);
if (value != null) {
const parts = value.split('-|-');

View File

@@ -167,9 +167,9 @@ describe('deriveAlertPrefill', () => {
it.each([
['above', AlertThresholdOperator.IS_ABOVE],
['above_or_equal', AlertThresholdOperator.IS_ABOVE],
['above_or_equal', AlertThresholdOperator.IS_ABOVE_OR_EQUAL_TO],
['below', AlertThresholdOperator.IS_BELOW],
['below_or_equal', AlertThresholdOperator.IS_BELOW],
['below_or_equal', AlertThresholdOperator.IS_BELOW_OR_EQUAL_TO],
['equal', AlertThresholdOperator.IS_EQUAL_TO],
['not_equal', AlertThresholdOperator.IS_NOT_EQUAL_TO],
])('maps panel operator %s → %s', (op, expected) => {

View File

@@ -104,20 +104,6 @@ function pickHighestDanger(
)[0];
}
// The alert UI has no inclusive operator; collapse "or equal" onto its strict variant.
function panelOperatorToAlertOperator(
operator: DashboardtypesComparisonOperatorDTO | undefined,
): AlertThresholdOperator | undefined {
switch (operator) {
case 'above_or_equal':
return normalizeOperator('above');
case 'below_or_equal':
return normalizeOperator('below');
default:
return normalizeOperator(operator);
}
}
export function deriveAlertPrefill(
panel: DashboardtypesPanelDTO,
query: Query,
@@ -135,7 +121,7 @@ export function deriveAlertPrefill(
const top = pickHighestDanger(readPanelThresholds(panel.spec.plugin));
if (top) {
prefill.operator = panelOperatorToAlertOperator(top.operator);
prefill.operator = normalizeOperator(top.operator);
prefill.threshold = {
id: uuid(),
label: 'critical',

2
go.mod
View File

@@ -4,7 +4,7 @@ go 1.25.7
require (
dario.cat/mergo v1.0.2
github.com/AfterShip/clickhouse-sql-parser v0.5.3
github.com/AfterShip/clickhouse-sql-parser v0.5.4
github.com/ClickHouse/clickhouse-go/v2 v2.44.0
github.com/DATA-DOG/go-sqlmock v1.5.2
github.com/SigNoz/clickhouse-go-mock v0.14.0

4
go.sum
View File

@@ -66,8 +66,8 @@ dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA=
dmitri.shuralyov.com/gpu/mtl v0.0.0-20190408044501-666a987793e9/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU=
filippo.io/edwards25519 v1.2.0 h1:crnVqOiS4jqYleHd9vaKZ+HKtHfllngJIiOpNpoJsjo=
filippo.io/edwards25519 v1.2.0/go.mod h1:xzAOLCNug/yB62zG1bQ8uziwrIqIuxhctzJT18Q77mc=
github.com/AfterShip/clickhouse-sql-parser v0.5.3 h1:6iap8XGjuSjD3w7r1UNrg66ljBugcv2P39s4eo/ZLRw=
github.com/AfterShip/clickhouse-sql-parser v0.5.3/go.mod h1:Qi3qvPTfZb/aFwI5V4WFOahgjsLJa4MzVijIAfwOhDw=
github.com/AfterShip/clickhouse-sql-parser v0.5.4 h1:yiCQaMq8EO+dpKdnpP9YYd/ne6MSuOXgsMsNL33NiTI=
github.com/AfterShip/clickhouse-sql-parser v0.5.4/go.mod h1:Qi3qvPTfZb/aFwI5V4WFOahgjsLJa4MzVijIAfwOhDw=
github.com/Azure/azure-sdk-for-go v68.0.0+incompatible h1:fcYLmCpyNYRnvJbPerq7U0hS+6+I79yEDJBqVNcqUzU=
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.21.0 h1:fou+2+WFTib47nS+nz/ozhEBnvU96bKHy6LjRsY4E28=
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.21.0/go.mod h1:t76Ruy8AHvUAC8GfMWJMa0ElSbuIcO03NLpynfbgsPA=

View File

@@ -24,7 +24,7 @@ type fieldPath string
// Slices and interfaces are not surfaced. Pointer fields are dereferenced.
func extractFieldMappings(data any) []fieldPath {
val := reflect.ValueOf(data)
if val.Kind() == reflect.Ptr {
if val.Kind() == reflect.Pointer {
if val.IsNil() {
return nil
}
@@ -60,7 +60,7 @@ func collectFieldMappings(val reflect.Value, prefix string) []fieldPath {
}
ft := field.Type
if ft.Kind() == reflect.Ptr {
if ft.Kind() == reflect.Pointer {
ft = ft.Elem()
}
@@ -73,7 +73,7 @@ func collectFieldMappings(val reflect.Value, prefix string) []fieldPath {
if ft.Kind() == reflect.Struct && ft.String() != "time.Time" {
paths = append(paths, fieldPath(key))
fv := val.Field(i)
if fv.Kind() == reflect.Ptr {
if fv.Kind() == reflect.Pointer {
if fv.IsNil() {
continue
}
@@ -95,7 +95,7 @@ func collectFieldMappings(val reflect.Value, prefix string) []fieldPath {
// flattened OTel-style label keys like "service.name" resolve naturally.
func structRootSet(data any) map[string]bool {
val := reflect.ValueOf(data)
if val.Kind() == reflect.Ptr {
if val.Kind() == reflect.Pointer {
if val.IsNil() {
return nil
}
@@ -121,7 +121,7 @@ func structRootSet(data any) map[string]bool {
continue
}
ft := field.Type
if ft.Kind() == reflect.Ptr {
if ft.Kind() == reflect.Pointer {
ft = ft.Elem()
}
if ft.Kind() == reflect.Struct && ft.String() != "time.Time" {

View File

@@ -85,6 +85,23 @@ func (provider *provider) addDashboardRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v2/dashboards/{id}/migrate", handler.New(provider.authzMiddleware.EditAccess(provider.dashboardHandler.MigrateV2), handler.OpenAPIDef{
ID: "MigrateDashboardV2",
Tags: []string{"dashboard"},
Summary: "Migrate dashboard to v2",
Description: "This endpoint retries the v1→v2 (Perses) migration on a dashboard still stored in the v1 schema and returns the v2-shape result. It is idempotent: a dashboard already in the v2 schema is returned unchanged.",
Request: nil,
RequestContentType: "",
Response: new(dashboardtypes.GettableDashboardV2),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
SecuritySchemes: newSecuritySchemes(types.RoleEditor),
})).Methods(http.MethodPost).GetError(); err != nil {
return err
}
if err := router.Handle("/api/v2/dashboards/{id}", handler.New(provider.authzMiddleware.ViewAccess(provider.dashboardHandler.GetV2), handler.OpenAPIDef{
ID: "GetDashboardV2",
Tags: []string{"dashboard"},

View File

@@ -1,4 +0,0 @@
// Package config provides the configuration management for the Signoz application.
// It includes functionality to define, load, and validate the application's configuration
// using various providers and formats.
package config

View File

@@ -1,3 +0,0 @@
// package error contains error related utilities. Use this package when
// a well-defined error has to be shown.
package errors

View File

@@ -14,6 +14,8 @@ var (
FeatureEnableAIObservability = featuretypes.MustNewName("enable_ai_observability")
FeatureEnableMetricsReduction = featuretypes.MustNewName("enable_metrics_reduction")
FeatureUseInfraMonitoringV2 = featuretypes.MustNewName("use_infra_monitoring_v2")
FeatureUsePrometheusClickhouseV2 = featuretypes.MustNewName("use_prometheus_clickhouse_v2")
)
func MustNewRegistry() featuretypes.Registry {
@@ -106,6 +108,14 @@ func MustNewRegistry() featuretypes.Registry {
DefaultVariant: featuretypes.MustNewName("disabled"),
Variants: featuretypes.NewBooleanVariants(),
},
&featuretypes.Feature{
Name: FeatureUsePrometheusClickhouseV2,
Kind: featuretypes.KindBoolean,
Stage: featuretypes.StageExperimental,
Description: "Runs PromQL queries on the clickhousev2 provider alongside the served engine result and logs any difference; serving is unaffected.",
DefaultVariant: featuretypes.MustNewName("disabled"),
Variants: featuretypes.NewBooleanVariants(),
},
)
if err != nil {
panic(err)

View File

@@ -1,3 +0,0 @@
// package http contains all http related functions such
// as servers, middlewares, routers and renders.
package http

View File

@@ -27,8 +27,12 @@ func TestCache(t *testing.T) {
}
go func() {
require.NoError(t, server.Serve(listener))
_ = server.Serve(listener)
}()
t.Cleanup(func() { _ = server.Close() })
client := &http.Client{Transport: &http.Transport{}}
t.Cleanup(client.CloseIdleConnections)
testCases := []struct {
name string
@@ -45,7 +49,7 @@ func TestCache(t *testing.T) {
req, err := http.NewRequest("GET", "http://"+listener.Addr().String(), nil)
require.NoError(t, err)
res, err := http.DefaultClient.Do(req)
res, err := client.Do(req)
require.NoError(t, err)
defer func() {
require.NoError(t, res.Body.Close())

View File

@@ -34,8 +34,12 @@ func TestTimeout(t *testing.T) {
}
go func() {
require.NoError(t, server.Serve(listener))
_ = server.Serve(listener)
}()
t.Cleanup(func() { _ = server.Close() })
client := &http.Client{Transport: &http.Transport{}}
t.Cleanup(client.CloseIdleConnections)
testCases := []struct {
name string
@@ -70,7 +74,7 @@ func TestTimeout(t *testing.T) {
require.NoError(t, err)
req.Header.Add(headerName, tc.header)
res, err := http.DefaultClient.Do(req)
res, err := client.Do(req)
require.NoError(t, err)
defer func() {
require.NoError(t, res.Body.Close())

View File

@@ -1,2 +0,0 @@
// package server contains an implementation of the http server.
package server

View File

@@ -1,5 +0,0 @@
// Package instrumentation provides utilities for initializing and managing
// OpenTelemetry resources, logging, tracing, and metering within the application. It
// leverages the OpenTelemetry SDK to facilitate the collection and
// export of telemetry data, to an OTLP (OpenTelemetry Protocol) endpoint.
package instrumentation

View File

@@ -5,6 +5,7 @@ import (
"log/slog"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
)
type exception struct{}
@@ -35,9 +36,9 @@ func (h *exception) Wrap(next LogHandler) LogHandler {
t, c, m, _, _, _ := errors.Unwrapb(foundErr)
newRecord.AddAttrs(
slog.String("exception.type", t.String()),
slog.String("exception.code", c.String()),
slog.String("exception.message", m),
slog.String(instrumentationtypes.ExceptionType, t.String()),
slog.String(instrumentationtypes.ExceptionCode, c.String()),
slog.String(instrumentationtypes.ExceptionMessage, m),
)
// Use the stacktrace captured at error creation time if available.
@@ -45,7 +46,7 @@ func (h *exception) Wrap(next LogHandler) LogHandler {
Stacktrace() string
}
if st, ok := foundErr.(stacktracer); ok && st.Stacktrace() != "" {
newRecord.AddAttrs(slog.String("exception.stacktrace", st.Stacktrace()))
newRecord.AddAttrs(slog.String(instrumentationtypes.ExceptionStacktrace, st.Stacktrace()))
}
return next.Handle(ctx, newRecord)

View File

@@ -4,6 +4,8 @@ import (
"context"
"log/slog"
"runtime"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
)
type source struct{}
@@ -17,9 +19,9 @@ func (h *source) Wrap(next LogHandler) LogHandler {
if record.PC != 0 {
frame, _ := runtime.CallersFrames([]uintptr{record.PC}).Next()
record.AddAttrs(
slog.String("code.filepath", frame.File),
slog.String("code.function", frame.Function),
slog.Int("code.lineno", frame.Line),
slog.String(instrumentationtypes.CodeFilePath, frame.File),
slog.String(instrumentationtypes.CodeFunctionName, frame.Function),
slog.Int(instrumentationtypes.CodeLineNumber, frame.Line),
)
}

View File

@@ -63,6 +63,9 @@ type Module interface {
GetV2(ctx context.Context, orgID valuer.UUID, id valuer.UUID) (*dashboardtypes.DashboardV2, error)
// MigrateV2 retries the v1→v2 migration on a dashboard still stored in the v1 schema.
MigrateV2(ctx context.Context, orgID valuer.UUID, id valuer.UUID) (*dashboardtypes.DashboardV2, error)
ListV2(ctx context.Context, orgID valuer.UUID, params *dashboardtypes.ListDashboardsV2Params) (*dashboardtypes.ListableDashboardV2, error)
ListForUserV2(ctx context.Context, orgID valuer.UUID, userID valuer.UUID, params *dashboardtypes.ListDashboardsV2Params) (*dashboardtypes.ListableDashboardForUserV2, error)
@@ -132,6 +135,8 @@ type Handler interface {
GetV2(http.ResponseWriter, *http.Request)
MigrateV2(http.ResponseWriter, *http.Request)
ListV2(http.ResponseWriter, *http.Request)
ListForUserV2(http.ResponseWriter, *http.Request)

View File

@@ -207,6 +207,38 @@ func (handler *handler) GetV2(rw http.ResponseWriter, r *http.Request) {
render.Success(rw, http.StatusOK, dashboard.ToGettableDashboardV2())
}
func (handler *handler) MigrateV2(rw http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
render.Error(rw, err)
return
}
orgID := valuer.MustNewUUID(claims.OrgID)
id := mux.Vars(r)["id"]
if id == "" {
render.Error(rw, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "id is missing in the path"))
return
}
dashboardID, err := valuer.NewUUID(id)
if err != nil {
render.Error(rw, err)
return
}
dashboard, err := handler.module.MigrateV2(ctx, orgID, dashboardID)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, dashboard.ToGettableDashboardV2())
}
func (handler *handler) LockV2(rw http.ResponseWriter, r *http.Request) {
handler.lockUnlockV2(rw, r, true)
}

View File

@@ -4,6 +4,7 @@ import (
"context"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/transition"
"github.com/SigNoz/signoz/pkg/types/coretypes"
"github.com/SigNoz/signoz/pkg/types/dashboardtypes"
"github.com/SigNoz/signoz/pkg/types/tagtypes"
@@ -121,6 +122,51 @@ func (module *module) GetV2(ctx context.Context, orgID valuer.UUID, id valuer.UU
return storable.ToDashboardV2(tags)
}
// MigrateV2 retries the v1→v2 migration on a dashboard still stored as v1 (one the
// bulk 103 migration skipped or failed). Idempotent: an already-v2 one is unchanged.
func (module *module) MigrateV2(ctx context.Context, orgID valuer.UUID, id valuer.UUID) (*dashboardtypes.DashboardV2, error) {
storable, err := module.store.Get(ctx, orgID, id)
if err != nil {
return nil, err
}
// Already migrated: return as-is.
if storable.IsV2() {
tags, err := module.tagModule.ListForResource(ctx, orgID, coretypes.KindDashboard, id)
if err != nil {
return nil, err
}
return storable.ToDashboardV2(tags)
}
// v1→v2 needs v5-shaped queries; run v4→v5 in place first.
transition.NewDashboardMigrateV5(module.settings.Logger(), nil, nil).Migrate(ctx, storable.Data)
v2, err := storable.ConvertV1ToV2()
if err != nil {
return nil, err
}
err = module.store.RunInTx(ctx, func(ctx context.Context) error {
resolvedTags, err := module.tagModule.SyncTags(ctx, orgID, coretypes.KindDashboard, v2.ID, tagtypes.NewPostableTagsFromTags(v2.Tags))
if err != nil {
return err
}
v2.Tags = resolvedTags
storableV2, err := v2.ToStorableDashboard()
if err != nil {
return err
}
return module.store.Update(ctx, orgID, storableV2)
})
if err != nil {
return nil, err
}
return v2, nil
}
func (module *module) UpdateV2(ctx context.Context, orgID valuer.UUID, id valuer.UUID, updatedBy string, updatable dashboardtypes.UpdatableDashboardV2) (*dashboardtypes.DashboardV2, error) {
if err := updatable.Validate(); err != nil {
return nil, err

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildClusterRecords assembles the page records. Node condition counts and
@@ -83,16 +84,36 @@ func buildClusterRecords(
return records
}
func (m *module) getTopClusterGroups(
func (m *module) getTopClusterGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableClusters,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getClustersTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.ClusterNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.ClusterNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.ClusterNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToClustersQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -126,13 +147,20 @@ func (m *module) getTopClusterGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getClustersTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableClusters) (map[string]map[string]string, error) {

View File

@@ -12,6 +12,7 @@ import (
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"golang.org/x/sync/errgroup"
)
// buildContainerRecords assembles the page records, merging kubeletstats
@@ -138,16 +139,36 @@ func buildContainerRecords(
return records
}
func (m *module) getTopContainerGroups(
func (m *module) getTopContainerGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableContainers,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getContainersTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.ContainerNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.ContainerNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.ContainerNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToContainersQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -181,13 +202,20 @@ func (m *module) getTopContainerGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getContainersTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableContainers) (map[string]map[string]string, error) {

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildDaemonSetRecords assembles the page records. Pod status counts come from
@@ -89,16 +90,36 @@ func buildDaemonSetRecords(
return records
}
func (m *module) getTopDaemonSetGroups(
func (m *module) getTopDaemonSetGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableDaemonSets,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getDaemonSetsTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.DaemonSetNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.DaemonSetNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.DaemonSetNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToDaemonSetsQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -132,13 +153,20 @@ func (m *module) getTopDaemonSetGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getDaemonSetsTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableDaemonSets) (map[string]map[string]string, error) {

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildDeploymentRecords assembles the page records. Pod status counts come from
@@ -81,16 +82,36 @@ func buildDeploymentRecords(
return records
}
func (m *module) getTopDeploymentGroups(
func (m *module) getTopDeploymentGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableDeployments,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getDeploymentsTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.DeploymentNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.DeploymentNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.DeploymentNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToDeploymentsQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -124,13 +145,20 @@ func (m *module) getTopDeploymentGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getDeploymentsTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableDeployments) (map[string]map[string]string, error) {

View File

@@ -13,6 +13,7 @@ import (
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"golang.org/x/sync/errgroup"
)
// getPerGroupHostStatusCounts computes the number of active and inactive hosts per group
@@ -248,19 +249,41 @@ func buildHostRecords(
return records
}
// getTopHostGroups runs a ranking query for the ordering metric, sorts the
// results, paginates, and backfills from metadataMap when the page extends
// past the metric-ranked groups.
func (m *module) getTopHostGroups(
// getTopHostGroupsAndMetadata fetches the group metadata and the ordering-metric
// ranking concurrently, then sorts the ranked results, paginates, and backfills
// from metadataMap when the page extends past the metric-ranked groups. Returns
// the page of groups and the metadata map (the caller needs it for Total and
// records). Callers must apply any req.Filter mutation before calling.
func (m *module) getTopHostGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableHosts,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getHostsTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.HostNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.HostNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.HostNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToHostsQueryNames[orderByKey]
// The last entry is the formula/query whose value we sort by.
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -295,13 +318,20 @@ func (m *module) getTopHostGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
// applyHostsActiveStatusFilter MODIFIES req.Filter.Expression to include an IN/NOT IN

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildJobRecords assembles the page records. Pod status counts come from
@@ -89,16 +90,36 @@ func buildJobRecords(
return records
}
func (m *module) getTopJobGroups(
func (m *module) getTopJobGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableJobs,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getJobsTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.JobNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.JobNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.JobNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToJobsQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -132,13 +153,20 @@ func (m *module) getTopJobGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getJobsTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableJobs) (map[string]map[string]string, error) {

View File

@@ -191,18 +191,13 @@ func (m *module) ListHosts(ctx context.Context, orgID valuer.UUID, req *inframon
return resp, nil
}
metadataMap, err := m.getHostsTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopHostGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopHostGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.HostRecord{}
return resp, nil
@@ -291,18 +286,13 @@ func (m *module) ListPods(ctx context.Context, orgID valuer.UUID, req *inframoni
return resp, nil
}
metadataMap, err := m.getPodsTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopPodGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopPodGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.PodRecord{}
return resp, nil
@@ -389,18 +379,13 @@ func (m *module) ListContainers(ctx context.Context, orgID valuer.UUID, req *inf
return resp, nil
}
metadataMap, err := m.getContainersTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopContainerGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopContainerGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.ContainerRecord{}
return resp, nil
@@ -493,18 +478,13 @@ func (m *module) ListNodes(ctx context.Context, orgID valuer.UUID, req *inframon
return resp, nil
}
metadataMap, err := m.getNodesTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopNodeGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopNodeGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.NodeRecord{}
return resp, nil
@@ -591,18 +571,13 @@ func (m *module) ListNamespaces(ctx context.Context, orgID valuer.UUID, req *inf
return resp, nil
}
metadataMap, err := m.getNamespacesTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopNamespaceGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopNamespaceGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.NamespaceRecord{}
return resp, nil
@@ -688,18 +663,13 @@ func (m *module) ListClusters(ctx context.Context, orgID valuer.UUID, req *infra
return resp, nil
}
metadataMap, err := m.getClustersTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopClusterGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopClusterGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.ClusterRecord{}
return resp, nil
@@ -799,18 +769,13 @@ func (m *module) ListVolumes(ctx context.Context, orgID valuer.UUID, req *infram
return resp, nil
}
metadataMap, err := m.getVolumesTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopVolumeGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopVolumeGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.VolumeRecord{}
return resp, nil
@@ -877,18 +842,13 @@ func (m *module) ListDeployments(ctx context.Context, orgID valuer.UUID, req *in
return resp, nil
}
metadataMap, err := m.getDeploymentsTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopDeploymentGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopDeploymentGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.DeploymentRecord{}
return resp, nil
@@ -974,18 +934,13 @@ func (m *module) ListStatefulSets(ctx context.Context, orgID valuer.UUID, req *i
return resp, nil
}
metadataMap, err := m.getStatefulSetsTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopStatefulSetGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopStatefulSetGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.StatefulSetRecord{}
return resp, nil
@@ -1073,18 +1028,13 @@ func (m *module) ListJobs(ctx context.Context, orgID valuer.UUID, req *inframoni
return resp, nil
}
metadataMap, err := m.getJobsTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopJobGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopJobGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.JobRecord{}
return resp, nil
@@ -1172,18 +1122,13 @@ func (m *module) ListDaemonSets(ctx context.Context, orgID valuer.UUID, req *inf
return resp, nil
}
metadataMap, err := m.getDaemonSetsTableMetadata(ctx, orgID, req)
pageGroups, metadataMap, err := m.getTopDaemonSetGroupsAndMetadata(ctx, orgID, req)
if err != nil {
return nil, err
}
resp.Total = len(metadataMap)
pageGroups, err := m.getTopDaemonSetGroups(ctx, orgID, req, metadataMap)
if err != nil {
return nil, err
}
if len(pageGroups) == 0 {
resp.Records = []inframonitoringtypes.DaemonSetRecord{}
return resp, nil

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildNamespaceRecords assembles the page records. Pod status counts come from
@@ -64,16 +65,36 @@ func buildNamespaceRecords(
return records
}
func (m *module) getTopNamespaceGroups(
func (m *module) getTopNamespaceGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableNamespaces,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getNamespacesTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.NamespaceNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.NamespaceNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.NamespaceNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToNamespacesQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -107,13 +128,20 @@ func (m *module) getTopNamespaceGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getNamespacesTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableNamespaces) (map[string]map[string]string, error) {

View File

@@ -12,6 +12,7 @@ import (
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"golang.org/x/sync/errgroup"
)
// buildNodeRecords assembles the page records. Condition counts come from
@@ -91,16 +92,36 @@ func buildNodeRecords(
return records
}
func (m *module) getTopNodeGroups(
func (m *module) getTopNodeGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableNodes,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getNodesTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.NodeNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.NodeNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.NodeNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToNodesQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -134,13 +155,20 @@ func (m *module) getTopNodeGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getNodesTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableNodes) (map[string]map[string]string, error) {

View File

@@ -13,6 +13,7 @@ import (
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"golang.org/x/sync/errgroup"
)
// buildPodRecords assembles the page records. Status counts come from
@@ -145,16 +146,40 @@ func buildPodRecords(
return records
}
func (m *module) getTopPodGroups(
// getTopPodGroupsAndMetadata fetches the group metadata and the ordering-metric
// ranking concurrently, then pages the ranked groups, backfilling from metadata
// when the page extends past the metric-ranked groups. Returns the page of
// groups and the metadata map (needed by the caller for Total and records).
func (m *module) getTopPodGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostablePods,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getPodsTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.PodNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.PodNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.PodNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToPodsQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -188,13 +213,20 @@ func (m *module) getTopPodGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getPodsTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostablePods) (map[string]map[string]string, error) {

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildStatefulSetRecords assembles the page records. Pod status counts come from
@@ -81,16 +82,36 @@ func buildStatefulSetRecords(
return records
}
func (m *module) getTopStatefulSetGroups(
func (m *module) getTopStatefulSetGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableStatefulSets,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getStatefulSetsTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.StatefulSetNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.StatefulSetNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.StatefulSetNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToStatefulSetsQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -124,13 +145,20 @@ func (m *module) getTopStatefulSetGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getStatefulSetsTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableStatefulSets) (map[string]map[string]string, error) {

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/types/inframonitoringtypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/valuer"
"golang.org/x/sync/errgroup"
)
// buildVolumeRecords assembles the page records. VolumeUsage is taken from the
@@ -68,16 +69,36 @@ func buildVolumeRecords(
return records
}
func (m *module) getTopVolumeGroups(
func (m *module) getTopVolumeGroupsAndMetadata(
ctx context.Context,
orgID valuer.UUID,
req *inframonitoringtypes.PostableVolumes,
metadataMap map[string]map[string]string,
) ([]map[string]string, error) {
orderByKey := req.OrderBy.Key.Name
) ([]map[string]string, map[string]map[string]string, error) {
var (
orderByKey string
metadataMap map[string]map[string]string
allMetricGroups []rankedGroup
)
orderByKey = req.OrderBy.Key.Name
g, gCtx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
metadataMap, err = m.getVolumesTableMetadata(gCtx, orgID, req)
return err
})
if orderByKey == inframonitoringtypes.PersistentVolumeClaimNameAttrKey {
return inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.PersistentVolumeClaimNameAttrKey), nil
if err := g.Wait(); err != nil {
return nil, nil, err
}
pageGroups := inframonitoringtypes.PaginateMetadataByName(metadataMap, req.GroupBy, req.OrderBy.Direction, req.Offset, req.Limit, inframonitoringtypes.PersistentVolumeClaimNameAttrKey)
return pageGroups, metadataMap, nil
}
queryNamesForOrderBy := orderByToVolumesQueryNames[orderByKey]
rankingQueryName := queryNamesForOrderBy[len(queryNamesForOrderBy)-1]
@@ -111,13 +132,20 @@ func (m *module) getTopVolumeGroups(
topReq.CompositeQuery.Queries = append(topReq.CompositeQuery.Queries, copied)
}
resp, err := m.querier.QueryRange(ctx, orgID, topReq)
if err != nil {
return nil, err
g.Go(func() error {
resp, err := m.querier.QueryRange(gCtx, orgID, topReq)
if err != nil {
return err
}
allMetricGroups = parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return nil
})
if err := g.Wait(); err != nil {
return nil, nil, err
}
allMetricGroups := parseAndSortGroups(resp, rankingQueryName, req.GroupBy, req.OrderBy.Direction)
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), nil
return paginateWithBackfill(allMetricGroups, metadataMap, req.GroupBy, req.Offset, req.Limit), metadataMap, nil
}
func (m *module) getVolumesTableMetadata(ctx context.Context, orgID valuer.UUID, req *inframonitoringtypes.PostableVolumes) (map[string]map[string]string, error) {

View File

@@ -0,0 +1,85 @@
package clickhouseprometheusv2
import (
"context"
"sync"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/prometheus/prometheus/model/labels"
"github.com/prometheus/prometheus/storage"
"github.com/prometheus/prometheus/util/annotations"
)
// statementRecorder collects the statements a PromQL evaluation would run.
// Safe for concurrent use: the engine may Select selectors concurrently.
type statementRecorder struct {
mu sync.Mutex
statements []prometheus.CapturedStatement
}
func (r *statementRecorder) record(query string, args []any) {
r.mu.Lock()
defer r.mu.Unlock()
r.statements = append(r.statements, prometheus.CapturedStatement{Query: query, Args: args})
}
func (r *statementRecorder) Statements() []prometheus.CapturedStatement {
r.mu.Lock()
defer r.mu.Unlock()
out := make([]prometheus.CapturedStatement, len(r.statements))
copy(out, r.statements)
return out
}
type captureQueryable struct {
client *client
recorder *statementRecorder
}
func (c *captureQueryable) Querier(mint, maxt int64) (storage.Querier, error) {
return &captureQuerier{
querier: querier{mint: mint, maxt: maxt, client: c.client},
recorder: c.recorder,
}, nil
}
// captureQuerier builds the same SQL as the live querier but records it and
// returns no data. The fingerprint filter always takes the subquery form:
// without executing the series lookup, the inline literal set is unknown.
type captureQuerier struct {
querier
recorder *statementRecorder
}
func (c *captureQuerier) Select(ctx context.Context, _ bool, hints *storage.SelectHints, matchers ...*labels.Matcher) storage.SeriesSet {
start, end := c.window(hints)
samplesQuery, args, err := buildSamplesQuery(start, end, metricNamesFromMatchers(matchers), matchers, c.lastSamplePerStepFor(ctx, hints))
if err != nil {
return storage.ErrSeriesSet(err)
}
c.recorder.record(samplesQuery, args)
return storage.EmptySeriesSet()
}
func (c *captureQuerier) LabelValues(context.Context, string, *storage.LabelHints, ...*labels.Matcher) ([]string, annotations.Annotations, error) {
return nil, nil, nil
}
func (c *captureQuerier) LabelNames(context.Context, *storage.LabelHints, ...*labels.Matcher) ([]string, annotations.Annotations, error) {
return nil, nil, nil
}
// metricNamesFromMatchers extracts the statically known metric name, if any.
// The live path derives names from the matched series; the capture path has
// no execution results, so only a __name__ equality contributes.
func metricNamesFromMatchers(matchers []*labels.Matcher) []string {
for _, m := range matchers {
if m.Name == metricNameLabel && m.Type == labels.MatchEqual && m.Value != "" {
return []string{m.Value}
}
}
return nil
}

View File

@@ -0,0 +1,180 @@
package clickhouseprometheusv2
import (
"context"
"encoding/json"
"log/slog"
"math"
"slices"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/prometheus/prometheus/model/labels"
promValue "github.com/prometheus/prometheus/model/value"
)
// seriesLookup holds a series lookup's result: matched fingerprints with
// their labels, and the distinct metric names seen on them.
type seriesLookup struct {
fingerprints map[uint64]labels.Labels
metricNames []string
}
// client executes the series, samples and raw queries against ClickHouse.
type client struct {
settings factory.ScopedProviderSettings
telemetryStore telemetrystore.TelemetryStore
lookbackMs int64
}
func newClient(settings factory.ScopedProviderSettings, telemetryStore telemetrystore.TelemetryStore, cfg prometheus.Config) *client {
lookback := cfg.LookbackDelta
if lookback <= 0 {
// Mirror the engine: promql defaults an unset lookback to 5m.
lookback = defaultLookbackDelta
}
return &client{
settings: settings,
telemetryStore: telemetryStore,
lookbackMs: lookback.Milliseconds(),
}
}
func (c *client) withContext(ctx context.Context, functionName string) context.Context {
return ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalMetrics.StringValue(),
instrumentationtypes.CodeNamespace: "clickhouse-prometheus-v2",
instrumentationtypes.CodeFunctionName: functionName,
})
}
func (c *client) selectSeries(ctx context.Context, query string, args []any) (*seriesLookup, error) {
ctx = c.withContext(ctx, "selectSeries")
rows, err := c.telemetryStore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return nil, err
}
defer rows.Close()
lookup := &seriesLookup{fingerprints: make(map[uint64]labels.Labels)}
names := make(map[string]struct{})
var fingerprint uint64
var labelsJSON string
for rows.Next() {
if err := rows.Scan(&fingerprint, &labelsJSON); err != nil {
return nil, err
}
lset, err := unmarshalLabels(labelsJSON)
if err != nil {
return nil, err
}
lookup.fingerprints[fingerprint] = lset
if name := lset.Get(metricNameLabel); name != "" {
names[name] = struct{}{}
}
}
if err := rows.Err(); err != nil {
return nil, err
}
for name := range names {
lookup.metricNames = append(lookup.metricNames, name)
}
slices.Sort(lookup.metricNames)
return lookup, nil
}
// unmarshalLabels parses the labels JSON column, dropping empty-valued
// labels: empty means "absent" in Prometheus, but stored attribute JSON can
// carry them.
func unmarshalLabels(s string) (labels.Labels, error) {
m := make(map[string]string)
if err := json.Unmarshal([]byte(s), &m); err != nil {
return labels.EmptyLabels(), err
}
builder := labels.NewScratchBuilder(len(m))
for k, v := range m {
if v == "" {
continue
}
builder.Add(k, v)
}
builder.Sort()
return builder.Labels(), nil
}
// selectSamples assembles per-series sample slices from a samples query (raw
// or last-sample-per-step; same column shape), whose rows arrive ordered by
// (fingerprint, unix_milli). Fingerprints missing from the lookup are skipped:
// the samples query's semi-join re-runs the series predicates and can match
// series registered after the lookup ran — the lookup is the read snapshot.
// Stale flags map to the engine's StaleNaN. Duplicate
// timestamps pass through: uniqueness is ingest's job, and v1 feeds them to
// the engine as-is over the same dirty data.
func (c *client) selectSamples(ctx context.Context, query string, args []any, lookup *seriesLookup) ([]*series, error) {
ctx = c.withContext(ctx, "selectSamples")
rows, err := c.telemetryStore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return nil, err
}
defer rows.Close()
var (
result []*series
current *series
fingerprint uint64
prevFp uint64
timestampMs int64
val float64
flags uint32
first = true
haveCurrent bool
staleMarker = math.Float64frombits(promValue.StaleNaN)
unknownCount int
)
for rows.Next() {
if err := rows.Scan(&fingerprint, &timestampMs, &val, &flags); err != nil {
return nil, err
}
if first || fingerprint != prevFp {
first = false
prevFp = fingerprint
lset, ok := lookup.fingerprints[fingerprint]
if !ok {
unknownCount++
haveCurrent = false
continue
}
current = &series{lset: lset}
result = append(result, current)
haveCurrent = true
}
if !haveCurrent {
continue
}
if flags&1 == 1 {
val = staleMarker
}
current.ts = append(current.ts, timestampMs)
current.vs = append(current.vs, val)
}
if err := rows.Err(); err != nil {
return nil, err
}
if unknownCount > 0 {
c.settings.Logger().DebugContext(ctx, "skipped samples of fingerprints missing from series lookup",
slog.Int("unknown_fingerprints", unknownCount))
}
return result, nil
}

View File

@@ -0,0 +1,71 @@
package clickhouseprometheusv2
import (
"context"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/prometheus/prometheus/storage"
)
// provider ties the package together: its own engine and parser, and the
// ClickHouse client behind the native storage.Querier. It stays unexported:
// callers hold the prometheus.Prometheus interface, which is the boundary
// between the two provider implementations.
type provider struct {
settings factory.ScopedProviderSettings
engine *prometheus.Engine
parser prometheus.Parser
client *client
}
var (
_ prometheus.Prometheus = (*provider)(nil)
_ prometheus.StatementCapturer = (*provider)(nil)
)
func NewFactory(telemetryStore telemetrystore.TelemetryStore) factory.ProviderFactory[prometheus.Prometheus, prometheus.Config] {
return factory.NewProviderFactory(factory.MustNewName("clickhousev2"), func(ctx context.Context, providerSettings factory.ProviderSettings, config prometheus.Config) (prometheus.Prometheus, error) {
return New(ctx, providerSettings, config, telemetryStore)
})
}
func New(_ context.Context, providerSettings factory.ProviderSettings, config prometheus.Config, telemetryStore telemetrystore.TelemetryStore) (prometheus.Prometheus, error) {
settings := factory.NewScopedProviderSettings(providerSettings, "github.com/SigNoz/signoz/pkg/prometheus/clickhouseprometheusv2")
engine := prometheus.NewEngine(settings.Logger(), config)
parser := prometheus.NewParser()
client := newClient(settings, telemetryStore, config)
return &provider{
settings: settings,
engine: engine,
parser: parser,
client: client,
}, nil
}
func (p *provider) Engine() *prometheus.Engine {
return p.engine
}
func (p *provider) Parser() prometheus.Parser {
return p.parser
}
func (p *provider) Storage() storage.Queryable {
return p
}
func (p *provider) Querier(mint, maxt int64) (storage.Querier, error) {
return &querier{mint: mint, maxt: maxt, client: p.client}, nil
}
// CapturingStorage implements prometheus.StatementCapturer: a storage that
// records each selector's SQL without executing it, for the preview path.
// A fresh recorder per call keeps concurrent dry-runs isolated.
func (p *provider) CapturingStorage() (storage.Queryable, prometheus.StatementRecorder) {
recorder := &statementRecorder{}
return &captureQueryable{client: p.client, recorder: recorder}, recorder
}

View File

@@ -0,0 +1,171 @@
package clickhouseprometheusv2
import (
"context"
"fmt"
"slices"
"time"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
"github.com/huandu/go-sqlbuilder"
"github.com/prometheus/prometheus/model/labels"
"github.com/prometheus/prometheus/storage"
"github.com/prometheus/prometheus/util/annotations"
)
// defaultLookbackDelta mirrors promql's default when the config leaves the
// lookback unset; the engine and the storage must agree on it for
// last-sample-per-step bucket anchoring.
const defaultLookbackDelta = 5 * time.Minute
// querier is a native storage.Querier over ClickHouse: Select builds SQL
// directly from the matchers and hints, with no remote-read protobuf layer.
type querier struct {
mint, maxt int64
client *client
}
var _ storage.Querier = (*querier)(nil)
func (q *querier) Select(ctx context.Context, sortSeries bool, hints *storage.SelectHints, matchers ...*labels.Matcher) storage.SeriesSet {
start, end := q.window(hints)
seriesQuery, seriesArgs, err := buildSeriesQuery(start, end, matchers)
if err != nil {
return storage.ErrSeriesSet(err)
}
lookup, err := q.client.selectSeries(ctx, seriesQuery, seriesArgs)
if err != nil {
return storage.ErrSeriesSet(err)
}
if len(lookup.fingerprints) == 0 {
return storage.EmptySeriesSet()
}
list, err := q.fetchSamples(ctx, start, end, matchers, lookup, q.lastSamplePerStepFor(ctx, hints))
if err != nil {
return storage.ErrSeriesSet(err)
}
// The engine assumes storages never emit duplicate label sets.
list = sortAndMerge(list)
return newSeriesSet(list)
}
func (q *querier) LabelValues(ctx context.Context, name string, hints *storage.LabelHints, matchers ...*labels.Matcher) ([]string, annotations.Annotations, error) {
sb := sqlbuilder.NewSelectBuilder()
if name == metricNameLabel {
sb.Select("DISTINCT metric_name AS value")
} else {
sb.Select(fmt.Sprintf("DISTINCT JSONExtractString(labels, %s) AS value", sb.Var(name)))
}
adjustedStart, _, table, _ := metricstelemetryschema.WhichTSTableToUse(uint64(q.mint), uint64(q.maxt), false, nil)
sb.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, table))
if err := applySeriesConditions(sb, int64(adjustedStart), q.maxt, matchers); err != nil {
return nil, nil, err
}
sb.Where("value != ''")
if hints != nil && hints.Limit > 0 {
sb.Limit(hints.Limit)
}
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
values, err := q.selectStrings(ctx, "LabelValues", query, args)
if err != nil {
return nil, nil, err
}
slices.Sort(values)
return values, nil, nil
}
func (q *querier) LabelNames(ctx context.Context, hints *storage.LabelHints, matchers ...*labels.Matcher) ([]string, annotations.Annotations, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("DISTINCT arrayJoin(JSONExtractKeys(labels)) AS name")
adjustedStart, _, table, _ := metricstelemetryschema.WhichTSTableToUse(uint64(q.mint), uint64(q.maxt), false, nil)
sb.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, table))
if err := applySeriesConditions(sb, int64(adjustedStart), q.maxt, matchers); err != nil {
return nil, nil, err
}
if hints != nil && hints.Limit > 0 {
sb.Limit(hints.Limit)
}
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
names, err := q.selectStrings(ctx, "LabelNames", query, args)
if err != nil {
return nil, nil, err
}
slices.Sort(names)
return names, nil, nil
}
func (q *querier) Close() error {
return nil
}
// window returns the per-selector fetch window from the hints (already
// adjusted for offset, @, range and lookback); mint/maxt span the union of
// all selectors and are the fallback.
func (q *querier) window(hints *storage.SelectHints) (int64, int64) {
if hints != nil && hints.Start != 0 && hints.End != 0 && hints.Start <= hints.End {
return hints.Start, hints.End
}
return q.mint, q.maxt
}
// lastSamplePerStepFor decides whether the fetch can keep only the last
// sample per step bucket (see lastSamplePerStep). Only instant selectors
// (hints.Range == 0) of subquery-free queries qualify: range selectors need
// every raw sample, and subquery selectors evaluate at the subquery's own
// step while hints carry the top-level step — the subquery-free proof
// arrives as QueryTraits in the context. The first evaluation timestamp is
// recovered as hints.Start + lookback - 1ms, inverting how the engine
// derives hints.Start.
func (q *querier) lastSamplePerStepFor(ctx context.Context, hints *storage.SelectHints) *lastSamplePerStep {
if hints == nil || hints.Range != 0 || hints.Start <= 0 {
return nil
}
traits, ok := prometheus.QueryTraitsFromContext(ctx)
if !ok || !traits.SubqueryFree {
return nil
}
firstEval := hints.Start + q.client.lookbackMs - 1
if firstEval > hints.End {
// Defensive: never anchor a bucket past the window.
firstEval = hints.End
}
return &lastSamplePerStep{firstEvalMs: firstEval, stepMs: hints.Step}
}
// fetchSamples runs the samples query for the matched series (see
// buildSamplesQuery).
func (q *querier) fetchSamples(ctx context.Context, start, end int64, matchers []*labels.Matcher, lookup *seriesLookup, lastPerStep *lastSamplePerStep) ([]*series, error) {
query, args, err := buildSamplesQuery(start, end, lookup.metricNames, matchers, lastPerStep)
if err != nil {
return nil, err
}
return q.client.selectSamples(ctx, query, args, lookup)
}
func (q *querier) selectStrings(ctx context.Context, fn, query string, args []any) ([]string, error) {
ctx = q.client.withContext(ctx, fn)
rows, err := q.client.telemetryStore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return nil, err
}
defer rows.Close()
var out []string
var v string
for rows.Next() {
if err := rows.Scan(&v); err != nil {
return nil, err
}
out = append(out, v)
}
if err := rows.Err(); err != nil {
return nil, err
}
return out, nil
}

View File

@@ -0,0 +1,184 @@
package clickhouseprometheusv2
import (
"sort"
"github.com/prometheus/prometheus/model/histogram"
"github.com/prometheus/prometheus/model/labels"
"github.com/prometheus/prometheus/storage"
"github.com/prometheus/prometheus/tsdb/chunkenc"
"github.com/prometheus/prometheus/util/annotations"
)
// series is one time series with samples as parallel slices ordered by
// timestamp. Deliberately not storage.NewListSeries: that boxes every sample
// as an interface value, a per-sample allocation this fetch path exists to
// avoid.
type series struct {
lset labels.Labels
ts []int64
vs []float64
}
var _ storage.Series = (*series)(nil)
func (s *series) Labels() labels.Labels {
return s.lset
}
func (s *series) Iterator(it chunkenc.Iterator) chunkenc.Iterator {
if fit, ok := it.(*floatIterator); ok {
fit.reset(s)
return fit
}
fit := &floatIterator{}
fit.reset(s)
return fit
}
// floatIterator implements chunkenc.Iterator over a series' sample slices.
type floatIterator struct {
s *series
i int
}
var _ chunkenc.Iterator = (*floatIterator)(nil)
func (it *floatIterator) reset(s *series) {
it.s = s
it.i = -1
}
func (it *floatIterator) Next() chunkenc.ValueType {
it.i++
if it.i >= len(it.s.ts) {
return chunkenc.ValNone
}
return chunkenc.ValFloat
}
func (it *floatIterator) Seek(t int64) chunkenc.ValueType { //nolint:govet // stdmethods flags io.Seeker; this is chunkenc.Iterator's Seek
if it.i < 0 {
it.i = 0
}
if it.i >= len(it.s.ts) {
return chunkenc.ValNone
}
// The current position, once valid, must not move backwards.
if it.s.ts[it.i] >= t {
return chunkenc.ValFloat
}
it.i += sort.Search(len(it.s.ts)-it.i, func(j int) bool {
return it.s.ts[it.i+j] >= t
})
if it.i >= len(it.s.ts) {
return chunkenc.ValNone
}
return chunkenc.ValFloat
}
func (it *floatIterator) At() (int64, float64) {
return it.s.ts[it.i], it.s.vs[it.i]
}
func (it *floatIterator) AtHistogram(*histogram.Histogram) (int64, *histogram.Histogram) {
return 0, nil
}
func (it *floatIterator) AtFloatHistogram(*histogram.FloatHistogram) (int64, *histogram.FloatHistogram) {
return 0, nil
}
func (it *floatIterator) AtT() int64 {
return it.s.ts[it.i]
}
// AtST returns the current start timestamp; not tracked by this storage.
func (it *floatIterator) AtST() int64 {
return 0
}
func (it *floatIterator) Err() error {
return nil
}
// seriesSet iterates a fully materialized, label-sorted list of series.
type seriesSet struct {
series []*series
i int
}
var _ storage.SeriesSet = (*seriesSet)(nil)
func newSeriesSet(list []*series) *seriesSet {
return &seriesSet{series: list, i: -1}
}
func (s *seriesSet) Next() bool {
s.i++
return s.i < len(s.series)
}
func (s *seriesSet) At() storage.Series {
return s.series[s.i]
}
func (s *seriesSet) Err() error {
return nil
}
func (s *seriesSet) Warnings() annotations.Annotations {
return nil
}
// sortAndMerge orders series by label set and merges identical label sets
// by timestamp (first sample wins ties): distinct fingerprints can carry
// identical label sets, and the engine assumes storages never emit
// duplicates.
func sortAndMerge(list []*series) []*series {
if len(list) < 2 {
return list
}
sort.Slice(list, func(i, j int) bool {
return labels.Compare(list[i].lset, list[j].lset) < 0
})
out := list[:1]
for _, s := range list[1:] {
last := out[len(out)-1]
if labels.Compare(last.lset, s.lset) != 0 {
out = append(out, s)
continue
}
merged := mergeSamples(last, s)
out[len(out)-1] = merged
}
return out
}
func mergeSamples(a, b *series) *series {
ts := make([]int64, 0, len(a.ts)+len(b.ts))
vs := make([]float64, 0, len(a.ts)+len(b.ts))
i, j := 0, 0
for i < len(a.ts) && j < len(b.ts) {
switch {
case a.ts[i] < b.ts[j]:
ts = append(ts, a.ts[i])
vs = append(vs, a.vs[i])
i++
case a.ts[i] > b.ts[j]:
ts = append(ts, b.ts[j])
vs = append(vs, b.vs[j])
j++
default:
ts = append(ts, a.ts[i])
vs = append(vs, a.vs[i])
i++
j++
}
}
ts = append(ts, a.ts[i:]...)
vs = append(vs, a.vs[i:]...)
ts = append(ts, b.ts[j:]...)
vs = append(vs, b.vs[j:]...)
return &series{lset: a.lset, ts: ts, vs: vs}
}

View File

@@ -0,0 +1,172 @@
package clickhouseprometheusv2
import (
"fmt"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
"github.com/huandu/go-sqlbuilder"
"github.com/prometheus/prometheus/model/labels"
)
// metricNameLabel is the reserved PromQL label holding the metric name.
const metricNameLabel = "__name__"
// buildSeriesQuery renders the series lookup: one row per matched fingerprint
// with its labels.
func buildSeriesQuery(start, end int64, matchers []*labels.Matcher) (string, []any, error) {
// The series tables hold one row per (fingerprint, bucket) at 1h/6h/1d/1w
// granularities; the schema package picks the table whose bucket fits the
// window and rounds the start down to the bucket boundary, so a window
// beginning mid-bucket still matches the bucket's row.
adjustedStart, _, table, _ := metricstelemetryschema.WhichTSTableToUse(uint64(start), uint64(end), false, nil)
sb := sqlbuilder.NewSelectBuilder()
sb.Select("fingerprint", "any(labels)")
sb.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, table))
if err := applySeriesConditions(sb, int64(adjustedStart), end, matchers); err != nil {
return "", nil, err
}
sb.GroupBy("fingerprint")
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return query, args, nil
}
// buildSamplesQuery renders the samples fetch for the matched series: a
// semi-join re-runs the series predicates against the shard-local series
// table (complete by fingerprint co-locality, see localTimeSeriesTable — a
// GLOBAL broadcast would ship the matched set to every shard, and ClickHouse
// materializes the subquery's set per shard before the scan, so it still
// engages the fingerprint primary-key column). metricNames (observed on the
// matched series when the selector had no __name__ equality) narrows the
// primary-key scan. A non-nil lastPerStep groups to the last sample per
// step bucket.
func buildSamplesQuery(start, end int64, metricNames []string, matchers []*labels.Matcher, lastPerStep *lastSamplePerStep) (string, []any, error) {
sb := sqlbuilder.NewSelectBuilder()
if lastPerStep != nil {
// Aliases must not shadow source columns: ClickHouse resolves aliases
// in WHERE too, and "max(unix_milli) AS unix_milli" would put an
// aggregate into the WHERE clause (error 184).
sb.Select("fingerprint", "max(unix_milli) AS ts", "argMax(value, unix_milli) AS val", "argMax(flags, unix_milli) AS fl")
} else {
sb.Select("fingerprint", "unix_milli", "value", "flags")
}
sb.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, metricstelemetryschema.SamplesV4TableName))
switch len(metricNames) {
case 0:
// No name constraint derivable; the primary-key prefix goes unused.
case 1:
sb.Where(sb.EQ("metric_name", metricNames[0]))
default:
sb.Where(sb.In("metric_name", sqlbuilder.List(metricNames)))
}
// Semantically redundant (the fingerprints already come from these
// temporalities) but engages the leading primary-key column.
sb.Where("temporality IN ['Cumulative', 'Unspecified']")
sub := sqlbuilder.NewSelectBuilder()
sub.Select("fingerprint")
adjustedStart, _, _, localTable := metricstelemetryschema.WhichTSTableToUse(uint64(start), uint64(end), false, nil)
sub.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, localTable))
if err := applySeriesConditions(sub, int64(adjustedStart), end, matchers); err != nil {
return "", nil, err
}
sb.Where(sb.In("fingerprint", sub))
sb.Where(sb.GTE("unix_milli", start), sb.LTE("unix_milli", end))
if lastPerStep != nil {
sb.GroupBy("fingerprint")
if expr := lastPerStep.bucketExpr(); expr != "" {
sb.GroupBy(expr)
}
sb.OrderBy("fingerprint", "ts")
} else {
sb.OrderBy("fingerprint", "unix_milli")
}
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return query, args, nil
}
// applySeriesConditions adds the WHERE conditions of a series table scan:
// __name__ matchers (all four types) translate to the metric_name column,
// every other matcher to a JSONExtractString condition on the labels column.
// An equality matcher against "" matches series without the label, mirroring
// PromQL, because JSONExtractString returns "" for missing keys. Regexes are
// anchored: PromQL matchers match the whole value, ClickHouse match()
// searches for a substring.
func applySeriesConditions(sb *sqlbuilder.SelectBuilder, start, end int64, matchers []*labels.Matcher) error {
for _, m := range matchers {
if m.Name != metricNameLabel {
continue
}
switch m.Type {
case labels.MatchEqual:
sb.Where(sb.EQ("metric_name", m.Value))
case labels.MatchNotEqual:
sb.Where(sb.NE("metric_name", m.Value))
case labels.MatchRegexp:
sb.Where(fmt.Sprintf("match(metric_name, %s)", sb.Var(anchorRegex(m.Value))))
case labels.MatchNotRegexp:
sb.Where(fmt.Sprintf("NOT match(metric_name, %s)", sb.Var(anchorRegex(m.Value))))
default:
return errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported matcher type %q for __name__", m.Type)
}
}
sb.Where("temporality IN ['Cumulative', 'Unspecified']")
// Inclusive upper bound: registration rows are hour-floored (and 6h/1d/1w
// for the rollup tables) by the exporter, so a series first registered in
// the bucket starting exactly at `end` would otherwise be invisible while
// its samples (<= end) are in range.
sb.Where(sb.GTE("unix_milli", start), sb.LTE("unix_milli", end))
for _, m := range matchers {
if m.Name == metricNameLabel {
continue
}
switch m.Type {
case labels.MatchEqual:
sb.Where(fmt.Sprintf("JSONExtractString(labels, %s) = %s", sb.Var(m.Name), sb.Var(m.Value)))
case labels.MatchNotEqual:
sb.Where(fmt.Sprintf("JSONExtractString(labels, %s) != %s", sb.Var(m.Name), sb.Var(m.Value)))
case labels.MatchRegexp:
sb.Where(fmt.Sprintf("match(JSONExtractString(labels, %s), %s)", sb.Var(m.Name), sb.Var(anchorRegex(m.Value))))
case labels.MatchNotRegexp:
sb.Where(fmt.Sprintf("NOT match(JSONExtractString(labels, %s), %s)", sb.Var(m.Name), sb.Var(anchorRegex(m.Value))))
default:
return errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported matcher type %q", m.Type)
}
}
return nil
}
// anchorRegex turns a PromQL regex into its fully-anchored form (see
// applySeriesConditions).
func anchorRegex(v string) string {
return "^(?:" + v + ")$"
}
// lastSamplePerStep reduces an instant-selector fetch to the last sample of
// each step bucket: bucket 0 is (start, firstEval], bucket i is
// (firstEval+(i-1)·step, firstEval+i·step]. Anchoring buckets at the first
// evaluation timestamp makes every bucket boundary an evaluation timestamp,
// so a non-final sample of a bucket can never be the latest sample in any
// (t-lookback, t] the engine resolves — the reduction is lossless. Real
// timestamps are preserved, so the engine's own lookback and staleness
// handling remain exact.
type lastSamplePerStep struct {
firstEvalMs int64
stepMs int64
}
func (t *lastSamplePerStep) bucketExpr() string {
if t.stepMs <= 0 {
// Instant query: a single evaluation at firstEval; one bucket.
return ""
}
return fmt.Sprintf(
"if(unix_milli <= %d, 0, intDiv(unix_milli - %d - 1, %d) + 1)",
t.firstEvalMs, t.firstEvalMs, t.stepMs,
)
}

View File

@@ -0,0 +1,115 @@
package clickhouseprometheusv2
import (
"testing"
"time"
"github.com/prometheus/prometheus/model/labels"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func mustMatcher(t *testing.T, mt labels.MatchType, name, value string) *labels.Matcher {
t.Helper()
m, err := labels.NewMatcher(mt, name, value)
require.NoError(t, err)
return m
}
func TestBuildSeriesQuery(t *testing.T) {
start := int64(1_700_000_000_000)
end := start + time.Hour.Milliseconds()
// The series table window rounds down to the table's bucket boundary.
adjustedStart := start - (start % time.Hour.Milliseconds())
t.Run("equality name and label matchers", func(t *testing.T) {
query, args, err := buildSeriesQuery(start, end, []*labels.Matcher{
mustMatcher(t, labels.MatchEqual, "__name__", "http_requests_total"),
mustMatcher(t, labels.MatchEqual, "job", "api"),
})
require.NoError(t, err)
assert.Equal(t,
"SELECT fingerprint, any(labels) FROM signoz_metrics.distributed_time_series_v4 WHERE metric_name = ? AND temporality IN ['Cumulative', 'Unspecified'] AND unix_milli >= ? AND unix_milli <= ? AND JSONExtractString(labels, ?) = ? GROUP BY fingerprint",
query,
)
assert.Equal(t, []any{"http_requests_total", adjustedStart, end, "job", "api"}, args)
})
t.Run("regex matchers are anchored", func(t *testing.T) {
_, args, err := buildSeriesQuery(start, end, []*labels.Matcher{
mustMatcher(t, labels.MatchEqual, "__name__", "up"),
mustMatcher(t, labels.MatchRegexp, "instance", "prod.*"),
mustMatcher(t, labels.MatchNotRegexp, "env", "dev|test"),
})
require.NoError(t, err)
assert.Equal(t, []any{"up", adjustedStart, end, "instance", "^(?:prod.*)$", "env", "^(?:dev|test)$"}, args)
})
t.Run("regex name matcher uses metric_name column", func(t *testing.T) {
query, args, err := buildSeriesQuery(start, end, []*labels.Matcher{
mustMatcher(t, labels.MatchRegexp, "__name__", "node_cpu.*|node_memory.*"),
})
require.NoError(t, err)
assert.Contains(t, query, "match(metric_name, ?)")
assert.NotContains(t, query, "JSONExtractString")
assert.Equal(t, []any{"^(?:node_cpu.*|node_memory.*)$", adjustedStart, end}, args)
})
t.Run("no name matcher omits metric_name condition", func(t *testing.T) {
query, _, err := buildSeriesQuery(start, end, []*labels.Matcher{
mustMatcher(t, labels.MatchEqual, "job", "api"),
})
require.NoError(t, err)
assert.NotContains(t, query, "metric_name")
})
}
func TestBuildSamplesQuery(t *testing.T) {
start := int64(1_700_000_000_000)
end := start + time.Hour.Milliseconds()
adjustedStart := start - (start % time.Hour.Milliseconds())
matchers := []*labels.Matcher{
mustMatcher(t, labels.MatchEqual, "__name__", "up"),
mustMatcher(t, labels.MatchEqual, "job", "api"),
}
t.Run("raw fetch filters by a shard-local semi-join", func(t *testing.T) {
query, args, err := buildSamplesQuery(start, end, []string{"up"}, matchers, nil)
require.NoError(t, err)
assert.Contains(t, query, "fingerprint IN (SELECT fingerprint FROM signoz_metrics.time_series_v4 WHERE ")
assert.NotContains(t, query, "GLOBAL IN")
assert.Contains(t, query, "ORDER BY fingerprint, unix_milli")
// Args follow placeholder order: samples metric name, the semi-join's
// series predicates, then the samples window bounds.
assert.Equal(t, []any{"up", "up", adjustedStart, end, "job", "api", start, end}, args)
})
t.Run("last-sample-per-step groups by step bucket anchored at first eval", func(t *testing.T) {
lastPerStep := &lastSamplePerStep{firstEvalMs: start + 299_999, stepMs: 60_000}
query, _, err := buildSamplesQuery(start, end, []string{"up"}, matchers, lastPerStep)
require.NoError(t, err)
assert.Contains(t, query, "argMax(value, unix_milli) AS val")
assert.Contains(t, query, "argMax(flags, unix_milli) AS fl")
assert.Contains(t, query, "GROUP BY fingerprint, if(unix_milli <= 1700000299999, 0, intDiv(unix_milli - 1700000299999 - 1, 60000) + 1)")
assert.Contains(t, query, "ORDER BY fingerprint, ts")
// Aliases must not shadow the source columns referenced in WHERE.
assert.NotContains(t, query, "AS unix_milli")
assert.NotContains(t, query, "AS value")
assert.NotContains(t, query, "AS flags")
})
t.Run("instant query keeps one bucket", func(t *testing.T) {
lastPerStep := &lastSamplePerStep{firstEvalMs: end, stepMs: 0}
query, _, err := buildSamplesQuery(start, end, []string{"up"}, matchers, lastPerStep)
require.NoError(t, err)
assert.Contains(t, query, "GROUP BY fingerprint ORDER BY fingerprint, ts")
assert.NotContains(t, query, "intDiv")
})
t.Run("multiple metric names from regex selector", func(t *testing.T) {
query, args, err := buildSamplesQuery(start, end, []string{"node_cpu", "node_memory"}, matchers, nil)
require.NoError(t, err)
assert.Contains(t, query, "metric_name IN (?, ?)")
assert.Equal(t, []any{"node_cpu", "node_memory", "up", adjustedStart, end, "job", "api", start, end}, args)
})
}

View File

@@ -24,6 +24,10 @@ type Config struct {
// Timeout is the maximum time a query is allowed to run before being aborted.
Timeout time.Duration `mapstructure:"timeout"`
// ProviderName selects the storage provider: "clickhouse" (default) or
// "clickhousev2".
ProviderName string `mapstructure:"provider"`
}
func NewConfigFactory() factory.ConfigFactory {
@@ -37,7 +41,8 @@ func newConfig() factory.Config {
Path: "",
MaxConcurrent: 20,
},
Timeout: 2 * time.Minute,
Timeout: 2 * time.Minute,
ProviderName: "clickhouse",
}
}
@@ -45,9 +50,15 @@ func (c Config) Validate() error {
if c.Timeout <= 0 {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "prometheus::timeout must be greater than 0")
}
if c.ProviderName != "" && c.ProviderName != "clickhouse" && c.ProviderName != "clickhousev2" {
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "prometheus::provider must be one of [clickhouse, clickhousev2], got %q", c.ProviderName)
}
return nil
}
func (c Config) Provider() string {
return "clickhouse"
if c.ProviderName == "" {
return "clickhouse"
}
return c.ProviderName
}

View File

@@ -35,3 +35,9 @@ type StatementRecorder interface {
type StatementCapturer interface {
CapturingStorage() (storage.Queryable, StatementRecorder)
}
// ProviderClickhouseV2 is the clickhousev2 provider name: the factory
// registration, the prometheus::provider config value and the
// X-SigNoz-PromQL-Provider request header all use it, so they cannot drift
// apart.
const ProviderClickhouseV2 = "clickhousev2"

49
pkg/prometheus/traits.go Normal file
View File

@@ -0,0 +1,49 @@
package prometheus
import (
"context"
"github.com/prometheus/prometheus/promql/parser"
)
type queryTraitsKey struct{}
// QueryTraits carries per-query facts a storage implementation cannot derive
// from SelectHints alone. Call sites that parse the PromQL expression attach
// traits to the context before handing it to the engine; storages treat a
// missing traits value as "unknown" and stay conservative.
type QueryTraits struct {
// SubqueryFree is true when the query contains no subquery expression.
// Subquery selectors are evaluated at the subquery's own step, but
// SelectHints.Step always carries the top-level step, so step-aligned
// storage optimizations (e.g. keeping only the last sample per step
// bucket) are safe only when this is true.
SubqueryFree bool
}
// DetectQueryTraits derives QueryTraits from a parsed PromQL expression.
func DetectQueryTraits(expr parser.Expr) QueryTraits {
subqueryFree := true
parser.Inspect(expr, func(node parser.Node, _ []parser.Node) error {
if _, ok := node.(*parser.SubqueryExpr); ok {
subqueryFree = false
}
return nil
})
return QueryTraits{SubqueryFree: subqueryFree}
}
// NewContextWithQueryTraits returns a context carrying the given traits.
func NewContextWithQueryTraits(ctx context.Context, traits QueryTraits) context.Context {
return context.WithValue(ctx, queryTraitsKey{}, traits)
}
// QueryTraitsFromContext returns the traits attached to ctx, if any.
//
// Context is used here, unlike for backend selection, because traits must
// cross the promql engine to reach storage.Querier.Select, and the engine's
// interfaces offer no other channel; the alternative is a Prometheus fork.
func QueryTraitsFromContext(ctx context.Context) (QueryTraits, bool) {
traits, ok := ctx.Value(queryTraitsKey{}).(QueryTraits)
return traits, ok
}

View File

@@ -57,6 +57,7 @@ func (handler *handler) QueryRange(rw http.ResponseWriter, req *http.Request) {
render.Error(rw, err)
return
}
queryRangeRequest.PromQLProvider = req.Header.Get("X-SigNoz-PromQL-Provider")
// Validate the query request
if err := queryRangeRequest.Validate(); err != nil {

View File

@@ -5,6 +5,7 @@ import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/statementbuilder"
)
const DefaultMaxConcurrentQueries = 8
@@ -19,6 +20,9 @@ type Config struct {
MaxConcurrentQueries int `yaml:"max_concurrent_queries" mapstructure:"max_concurrent_queries"`
// LogTraceIDWindowPadding is the padding added to narrowed down timerange from trace summary to logs with trace_id filter.
LogTraceIDWindowPadding time.Duration `yaml:"log_trace_id_window_padding" mapstructure:"log_trace_id_window_padding"`
// Keys sit under querier.skip_resource_fingerprint.
statementbuilder.Config `mapstructure:",squash" yaml:",squash"`
}
// NewConfigFactory creates a new config factory for querier.
@@ -29,10 +33,11 @@ func NewConfigFactory() factory.ConfigFactory {
func newConfig() factory.Config {
return Config{
// Default values
CacheTTL: 168 * time.Hour,
FluxInterval: 5 * time.Minute,
CacheTTL: 168 * time.Hour,
FluxInterval: 5 * time.Minute,
MaxConcurrentQueries: DefaultMaxConcurrentQueries,
LogTraceIDWindowPadding: 5 * time.Minute,
Config: statementbuilder.NewConfig(),
}
}
@@ -50,6 +55,10 @@ func (c Config) Validate() error {
if c.LogTraceIDWindowPadding < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "log_trace_id_window_padding must not be negative, got %v", c.LogTraceIDWindowPadding)
}
// Embedded Validate is shadowed by this one; call it explicitly.
if err := c.Config.Validate(); err != nil {
return err
}
return nil
}

View File

@@ -286,7 +286,7 @@ func isNumericKind(t reflect.Type) bool {
if t == nil {
return false
}
for t.Kind() == reflect.Ptr || t.Kind() == reflect.UnsafePointer {
for t.Kind() == reflect.Pointer || t.Kind() == reflect.UnsafePointer {
t = t.Elem()
}
switch t.Kind() {
@@ -367,7 +367,7 @@ func derefValue(v any) any {
val := reflect.ValueOf(v)
for val.Kind() == reflect.Ptr {
for val.Kind() == reflect.Pointer {
if val.IsNil() {
return nil
}

View File

@@ -61,7 +61,7 @@ func getPointerValue(v any) any {
// Use reflection to check if the pointer is nil
rv := reflect.ValueOf(v)
if rv.Kind() == reflect.Ptr && rv.IsNil() {
if rv.Kind() == reflect.Pointer && rv.IsNil() {
return nil
}

View File

@@ -231,7 +231,7 @@ func (q *querier) buildPreviewProviders(
sub.CompositeQuery = qbtypes.CompositeQuery{Queries: []qbtypes.QueryEnvelope{query}}
}
built, _, bErr := q.buildQueries(orgID, &sub, deps, missingMetricQuerySet, event)
built, _, bErr := q.buildQueries(orgID, &sub, deps, missingMetricQuerySet, event, promqlOptions{})
if bErr != nil {
errs[name] = bErr
continue

View File

@@ -8,9 +8,12 @@ import (
"regexp"
"sort"
"strings"
"sync"
"text/template"
"time"
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/prometheus/prometheus/model/labels"
"github.com/prometheus/prometheus/promql"
"github.com/prometheus/prometheus/promql/parser"
@@ -98,6 +101,24 @@ type promqlQuery struct {
tr qbv5.TimeRange
requestType qbv5.RequestType
vars map[string]qbv5.VariableItem
opts promqlOptions
}
// promqlOptions is how a PromQL query relates to the clickhousev2 provider
// (see querier.promqlOptions for where the fields come from and why they are
// flag-gated). Both providers are nil for a plain request, so a plain
// request costs nothing extra.
type promqlOptions struct {
// shadow, when set, runs the query on this provider after serving and
// logs any result difference; the response is never affected.
shadow prometheus.Prometheus
// shadowSlots is the querier-wide admission for shadow runs, shared by
// every query so the bound holds per process.
shadowSlots chan struct{}
// serve, when set, serves the response from this provider instead of the
// default path. Comparison callers fetch the default and the pinned
// result as two API calls and diff them.
serve prometheus.Prometheus
}
var _ qbv5.Query = (*promqlQuery)(nil)
@@ -110,6 +131,7 @@ func newPromqlQuery(
tr qbv5.TimeRange,
requestType qbv5.RequestType,
variables map[string]qbv5.VariableItem,
opts promqlOptions,
) *promqlQuery {
return &promqlQuery{
logger: logger,
@@ -119,10 +141,19 @@ func newPromqlQuery(
tr: tr,
requestType: requestType,
vars: variables,
opts: opts,
}
}
func (q *promqlQuery) Fingerprint() string {
// A pinned request must not share cache entries with default serving: a
// cached default result would satisfy the pin without running the pinned
// provider, and a pinned result would poison normal serving. No
// fingerprint means no caching at all — the pin exists to observe a
// provider, so a cache in front of it defeats the point.
if q.opts.serve != nil {
return ""
}
if q.requestType != qbv5.RequestTypeTimeSeries {
return ""
}
@@ -252,7 +283,16 @@ func (q *promqlQuery) PreviewStatements(ctx context.Context) ([]prometheus.Captu
start := int64(querybuilder.ToNanoSecs(q.tr.From))
end := int64(querybuilder.ToNanoSecs(q.tr.To))
// Attach the same query traits as Execute so the captured statements
// match what the live path would run.
if expr, parseErr := q.parser.ParseExpr(rendered); parseErr == nil {
ctx = prometheus.NewContextWithQueryTraits(ctx, prometheus.DetectQueryTraits(expr))
}
capStorage, recorder := storer.CapturingStorage()
if capStorage == nil {
return nil, nil
}
qry, err := q.promEngine.Engine().NewRangeQuery(
ctx,
capStorage,
@@ -296,6 +336,41 @@ func (q *promqlQuery) Execute(ctx context.Context) (*qbv5.Result, error) {
return nil, err
}
// Attach query traits so the storage can prove step-aligned optimizations
// safe (see prometheus.QueryTraits). A parse failure surfaces below via
// the engine with the enhanced error message.
if expr, parseErr := q.parser.ParseExpr(query); parseErr == nil {
ctx = prometheus.NewContextWithQueryTraits(ctx, prometheus.DetectQueryTraits(expr))
}
// Accumulate ClickHouse-side scan stats across every storage query this
// evaluation issues: progress options propagate to each ClickHouse query
// through the context.
var statsMu sync.Mutex
var rowsScanned, bytesScanned uint64
ctx = clickhouse.Context(ctx, clickhouse.WithProgress(func(p *clickhouse.Progress) {
statsMu.Lock()
rowsScanned += p.Rows
bytesScanned += p.Bytes
statsMu.Unlock()
}))
began := time.Now()
// A pinned provider serves directly from it: comparison callers fetch
// the default result and the pinned result as two API calls and diff
// them.
if q.opts.serve != nil {
matrix, err := q.serveFromProvider(ctx, query, start, end)
if err != nil {
if enhanced := tryEnhancePromQLExecError(err); enhanced != nil {
return nil, enhanced
}
return nil, err
}
return q.toResult(matrix, nil, began, &statsMu, &rowsScanned, &bytesScanned), nil
}
qry, err := q.promEngine.Engine().NewRangeQuery(
ctx,
q.promEngine.Storage(),
@@ -331,6 +406,34 @@ func (q *promqlQuery) Execute(ctx context.Context) (*qbv5.Result, error) {
return nil, errors.WrapInternalf(promErr, errors.CodeInternal, "error getting matrix from promql query %q", query)
}
if q.opts.shadow != nil {
// Shadows detach from the request, so without admission a dashboard
// burst would stack unbounded ClickHouse work for up to the shadow
// timeout — the concurrency pattern behind the original outages.
// Non-blocking: at the cap the comparison is skipped, not queued;
// a sampled shadow stream is exactly as useful for rollout evidence.
select {
case q.opts.shadowSlots <- struct{}{}:
// The engine pools the result's sample slices on Close; the
// shadow comparison needs a stable copy of what was served.
served := copyMatrix(matrix)
servedIn := time.Since(began)
go func() {
defer func() { <-q.opts.shadowSlots }()
q.runShadowCompare(context.WithoutCancel(ctx), query, start, end, served, servedIn)
}()
default:
q.logger.DebugContext(ctx, "promql shadow skipped: at concurrency cap", slog.String("query", query))
}
}
warnings, _ := res.Warnings.AsStrings(query, 10, 0)
return q.toResult(matrix, warnings, began, &statsMu, &rowsScanned, &bytesScanned), nil
}
// toResult converts an evaluated matrix into the v5 result shape, attaching
// the ClickHouse scan stats accumulated during evaluation.
func (q *promqlQuery) toResult(matrix promql.Matrix, warnings []string, began time.Time, statsMu *sync.Mutex, rowsScanned, bytesScanned *uint64) *qbv5.Result {
// Hide only known SigNoz storage keys: label names are user data and may
// legitimately start with "__" (e.g. __address__), so a blanket dunder
// strip mangles user labelsets. The __scope./__resource. prefixes cover
@@ -366,7 +469,13 @@ func (q *promqlQuery) Execute(ctx context.Context) (*qbv5.Result, error) {
series = append(series, &s)
}
warnings, _ := res.Warnings.AsStrings(query, 10, 0)
statsMu.Lock()
stats := qbv5.ExecStats{
RowsScanned: *rowsScanned,
BytesScanned: *bytesScanned,
DurationMS: uint64(time.Since(began).Milliseconds()),
}
statsMu.Unlock()
tsData := &qbv5.TimeSeriesData{
QueryName: q.query.Name,
@@ -400,6 +509,6 @@ func (q *promqlQuery) Execute(ctx context.Context) (*qbv5.Result, error) {
Type: q.requestType,
Value: payload,
Warnings: warnings,
// TODO: map promql stats?
}, nil
Stats: stats,
}
}

View File

@@ -7,6 +7,7 @@ import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/SigNoz/signoz/pkg/prometheus/prometheustest"
qbv5 "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/stretchr/testify/assert"
)
@@ -440,3 +441,15 @@ func TestQuotedMetricOutsideBracesPattern(t *testing.T) {
})
}
}
// A pinned request must not share cache entries with default serving: a
// cached default result would satisfy the pin without running the pinned
// provider.
func TestFingerprint_PinnedProviderBypassesCache(t *testing.T) {
q := &promqlQuery{
logger: slog.Default(),
query: qbv5.PromQuery{Query: "up"},
opts: promqlOptions{serve: &prometheustest.Provider{}},
}
assert.Empty(t, q.Fingerprint())
}

View File

@@ -0,0 +1,175 @@
package querier
import (
"context"
"fmt"
"log/slog"
"math"
"sort"
"time"
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/prometheus/prometheus/model/labels"
"github.com/prometheus/prometheus/promql"
)
// shadowTimeout bounds a shadow evaluation; a shadow run must never outlive
// the request by much or pile up.
const shadowTimeout = 2 * time.Minute
// runShadowCompare executes the query on the clickhousev2 provider exactly
// as it would serve (the engine over the v2 querier), compares against the
// served result and logs the outcome. Serving is never affected: this runs
// after the response, off the request context, and only logs. The mismatch
// and failure logs are the rollout evidence — serving cuts over to v2 only
// after they stay clean.
func (q *promqlQuery) runShadowCompare(ctx context.Context, query string, startNs, endNs int64, served promql.Matrix, servedIn time.Duration) {
defer func() {
if r := recover(); r != nil {
q.logger.ErrorContext(ctx, "promql shadow comparison panicked", slog.Any("panic", r), slog.String("query", query))
}
}()
ctx, cancel := context.WithTimeout(ctx, shadowTimeout)
defer cancel()
// The request context carries the served response's scan-stats progress
// callback; without replacing it the shadow's ClickHouse progress would
// race into the served stats. The response itself was already sent.
ctx = clickhouse.Context(ctx, clickhouse.WithProgress(func(*clickhouse.Progress) {}))
if expr, parseErr := q.parser.ParseExpr(query); parseErr == nil {
ctx = prometheus.NewContextWithQueryTraits(ctx, prometheus.DetectQueryTraits(expr))
}
start, end := time.Unix(0, startNs), time.Unix(0, endNs)
began := time.Now()
shadow, err := executeOnProvider(ctx, q.opts.shadow, query, start, end, q.query.Step.Duration)
shadowIn := time.Since(began)
logAttrs := []any{
slog.String("query", query),
slog.Int64("start_ms", startNs/int64(time.Millisecond)),
slog.Int64("end_ms", endNs/int64(time.Millisecond)),
slog.Duration("step", q.query.Step.Duration),
slog.Duration("served_in", servedIn),
slog.Duration("shadow_in", shadowIn),
}
if err != nil {
// A shadow failure would be a serving failure after rollout; surface
// it at the same level as a result mismatch.
q.logger.WarnContext(ctx, "promql shadow execution failed", append(logAttrs, slog.Any("error", err))...)
return
}
servedNorm := normalizeShadowMatrix(served)
shadowNorm := normalizeShadowMatrix(shadow)
if diff := diffShadowMatrices(servedNorm, shadowNorm); diff != "" {
q.logger.WarnContext(ctx, "promql shadow comparison mismatch", append(logAttrs,
slog.String("diff", diff),
slog.Int("served_series", len(servedNorm)),
slog.Int("shadow_series", len(shadowNorm)),
)...)
return
}
// Matches log the timings: served_in vs shadow_in across the fleet is
// the perf evidence for the cutover, gathered for free.
q.logger.DebugContext(ctx, "promql shadow comparison matched", logAttrs...)
}
// serveFromProvider evaluates the query the way the pinned provider would
// serve it.
func (q *promqlQuery) serveFromProvider(ctx context.Context, query string, startNs, endNs int64) (promql.Matrix, error) {
return executeOnProvider(ctx, q.opts.serve, query, time.Unix(0, startNs), time.Unix(0, endNs), q.query.Step.Duration)
}
// executeOnProvider evaluates the query the way the provider would serve it:
// the engine over the provider's storage. The returned matrix is an owned
// copy.
func executeOnProvider(ctx context.Context, prov prometheus.Prometheus, query string, start, end time.Time, step time.Duration) (promql.Matrix, error) {
qry, err := prov.Engine().NewRangeQuery(ctx, prov.Storage(), nil, query, start, end, step)
if err != nil {
return nil, err
}
defer qry.Close()
res := qry.Exec(ctx)
if res.Err != nil {
return nil, res.Err
}
matrix, err := res.Matrix()
if err != nil {
return nil, err
}
// Close returns the result's sample slices to the engine pool.
return copyMatrix(matrix), nil
}
func copyMatrix(matrix promql.Matrix) promql.Matrix {
out := make(promql.Matrix, 0, len(matrix))
for _, s := range matrix {
floats := make([]promql.FPoint, len(s.Floats))
copy(floats, s.Floats)
out = append(out, promql.Series{Metric: s.Metric.Copy(), Floats: floats})
}
return out
}
// normalizeShadowMatrix sorts by label set for order-independent
// comparison. Both providers now resolve series identity the same way
// (empty-valued labels dropped at read, no synthetic fingerprint label
// since the v1 series-identity fix), so labels need no normalization.
func normalizeShadowMatrix(matrix promql.Matrix) promql.Matrix {
out := make(promql.Matrix, 0, len(matrix))
out = append(out, matrix...)
sort.Slice(out, func(i, j int) bool { return labels.Compare(out[i].Metric, out[j].Metric) < 0 })
return out
}
// diffShadowMatrices returns a description of the first difference, or "".
// Values compare with relative tolerance: spatial aggregations accumulate
// floats in storage order, which differs between the providers in the last
// ULP.
func diffShadowMatrices(served, shadow promql.Matrix) string {
const relTol = 1e-9
if len(served) != len(shadow) {
return fmt.Sprintf("series count: served=%d shadow=%d", len(served), len(shadow))
}
for i := range served {
if labels.Compare(served[i].Metric, shadow[i].Metric) != 0 {
return fmt.Sprintf("series %d labels: served=%s shadow=%s", i, served[i].Metric, shadow[i].Metric)
}
if len(served[i].Floats) != len(shadow[i].Floats) {
return fmt.Sprintf("series %s points: served=%d shadow=%d", served[i].Metric, len(served[i].Floats), len(shadow[i].Floats))
}
for j := range served[i].Floats {
a, b := served[i].Floats[j], shadow[i].Floats[j]
if a.T != b.T {
return fmt.Sprintf("series %s point %d ts: served=%d shadow=%d", served[i].Metric, j, a.T, b.T)
}
// NaN and infinities first: NaN != NaN and Inf-Inf arithmetic
// would otherwise make one-sided NaN and Inf-vs-finite compare
// as equal (NaN > x and Inf > Inf are both false).
if math.IsNaN(a.F) || math.IsNaN(b.F) {
if math.IsNaN(a.F) != math.IsNaN(b.F) {
return fmt.Sprintf("series %s @%d value: served=%v shadow=%v", served[i].Metric, a.T, a.F, b.F)
}
continue
}
if math.IsInf(a.F, 0) || math.IsInf(b.F, 0) {
if a.F != b.F {
return fmt.Sprintf("series %s @%d value: served=%v shadow=%v", served[i].Metric, a.T, a.F, b.F)
}
continue
}
diff := math.Abs(a.F - b.F)
scale := math.Max(math.Abs(a.F), math.Abs(b.F))
if diff > relTol*math.Max(scale, 1e-300) && diff > 1e-12 {
return fmt.Sprintf("series %s @%d value: served=%v shadow=%v", served[i].Metric, a.T, a.F, b.F)
}
}
}
return ""
}

View File

@@ -0,0 +1,67 @@
package querier
import (
"math"
"testing"
"github.com/prometheus/prometheus/model/labels"
"github.com/prometheus/prometheus/promql"
"github.com/stretchr/testify/assert"
)
func TestNormalizeShadowMatrix(t *testing.T) {
matrix := promql.Matrix{
{
Metric: labels.FromStrings("__name__", "up", "job", "api"),
Floats: []promql.FPoint{{T: 1000, F: 1}},
},
{
Metric: labels.FromStrings("a", "1"),
Floats: []promql.FPoint{{T: 1000, F: 2}},
},
}
norm := normalizeShadowMatrix(matrix)
// sorted by label set; labels pass through untouched — both providers
// resolve series identity identically since the v1 series-identity fix
assert.Equal(t, labels.FromStrings("__name__", "up", "job", "api"), norm[0].Metric)
assert.Equal(t, labels.FromStrings("a", "1"), norm[1].Metric)
}
func TestDiffShadowMatrices(t *testing.T) {
series := func(v float64) promql.Matrix {
return promql.Matrix{{Metric: labels.FromStrings("a", "1"), Floats: []promql.FPoint{{T: 1000, F: v}}}}
}
assert.Empty(t, diffShadowMatrices(series(1.5), series(1.5)))
// last-ULP differences from storage-order float accumulation are expected
assert.Empty(t, diffShadowMatrices(series(0.08888888888888889), series(0.08888888888888888)))
assert.Empty(t, diffShadowMatrices(series(math.NaN()), series(math.NaN())))
assert.Contains(t, diffShadowMatrices(series(1.5), series(1.6)), "value")
assert.Contains(t, diffShadowMatrices(series(1.5), promql.Matrix{}), "series count")
assert.Contains(t, diffShadowMatrices(
series(1.5),
promql.Matrix{{Metric: labels.FromStrings("a", "2"), Floats: []promql.FPoint{{T: 1000, F: 1.5}}}},
), "labels")
assert.Contains(t, diffShadowMatrices(
series(1.5),
promql.Matrix{{Metric: labels.FromStrings("a", "1"), Floats: []promql.FPoint{{T: 2000, F: 1.5}}}},
), "ts")
}
// One-sided NaN makes every float comparison false, and Inf-Inf arithmetic
// yields Inf > Inf == false; without explicit handling both divergences log
// as matched — a shadow comparator that cannot see them would green-light a
// broken rollout.
func TestDiffShadowMatrices_SpecialFloats(t *testing.T) {
point := func(v float64) promql.Matrix {
return promql.Matrix{{Metric: labels.FromStrings("a", "1"), Floats: []promql.FPoint{{T: 1000, F: v}}}}
}
assert.NotEmpty(t, diffShadowMatrices(point(math.NaN()), point(1.5)), "one-sided NaN must diff")
assert.NotEmpty(t, diffShadowMatrices(point(1.5), point(math.NaN())), "one-sided NaN must diff either way")
assert.NotEmpty(t, diffShadowMatrices(point(math.Inf(1)), point(1.5)), "Inf vs finite must diff")
assert.NotEmpty(t, diffShadowMatrices(point(math.Inf(1)), point(math.Inf(-1))), "opposite infinities must diff")
assert.Empty(t, diffShadowMatrices(point(math.Inf(1)), point(math.Inf(1))), "equal infinities match")
assert.Empty(t, diffShadowMatrices(point(math.NaN()), point(math.NaN())), "both NaN match")
}

View File

@@ -24,6 +24,7 @@ import (
"github.com/SigNoz/signoz/pkg/statsreporter"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
"github.com/SigNoz/signoz/pkg/types/metrictypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -46,11 +47,19 @@ type Querier interface {
}
type querier struct {
logger *slog.Logger
fl flagger.Flagger
telemetryStore telemetrystore.TelemetryStore
metadataStore telemetrytypes.MetadataStore
promEngine prometheus.Prometheus
logger *slog.Logger
fl flagger.Flagger
telemetryStore telemetrystore.TelemetryStore
metadataStore telemetrytypes.MetadataStore
promEngine prometheus.Prometheus
// promV2 is the clickhousev2 prometheus provider, wired only when the
// serving provider is the default one (nil otherwise). It reads the same
// ClickHouse data through a different implementation; PromQL queries
// shadow-compare against it behind the use_prometheus_clickhouse_v2 flag
// and can be pinned to it for a response (see promqlOptions). It never
// serves by default — that cutover happens only after the shadow logs
// stay clean.
promV2 prometheus.Prometheus
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation]
logStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation]
auditStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation]
@@ -61,8 +70,16 @@ type querier struct {
liveDataRefresh time.Duration
builderConfig builderConfig
maxConcurrentQueries int
// shadowSlots bounds concurrent shadow comparisons per process; shadows
// detach from their requests, so nothing else limits how many pile up.
shadowSlots chan struct{}
}
// maxConcurrentShadows is deliberately small: a shadow is a full extra
// ClickHouse evaluation, and a sampled stream of comparisons is exactly as
// useful for rollout evidence as an exhaustive one under load.
const maxConcurrentShadows = 8
var _ Querier = (*querier)(nil)
func New(
@@ -70,6 +87,7 @@ func New(
telemetryStore telemetrystore.TelemetryStore,
metadataStore telemetrytypes.MetadataStore,
promEngine prometheus.Prometheus,
promV2 prometheus.Prometheus,
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation],
logStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation],
auditStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation],
@@ -91,6 +109,7 @@ func New(
telemetryStore: telemetryStore,
metadataStore: metadataStore,
promEngine: promEngine,
promV2: promV2,
traceStmtBuilder: traceStmtBuilder,
logStmtBuilder: logStmtBuilder,
auditStmtBuilder: auditStmtBuilder,
@@ -103,6 +122,7 @@ func New(
logTraceIDWindowPaddingMS: uint64(logTraceIDWindowPadding.Milliseconds()),
},
maxConcurrentQueries: maxConcurrentQueries,
shadowSlots: make(chan struct{}, maxConcurrentShadows),
}
}
@@ -142,7 +162,11 @@ func (q *querier) QueryRange(ctx context.Context, orgID valuer.UUID, req *qbtype
missingMetricQuerySet[name] = true
}
queries, steps, err := q.buildQueries(orgID, req, dependencyQueries, missingMetricQuerySet, event)
promqlOpts, err := q.promqlOptions(ctx, orgID, req)
if err != nil {
return nil, err
}
queries, steps, err := q.buildQueries(orgID, req, dependencyQueries, missingMetricQuerySet, event, promqlOpts)
if err != nil {
return nil, err
}
@@ -185,12 +209,41 @@ func (q *querier) QueryRange(ctx context.Context, orgID valuer.UUID, req *qbtype
return qbResp, qbErr
}
// promqlOptions derives the PromQL execution options for a request. With the
// org's use_prometheus_clickhouse_v2 flag on, queries are shadow-compared
// against the clickhousev2 provider (serving unaffected, diffs logged; see
// promql_shadow.go). The X-SigNoz-PromQL-Provider header may instead pin the
// response to that provider — integration tests and support fetch both
// results for comparison — so it is deliberately flag-gated too: without the
// gate the header would be an unaudited switch onto a provider still under
// validation.
func (q *querier) promqlOptions(ctx context.Context, orgID valuer.UUID, req *qbtypes.QueryRangeRequest) (promqlOptions, error) {
enabled := q.fl.BooleanOrEmpty(ctx, flagger.FeatureUsePrometheusClickhouseV2, featuretypes.NewFlaggerEvaluationContext(orgID))
if req.PromQLProvider == "" {
if enabled && q.promV2 != nil {
return promqlOptions{shadow: q.promV2, shadowSlots: q.shadowSlots}, nil
}
return promqlOptions{}, nil
}
if req.PromQLProvider != prometheus.ProviderClickhouseV2 {
return promqlOptions{}, errors.NewInvalidInputf(errors.CodeInvalidInput, "unknown promql provider %q", req.PromQLProvider)
}
if !enabled {
return promqlOptions{}, errors.NewInvalidInputf(errors.CodeInvalidInput, "promql provider %q requires the use_prometheus_clickhouse_v2 flag", req.PromQLProvider)
}
if q.promV2 == nil {
return promqlOptions{}, errors.NewInvalidInputf(errors.CodeInvalidInput, "promql provider %q is not available", req.PromQLProvider)
}
return promqlOptions{serve: q.promV2}, nil
}
func (q *querier) buildQueries(
orgID valuer.UUID,
req *qbtypes.QueryRangeRequest,
dependencyQueries map[string]bool,
missingMetricQuerySet map[string]bool,
event *qbtypes.QBEvent,
promqlOpts promqlOptions,
) (map[string]qbtypes.Query, map[string]qbtypes.Step, error) {
tmplVars := req.Variables
@@ -215,7 +268,7 @@ func (q *querier) buildQueries(
if !ok {
return nil, nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "invalid promql query spec %T", query.Spec)
}
promqlQuery := newPromqlQuery(q.logger, q.promEngine, promQuery, qbtypes.TimeRange{From: req.Start, To: req.End}, req.RequestType, tmplVars)
promqlQuery := newPromqlQuery(q.logger, q.promEngine, promQuery, qbtypes.TimeRange{From: req.Start, To: req.End}, req.RequestType, tmplVars, promqlOpts)
queries[promQuery.Name] = promqlQuery
steps[promQuery.Name] = promQuery.Step
case qbtypes.QueryTypeClickHouseSQL:
@@ -858,7 +911,7 @@ func (q *querier) createRangedQuery(_ valuer.UUID, originalQuery qbtypes.Query,
switch qt := originalQuery.(type) {
case *promqlQuery:
queryCopy := qt.query.Copy()
return newPromqlQuery(q.logger, q.promEngine, queryCopy, timeRange, qt.requestType, qt.vars)
return newPromqlQuery(q.logger, qt.promEngine, queryCopy, timeRange, qt.requestType, qt.vars, qt.opts)
case *chSQLQuery:
queryCopy := qt.query.Copy()

View File

@@ -48,6 +48,7 @@ func TestQueryRange_MetricTypeMissing(t *testing.T) {
nil, // telemetryStore
metadataStore,
nil, // prometheus
nil, // promV2
nil, // traceStmtBuilder
nil, // logStmtBuilder
nil, // auditStmtBuilder
@@ -120,6 +121,7 @@ func TestQueryRange_MetricTypeFromStore(t *testing.T) {
telemetryStore,
metadataStore,
nil, // prometheus
nil, // promV2
nil, // traceStmtBuilder
nil, // logStmtBuilder
nil, // auditStmtBuilder

View File

@@ -18,6 +18,7 @@ import (
func NewFactory(
telemetryStore telemetrystore.TelemetryStore,
prometheus prometheus.Prometheus,
promV2 prometheus.Prometheus,
metadataStore telemetrytypes.MetadataStore,
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation],
logStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation],
@@ -40,6 +41,7 @@ func NewFactory(
telemetryStore,
metadataStore,
prometheus,
promV2,
traceStmtBuilder,
logStmtBuilder,
auditStmtBuilder,

View File

@@ -128,7 +128,7 @@ func NewTestManager(t *testing.T, testOpts *TestManagerOptions) *Manager {
meterStmtBuilder, err := meterstatementbuilder.NewFactory(metadataStore, flagger).New(ctx, providerSettings, cfg)
require.NoError(t, err)
bucketCache := querier.NewBucketCache(providerSettings, cache, 0, 0)
providerFactory := signozquerier.NewFactory(telemetryStore, prometheus, metadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, flagger)
providerFactory := signozquerier.NewFactory(telemetryStore, prometheus, nil, metadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, flagger)
mockQuerier, err := providerFactory.New(context.Background(), providerSettings, querier.Config{})
require.NoError(t, err)

View File

@@ -40,6 +40,7 @@ func prepareQuerierForMetrics(t *testing.T, telemetryStore telemetrystore.Teleme
telemetryStore,
metadataStore,
nil, // prometheus
nil, // promV2
nil, // traceStmtBuilder
nil, // logStmtBuilder
nil, // auditStmtBuilder
@@ -74,6 +75,7 @@ func prepareQuerierForLogs(t *testing.T, telemetryStore telemetrystore.Telemetry
telemetryStore,
metadataStore,
nil, // prometheus
nil, // promV2
nil, // traceStmtBuilder
logStmtBuilder,
nil, // auditStmtBuilder
@@ -109,6 +111,7 @@ func prepareQuerierForTraces(t *testing.T, telemetryStore telemetrystore.Telemet
telemetryStore,
metadataStore,
nil, // prometheus
nil, // promV2
traceStmtBuilder,
nil, // logStmtBuilder
nil, // auditStmtBuilder

View File

@@ -3,6 +3,8 @@ package querybuilder
import (
"context"
"log/slog"
"maps"
"slices"
"strings"
chparser "github.com/AfterShip/clickhouse-sql-parser/parser"
@@ -25,7 +27,23 @@ var internalDatabases = map[string]struct{}{
"information_schema": {},
}
// The parser's grammar has gaps against SQL that ClickHouse itself accepts. See TestErrIfStatementIsNotValid_ShouldPassButFails.
// generatorTableFunctions compute their rows from their arguments alone. They open no file or socket, reach no other host, and name no table, database or dictionary, so none of them can read through anything the rules here exist to protect. Can be used to build a dense axis to join a sparse series against. Every other table function is refused.
//
// Keyed by the lowercased name so that matching is case-insensitive, valued by the spelling to name it back to the caller.
//
// TODO(@therealpandey): take a deployment level allow list on top of this, so an operator can permit more without a release.
var generatorTableFunctions = map[string]string{
"numbers": "numbers",
"numbers_mt": "numbers_mt",
"zeros": "zeros",
"zeros_mt": "zeros_mt",
"generateseries": "generateSeries",
"generate_series": "generate_series",
}
var generatorTableFunctionsMessage = "allowed table functions are " + strings.Join(slices.Sorted(maps.Values(generatorTableFunctions)), ", ")
// The parser's grammar has gaps against SQL that ClickHouse itself accepts.
func ErrIfStatementIsNotValid(query string) (err error) {
defer func() {
// The parser has a history of panicking on malformed input rather than returning an error.
@@ -52,8 +70,17 @@ func ErrIfStatementIsNotValid(query string) (err error) {
visitor := &chparser.DefaultASTVisitor{Visit: func(node chparser.Expr) error {
switch expr := node.(type) {
case *chparser.TableFunctionExpr:
// Source table functions remain usable in ClickHouse read-only mode.
return errors.NewInvalidInputf(CodeClickHouseSQLTableFunction, "ClickHouse table functions are not allowed in SQL queries: %s", chparser.Format(expr.Name))
// Source table functions remain usable in ClickHouse read-only mode. Arguments are
// visited before this, so a read smuggled into one is already refused by the time
// an allowed generator gets here.
name := chparser.Format(expr.Name)
if _, ok := generatorTableFunctions[strings.ToLower(name)]; ok {
return nil
}
return errors.
NewInvalidInputf(CodeClickHouseSQLTableFunction, "ClickHouse table functions are not allowed in SQL queries: %s", name).
WithAdditional(generatorTableFunctionsMessage)
case *chparser.TableIdentifier:
// Reading these is unaffected by ClickHouse read-only mode.

View File

@@ -2,12 +2,11 @@ package querybuilder
import (
"testing"
"time"
"github.com/SigNoz/signoz/pkg/errors"
chparser "github.com/AfterShip/clickhouse-sql-parser/parser"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestErrIfStatementIsNotValid_Pass(t *testing.T) {
@@ -21,6 +20,9 @@ func TestErrIfStatementIsNotValid_Pass(t *testing.T) {
{"CommonTableExpression", "WITH t AS (SELECT fingerprint FROM signoz_metrics.time_series_v4) SELECT * FROM t"},
{"Join", "SELECT * FROM t1 LEFT JOIN t2 ON t1.a = t2.b"},
{"GlobalIn", "SELECT a FROM t WHERE a GLOBAL IN (SELECT b FROM t2)"},
// GLOBAL parsed only when the join type was omitted, and only before IN. https://github.com/AfterShip/clickhouse-sql-parser/pull/293
{"GlobalLeftJoin", "SELECT * FROM t1 GLOBAL LEFT JOIN t2 ON t1.a = t2.a"},
{"GlobalNotIn", "SELECT a FROM t WHERE a GLOBAL NOT IN (SELECT b FROM t2)"},
{"Union", "SELECT * FROM t UNION ALL SELECT * FROM t2"},
{"Intersect", "SELECT * FROM t INTERSECT SELECT * FROM t2"},
{"WindowFunction", "SELECT sum(v) OVER (PARTITION BY a ORDER BY t) FROM t"},
@@ -36,18 +38,50 @@ func TestErrIfStatementIsNotValid_Pass(t *testing.T) {
// order by interval
{"OrderByInterval", "SELECT toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS interval ORDER BY interval"},
{"OrderByIntervalAndDirection", "SELECT toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS `interval` ORDER BY `interval` ASC"},
// Unspaced, so rejected until the parser stopped lexing a signed literal after a
// closing bracket. The spaced form above no longer needs to be spaced.
// https://github.com/AfterShip/clickhouse-sql-parser/issues/286
// `interval` is a unit keyword, so unquoting it was rejected everywhere the parser
// expected a plain identifier. https://github.com/AfterShip/clickhouse-sql-parser/pull/296
{"OrderByUnquotedIntervalAsc", "SELECT toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS interval FROM t GROUP BY interval ORDER BY interval ASC"},
{"OrderByUnquotedIntervalDesc", "SELECT toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS interval FROM t GROUP BY interval ORDER BY interval DESC"},
{"UnquotedIntervalInGroupByTuple", "SELECT a FROM t GROUP BY (`service.name`, `service.version`, interval)"},
{"UnquotedIntervalProductionQuery", "SELECT toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS interval, resource_string_service$$name AS `service.name`, attributes_string['http.route'] AS `http.route`, quantile(0.95)(duration_nano) / 1000000000 AS value FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_string_service$$name = 'svc-a' AND resources_string['deployment.environment'] = 'dev' AND attributes_string['http.route'] = '/v1' AND http_method = 'POST' AND timestamp BETWEEN toDateTime(1784601720) AND toDateTime(1784602620) AND ts_bucket_start BETWEEN 1784601720 - 1800 AND 1784602620 GROUP BY `service.name`, `http.route`, interval ORDER BY interval ASC"},
// Separating the two readings of INTERVAL needs backtracking as per the current implementation which could have performance regressions.
// https://github.com/AfterShip/clickhouse-sql-parser/pull/296#issuecomment-5150316367
{"UnquotedIntervalRepeatedThirtyTimes", "SELECT interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval + interval AS total FROM t WHERE interval > 0 ORDER BY interval ASC"},
{"SignedLiteralAfterClosingParenUnspaced", "SELECT now() AS ts, toFloat64(count()) AS value FROM ( SELECT attributes_string['TableName'] AS T, attributes_string['MissingId'] AS M, max(fromUnixTimestamp64Nano(timestamp)) AS last_seen, dateDiff('minute', min(fromUnixTimestamp64Nano(timestamp)), max(fromUnixTimestamp64Nano(timestamp))) AS age_min FROM signoz_logs.distributed_logs_v2 WHERE body='missing_map_record' AND timestamp >= (toUnixTimestamp(now())-3600)*1000000000 GROUP BY T, M ) WHERE age_min >= 20 AND last_seen >= now() - toIntervalMinute(8)"},
{"SignedLiteralAfterClosingParenMinimal", "SELECT (1)-1"},
{"TrimFunction", "SELECT trimBoth('/api/endpoint/', '/');"},
// The SQL-standard keyword-separated argument forms, which took commas only. https://github.com/AfterShip/clickhouse-sql-parser/pull/290
{"StandardTrimSyntax", "SELECT trim(BOTH ' ' FROM body) FROM t"},
{"StandardSubstringSyntax", "SELECT substring(body FROM 2 FOR 3) FROM t"},
{"StandardOverlaySyntax", "SELECT overlay(body PLACING 'x' FROM 2) FROM t"},
// Row generators compute their rows from their arguments, so they read through nothing. This is the shape they get used for: a dense interval axis to CROSS JOIN a sparse series against.
{"NumbersTableFunction", "SELECT intervals.interval AS interval, active.cluster AS cluster, toFloat64(if(ts_data.has_data = 0, 0, 1)) AS value FROM ( SELECT DISTINCT JSONExtractString(labels, 'k8s.cluster.name') AS cluster FROM signoz_metrics.distributed_time_series_v4 WHERE metric_name = 'my_metric' AND unix_milli >= toUnixTimestamp(now() - INTERVAL 30 DAY) * 1000 HAVING cluster != '' ) AS active CROSS JOIN ( SELECT toStartOfInterval( toDateTime(toUnixTimestamp(now() - INTERVAL 30 MINUTE) + number * 60), INTERVAL 1 MINUTE ) AS interval FROM numbers(31) ) AS intervals LEFT JOIN ( SELECT toStartOfInterval( toDateTime(intDiv(s.unix_milli, 1000)), INTERVAL 1 MINUTE ) AS interval, JSONExtractString(ts.labels, 'k8s.cluster.name') AS cluster, 1 AS has_data FROM signoz_metrics.distributed_samples_v4 s INNER JOIN ( SELECT DISTINCT fingerprint, labels FROM signoz_metrics.distributed_time_series_v4 WHERE metric_name = 'my_metric' ) AS ts ON s.fingerprint = ts.fingerprint WHERE s.metric_name = 'my_metric' AND s.unix_milli >= toUnixTimestamp(now() - INTERVAL 30 MINUTE) * 1000 GROUP BY interval, cluster ) AS ts_data ON active.cluster = ts_data.cluster AND intervals.interval = ts_data.interval ORDER BY interval ASC"},
{"NumbersMtTableFunction", "SELECT * FROM numbers_mt(31)"},
{"ZerosTableFunction", "SELECT * FROM zeros(31)"},
{"ZerosMtTableFunction", "SELECT * FROM zeros_mt(31)"},
{"GenerateSeriesTableFunction", "SELECT * FROM generateSeries(1, 10)"},
{"GenerateSeriesSnakeCaseTableFunction", "SELECT * FROM generate_series(1, 10)"},
{"GeneratorTableFunctionUppercase", "SELECT * FROM NUMBERS(31)"},
{"GeneratorTableFunctionParenthesisedArgument", "SELECT * FROM NUMBERS((31))"},
{"GeneratorTableFunctionInJoin", "SELECT * FROM signoz_logs.distributed_logs_v2 AS l CROSS JOIN numbers(31) AS n"},
{"GeneratorTableFunctionInCommonTableExpression", "WITH axis AS (SELECT number FROM numbers(31)) SELECT * FROM axis"},
{"GeneratorTableFunctionInWhereSubquery", "SELECT * FROM t WHERE a IN (SELECT number FROM numbers(31))"},
{"GeneratorTableFunctionInUnion", "SELECT number FROM numbers(31) UNION ALL SELECT number FROM zeros(31)"},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
err := ErrIfStatementIsNotValid(testCase.query)
assert.NoError(t, err)
// Bounded rather than called directly: a parser that backtracks without memoising
// hangs instead of returning. Every case here parses in well under a millisecond.
errC := make(chan error, 1)
go func() { errC <- ErrIfStatementIsNotValid(testCase.query) }()
select {
case err := <-errC:
assert.NoError(t, err)
case <-time.After(10 * time.Second):
assert.Fail(t, "timed out, which means the parser is no longer bounding its backtracking")
}
})
}
}
@@ -70,6 +104,8 @@ func TestErrIfStatementIsNotValid_Fail(t *testing.T) {
{"CreateTable", "CREATE TABLE evil (a Int) ENGINE = Memory", CodeClickHouseSQLNotSelect},
{"Grant", "GRANT ALL ON *.* TO admin", CodeClickHouseSQLNotSelect},
{"Set", "SET readonly = 0", CodeClickHouseSQLNotSelect},
// The parser still dereferences nil on a DEFAULT expression it cannot read, so the recover is what turns this into a rejection rather than a crash.
{"UnparseableDefaultExpression", "CREATE TABLE t (a String DEFAULT foo(b FROM 2)) ENGINE = Memory", CodeClickHouseSQLParserPanic},
// These the parser rejects outright rather than classifying.
{"ShowGrants", "SHOW GRANTS", CodeClickHouseSQLUnparseable},
{"IntoOutfile", "SELECT * FROM t INTO OUTFILE '/tmp/x.csv'", CodeClickHouseSQLUnparseable},
@@ -81,6 +117,20 @@ func TestErrIfStatementIsNotValid_Fail(t *testing.T) {
{"TableFunctionInCommonTableExpression", "WITH c AS (SELECT * FROM url('http://x', CSV, 'a String')) SELECT * FROM c", CodeClickHouseSQLTableFunction},
{"TableFunctionInWhereSubquery", "SELECT * FROM t WHERE a IN (SELECT * FROM file('/etc/passwd', CSV, 'a String'))", CodeClickHouseSQLTableFunction},
{"TableFunctionInUnion", "SELECT * FROM t UNION ALL SELECT * FROM url('http://x', CSV, 'a String')", CodeClickHouseSQLTableFunction},
// These reach the internal databases without ever naming one, so the table-function rule is the only thing that sees them.
{"MergeTableFunction", "SELECT * FROM merge('system', '.*')", CodeClickHouseSQLTableFunction},
{"RemoteTableFunction", "SELECT * FROM remote('other-host', 'system.users')", CodeClickHouseSQLTableFunction},
{"ClusterTableFunction", "SELECT * FROM cluster('c', 'system.users')", CodeClickHouseSQLTableFunction},
// Pure, but excluded: generateRandom streams rows the arguments do not bound, and values has no use here that an array literal does not already cover.
{"GenerateRandomTableFunction", "SELECT * FROM generateRandom('a UInt64')", CodeClickHouseSQLTableFunction},
{"ValuesTableFunction", "SELECT * FROM values('a UInt64', 1, 2)", CodeClickHouseSQLTableFunction},
// Arguments are visited before the table function itself, so allowing a generator does not give anyone a wrapper to smuggle a read through.
{"InternalDatabaseInsideAllowedTableFunction", "SELECT * FROM numbers((SELECT count() FROM system.users))", CodeClickHouseSQLInternalDatabase},
{"InternalDatabaseJoinedOntoAllowedTableFunction", "SELECT * FROM numbers(31) AS n JOIN system.users AS u ON 1 = 1", CodeClickHouseSQLInternalDatabase},
{"InternalDatabaseUnionedWithAllowedTableFunction", "SELECT number FROM numbers(31) UNION ALL SELECT name FROM system.users", CodeClickHouseSQLInternalDatabase},
{"RefusedTableFunctionJoinedOntoAllowedTableFunction", "SELECT * FROM numbers(31) AS n JOIN url('http://x', CSV, 'a String') AS u ON 1 = 1", CodeClickHouseSQLTableFunction},
{"RefusedTableFunctionInsideAllowedTableFunction", "SELECT * FROM numbers((SELECT count() FROM file('/etc/passwd', CSV, 'a String')))", CodeClickHouseSQLTableFunction},
{"InternalDatabaseInsideAllowedTableFunctionCommonTableExpression", "WITH axis AS (SELECT * FROM numbers((SELECT count() FROM system.users))) SELECT * FROM axis", CodeClickHouseSQLInternalDatabase},
// Internal databases, which hold grants and server metadata rather than telemetry.
{"SystemUsers", "SELECT * FROM system.users", CodeClickHouseSQLInternalDatabase},
{"SystemUppercase", "SELECT * FROM SYSTEM.USERS", CodeClickHouseSQLInternalDatabase},
@@ -103,50 +153,3 @@ func TestErrIfStatementIsNotValid_Fail(t *testing.T) {
})
}
}
// Queries the parser cannot read. ClickHouse runs all of them.
func TestErrIfStatementIsNotValid_ShouldPassButFails(t *testing.T) {
testCases := []struct {
name string
query string
// The construct the parser stops after, which is the one it cannot read.
expectedStopsAfter string
// The same construct written so the parser accepts it.
fix string
}{
{
name: "IntervalAliasInOrderBy",
query: "SELECT toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS interval, resource_string_service$$name AS `service.name`, attributes_string['http.route'] AS `http.route`, quantile(0.95)(duration_nano) / 1000000000 AS value FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_string_service$$name = 'svc-a' AND resources_string['deployment.environment'] = 'dev' AND attributes_string['http.route'] = '/v1' AND http_method = 'POST' AND timestamp BETWEEN toDateTime(1784601720) AND toDateTime(1784602620) AND ts_bucket_start BETWEEN 1784601720 - 1800 AND 1784602620 GROUP BY `service.name`, `http.route`, interval ORDER BY interval ASC",
expectedStopsAfter: "ORDER BY interval ASC",
fix: "SELECT count() AS interval FROM t ORDER BY `interval` ASC",
},
{
name: "IntervalAliasInOrderByDesc",
query: "SELECT count() AS value, toStartOfInterval(timestamp, INTERVAL 1 MINUTE) AS interval, serviceName, resourceTagsMap['deployment.environment'] AS environment, exceptionStacktrace FROM signoz_traces.distributed_signoz_error_index_v2 WHERE exceptionType != 'OSError' AND resourceTagsMap['deployment.environment'] = 'staging' AND timestamp BETWEEN toDateTime(1785186300) AND toDateTime(1785186600) GROUP BY serviceName, interval, environment, exceptionStacktrace ORDER BY interval DESC",
expectedStopsAfter: "ORDER BY interval DESC",
fix: "SELECT count() AS interval FROM t ORDER BY `interval` DESC",
},
{
name: "StandardTrimSyntax",
query: "SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 5 MINUTE) AS interval, resources_string['host.name'] as host_name, toFloat64(countIf( lower(trim(BOTH ' ' FROM replaceOne( JSONExtractString(body, 'Action'), 'health_status: ', '' ))) IN ('unhealthy','starting','failing') )) as value FROM signoz_logs.distributed_logs_v2 WHERE timestamp BETWEEN 1784602320000000000 AND 1784602620000000000 AND ts_bucket_start BETWEEN 1784602320 - 300 AND 1784602620 AND JSONExtractString(body, 'Type') = 'container' AND JSONExtractString(body, 'Actor', 'Attributes', 'name') IS NOT NULL AND resources_string['host.name'] IS NOT NULL AND resources_string['host.name'] = 'aihub-nightly' GROUP BY interval, host_name ORDER BY interval, host_name",
expectedStopsAfter: "trim(BOTH '",
fix: "SELECT trimBoth(replaceOne( JSONExtractString(body, 'Action'), 'health_status: ', '' ), ' ')",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
err := ErrIfStatementIsNotValid(testCase.query)
var parseErr *chparser.ParseError
require.ErrorAs(t, err, &parseErr, "expected a parser failure rather than a rule violation")
// The parser reports the offset it stopped at, which sits just past the construct
// it choked on, so the text leading up to it is what needs looking at.
consumed := testCase.query[:parseErr.Pos]
assert.Equal(t, testCase.expectedStopsAfter, consumed[max(0, len(consumed)-len(testCase.expectedStopsAfter)):])
assert.NoError(t, ErrIfStatementIsNotValid(testCase.fix))
})
}
}

View File

@@ -39,7 +39,6 @@ import (
"github.com/SigNoz/signoz/pkg/sqlmigrator"
"github.com/SigNoz/signoz/pkg/sqlschema"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/statsreporter"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/tokenizer"
@@ -98,9 +97,6 @@ type Config struct {
// Querier config
Querier querier.Config `mapstructure:"querier"`
// StatementBuilder config
StatementBuilder statementbuilder.Config `mapstructure:"statementbuilder"`
// Ruler config
Ruler ruler.Config `mapstructure:"ruler"`
@@ -170,7 +166,6 @@ func NewConfig(ctx context.Context, logger *slog.Logger, resolverConfig config.R
prometheus.NewConfigFactory(),
alertmanager.NewConfigFactory(),
querier.NewConfigFactory(),
statementbuilder.NewConfigFactory(),
ruler.NewConfigFactory(),
emailing.NewConfigFactory(),
sharder.NewConfigFactory(),

View File

@@ -45,6 +45,7 @@ import (
"github.com/SigNoz/signoz/pkg/pprof/nooppprof"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/SigNoz/signoz/pkg/prometheus/clickhouseprometheus"
"github.com/SigNoz/signoz/pkg/prometheus/clickhouseprometheusv2"
"github.com/SigNoz/signoz/pkg/querier"
"github.com/SigNoz/signoz/pkg/querier/signozquerier"
"github.com/SigNoz/signoz/pkg/sharder"
@@ -231,6 +232,7 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewMigrateDashboardsV1ToV2Factory(sqlstore, sqlschema, dashboardStore, tagModule),
sqlmigration.NewFillDashboardMeterSourceFactory(sqlstore, dashboardStore),
sqlmigration.NewUpdateRoleTransactionGroupsFactory(),
sqlmigration.NewFillDashboardSpecCollectionsFactory(sqlstore, dashboardStore),
)
}
@@ -248,6 +250,7 @@ func NewTelemetryStoreProviderFactories() factory.NamedMap[factory.ProviderFacto
func NewPrometheusProviderFactories(telemetryStore telemetrystore.TelemetryStore) factory.NamedMap[factory.ProviderFactory[prometheus.Prometheus, prometheus.Config]] {
return factory.MustNewNamedMap(
clickhouseprometheus.NewFactory(telemetryStore),
clickhouseprometheusv2.NewFactory(telemetryStore),
)
}
@@ -289,9 +292,9 @@ func NewStatsReporterProviderFactories(aggregator statsreporter.Aggregator, orgG
)
}
func NewQuerierProviderFactories(telemetryStore telemetrystore.TelemetryStore, prometheus prometheus.Prometheus, metadataStore telemetrytypes.MetadataStore, traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation], logStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation], auditStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation], metricStmtBuilder qbtypes.StatementBuilder[qbtypes.MetricAggregation], meterStmtBuilder qbtypes.StatementBuilder[qbtypes.MetricAggregation], traceOperatorStmtBuilder qbtypes.TraceOperatorStatementBuilder, bucketCache querier.BucketCache, flagger flagger.Flagger) factory.NamedMap[factory.ProviderFactory[querier.Querier, querier.Config]] {
func NewQuerierProviderFactories(telemetryStore telemetrystore.TelemetryStore, prometheus prometheus.Prometheus, promV2 prometheus.Prometheus, metadataStore telemetrytypes.MetadataStore, traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation], logStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation], auditStmtBuilder qbtypes.StatementBuilder[qbtypes.LogAggregation], metricStmtBuilder qbtypes.StatementBuilder[qbtypes.MetricAggregation], meterStmtBuilder qbtypes.StatementBuilder[qbtypes.MetricAggregation], traceOperatorStmtBuilder qbtypes.TraceOperatorStatementBuilder, bucketCache querier.BucketCache, flagger flagger.Flagger) factory.NamedMap[factory.ProviderFactory[querier.Querier, querier.Config]] {
return factory.MustNewNamedMap(
signozquerier.NewFactory(telemetryStore, prometheus, metadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, flagger),
signozquerier.NewFactory(telemetryStore, prometheus, promV2, metadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, flagger),
)
}

View File

@@ -40,6 +40,7 @@ import (
"github.com/SigNoz/signoz/pkg/modules/tag/impltag"
"github.com/SigNoz/signoz/pkg/modules/user/impluser"
"github.com/SigNoz/signoz/pkg/prometheus"
"github.com/SigNoz/signoz/pkg/prometheus/clickhouseprometheusv2"
"github.com/SigNoz/signoz/pkg/querier"
"github.com/SigNoz/signoz/pkg/queryparser"
"github.com/SigNoz/signoz/pkg/ruler"
@@ -122,7 +123,7 @@ func newQueryStack(
) {
metadataStore := telemetrymetadata.NewTelemetryMetaStore(settings, telemetryStore, fl)
cfg := config.StatementBuilder
cfg := config.Querier.Config
traceStmtBuilder, err := tracesstatementbuilder.NewFactory(telemetryStore, metadataStore, fl).New(ctx, settings, cfg)
if err != nil {
return nil, nil, nil, nil, nil, nil, nil, nil, err
@@ -299,6 +300,11 @@ func New(
retentionGetter := implretention.NewGetter(implretention.NewStore(sqlstore))
// promV2 is the clickhousev2 provider handed to the querier for shadow
// comparison and pinned serving (declared before the serving provider,
// whose variable shadows the package name below).
var promV2 prometheus.Prometheus
// Initialize prometheus from the available prometheus provider factories
prometheus, err := factory.NewProviderFromNamedMap(
ctx,
@@ -311,6 +317,23 @@ func New(
return nil, err
}
// With the default provider, also stand up the clickhousev2 provider for
// the querier: PromQL queries shadow-compare against it behind the
// use_prometheus_clickhouse_v2 flag (see pkg/querier/promql_shadow.go).
// It never serves by default. An explicit
// prometheus::provider: clickhousev2 makes v2 the serving provider
// outright, so there is nothing to compare against.
if config.Prometheus.Provider() == "clickhouse" {
v2Config := config.Prometheus
// The v2 engine only evaluates shadow and pinned queries; disable its
// active query tracker so two trackers never share a file.
v2Config.ActiveQueryTrackerConfig.Enabled = false
promV2, err = clickhouseprometheusv2.New(ctx, providerSettings, v2Config, telemetrystore)
if err != nil {
return nil, err
}
}
// Assemble the query stack (metadata store, statement builders, bucket cache) once,
// and reuse the single metadata store everywhere downstream.
telemetryMetadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, err := newQueryStack(ctx, providerSettings, config, telemetrystore, cache, flagger)
@@ -323,7 +346,7 @@ func New(
ctx,
providerSettings,
config.Querier,
NewQuerierProviderFactories(telemetrystore, prometheus, telemetryMetadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, flagger),
NewQuerierProviderFactories(telemetrystore, prometheus, promV2, telemetryMetadataStore, traceStmtBuilder, logStmtBuilder, auditStmtBuilder, metricStmtBuilder, meterStmtBuilder, traceOperatorStmtBuilder, bucketCache, flagger),
config.Querier.Provider(),
)
if err != nil {

View File

@@ -0,0 +1,124 @@
package sqlmigration
import (
"context"
"log/slog"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/types"
"github.com/SigNoz/signoz/pkg/types/dashboardtypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
)
// Required, non-nullable v2 spec fields, mapped to their empty value.
var nullableSpecCollections = map[string]any{
"variables": []any{},
"panels": map[string]any{},
"layouts": []any{},
}
type fillDashboardSpecCollections struct {
sqlstore sqlstore.SQLStore
dashboardStore dashboardtypes.Store
settings factory.ProviderSettings
}
func NewFillDashboardSpecCollectionsFactory(sqlstore sqlstore.SQLStore, dashboardStore dashboardtypes.Store) factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(
factory.MustNewName("fill_dashboard_spec_collections"),
func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &fillDashboardSpecCollections{sqlstore: sqlstore, dashboardStore: dashboardStore, settings: ps}, nil
},
)
}
func (migration *fillDashboardSpecCollections) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
// Up replaces a missing or null spec.variables / spec.panels / spec.layouts with the
// empty collection. One transaction; v1 dashboards are skipped.
func (migration *fillDashboardSpecCollections) Up(ctx context.Context, _ *bun.DB) error {
return migration.sqlstore.RunInTxCtx(ctx, nil, func(ctx context.Context) error {
var orgIDs []string
if err := migration.sqlstore.BunDBCtx(ctx).NewSelect().Model((*types.Organization)(nil)).Column("id").Scan(ctx, &orgIDs); err != nil {
return err
}
for _, id := range orgIDs {
orgID, err := valuer.NewUUID(id)
if err != nil {
return err
}
if err := migration.fillOrg(ctx, orgID); err != nil {
return err
}
}
return nil
})
}
// fillOrg fills every v2 dashboard in the org that needs it, inside the caller's transaction.
func (migration *fillDashboardSpecCollections) fillOrg(ctx context.Context, orgID valuer.UUID) error {
// List, not ListV2: ListV2 paginates and excludes system dashboards; a migration needs every row.
storables, err := migration.dashboardStore.List(ctx, orgID)
if err != nil {
return err
}
logger := migration.settings.Logger
var stillInV1, malformedSpec, skippedNoNulls, migrated int
for _, storable := range storables {
if !storable.IsV2() {
stillInV1++
continue
}
// Raw data, not ToDashboardV2: decoding validates, and these are the rows it rejects.
spec, ok := storable.Data["spec"].(map[string]any)
if !ok {
malformedSpec++
logger.WarnContext(ctx, "v2 dashboard has no spec object; leaving it untouched", slog.String("org_id", orgID.String()), slog.String("dashboard_id", storable.ID.String()))
continue
}
if !fillSpecCollections(spec) {
skippedNoNulls++
continue
}
if err := migration.dashboardStore.Update(ctx, orgID, storable); err != nil {
return err
}
migrated++
}
logger.InfoContext(ctx, "filled required collections on v2 dashboards",
slog.String("org_id", orgID.String()),
slog.Int("total", len(storables)),
slog.Int("still_in_v1", stillInV1),
slog.Int("malformed_spec", malformedSpec),
slog.Int("skipped_no_nulls", skippedNoNulls),
slog.Int("migrated", migrated),
)
return nil
}
// fillSpecCollections empties each absent or null required collection, reporting whether
// anything changed. A present value is left alone whatever its shape, so a malformed one
// still surfaces as a validation error.
func fillSpecCollections(spec map[string]any) bool {
changed := false
for field, empty := range nullableSpecCollections {
if value, present := spec[field]; !present || value == nil {
spec[field] = empty
changed = true
}
}
return changed
}
func (migration *fillDashboardSpecCollections) Down(context.Context, *bun.DB) error {
return nil
}

View File

@@ -2,7 +2,6 @@ package statementbuilder
import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
)
// SkipResourceFingerprint configures when the resource fingerprint subquery is skipped in favor of main-table filtering.
@@ -18,12 +17,7 @@ type Config struct {
SkipResourceFingerprint SkipResourceFingerprint `yaml:"skip_resource_fingerprint" mapstructure:"skip_resource_fingerprint"`
}
// NewConfigFactory creates a new config factory for the statement builders.
func NewConfigFactory() factory.ConfigFactory {
return factory.NewConfigFactory(factory.MustNewName("statementbuilder"), newConfig)
}
func newConfig() factory.Config {
func NewConfig() Config {
return Config{
SkipResourceFingerprint: SkipResourceFingerprint{
Enabled: false,

View File

@@ -8,6 +8,7 @@ import (
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -51,7 +52,8 @@ func TestStatementBuilderGroupByUnknownKey(t *testing.T) {
statementBuilder := NewLogQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore, fm, cb, aggExprRewriter,
logstelemetryschema.DefaultFullTextColumn, fl, nil, false, 100000,
logstelemetryschema.DefaultFullTextColumn, fl, nil,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
query := qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{

View File

@@ -11,6 +11,7 @@ import (
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -1346,8 +1347,7 @@ func buildJSONTestStatementBuilder(t *testing.T, addIndexes bool) (*logQueryStat
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
return statementBuilder, mockMetadataStore

View File

@@ -62,7 +62,7 @@ func NewFactory(
aggExprRewriter := querybuilder.NewAggExprRewriter(settings, logstelemetryschema.DefaultFullTextColumn, fm, cb, fl)
return NewLogQueryStatementBuilder(
settings, metadataStore, fm, cb, aggExprRewriter, logstelemetryschema.DefaultFullTextColumn,
fl, telemetryStore, cfg.SkipResourceFingerprint.Enabled, cfg.SkipResourceFingerprint.Threshold,
fl, telemetryStore, cfg,
), nil
},
)
@@ -77,8 +77,7 @@ func NewLogQueryStatementBuilder(
fullTextColumn *telemetrytypes.TelemetryFieldKey,
fl flagger.Flagger,
telemetryStore telemetrystore.TelemetryStore,
skipResourceFingerprintEnable bool,
skipResourceFingerprintThreshold uint64,
cfg statementbuilder.Config,
) *logQueryStatementBuilder {
logsSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema")
@@ -92,7 +91,7 @@ func NewLogQueryStatementBuilder(
fullTextColumn,
fl,
telemetryStore,
skipResourceFingerprintThreshold,
cfg.SkipResourceFingerprint.Threshold,
)
return &logQueryStatementBuilder{
@@ -103,7 +102,7 @@ func NewLogQueryStatementBuilder(
resourceFilterResolver: resourceFilterResolver,
aggExprRewriter: aggExprRewriter,
fl: fl,
skipResourceFingerprintEnabled: skipResourceFingerprintEnable,
skipResourceFingerprintEnabled: cfg.SkipResourceFingerprint.Enabled,
fullTextColumn: fullTextColumn,
}
}

View File

@@ -11,6 +11,7 @@ import (
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
@@ -231,8 +232,7 @@ func TestStatementBuilderTimeSeries(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
for _, c := range cases {
@@ -374,8 +374,7 @@ func TestStatementBuilderListQuery(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
for _, c := range cases {
@@ -525,8 +524,7 @@ func TestStatementBuilderListQueryResourceTests(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
for _, c := range cases {
@@ -603,8 +601,7 @@ func TestStatementBuilderTimeSeriesBodyGroupBy(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
for _, c := range cases {
@@ -700,8 +697,7 @@ func TestStatementBuilderListQueryServiceCollision(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
for _, c := range cases {
@@ -926,8 +922,7 @@ func TestAdjustKey(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
for _, c := range cases {
@@ -1073,8 +1068,7 @@ func TestStmtBuilderBodyField(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, c.requestType, c.query, nil)
@@ -1174,8 +1168,7 @@ func TestStmtBuilderBodyFullTextSearch(t *testing.T) {
logstelemetryschema.DefaultFullTextColumn,
fl,
nil,
false,
100000,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, c.requestType, c.query, nil)
@@ -1296,7 +1289,6 @@ func newSkipResourceFingerprintLogsBuilder(
logstelemetryschema.DefaultFullTextColumn,
fl,
telemetryStore,
skipEnable,
threshold,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: skipEnable, Threshold: threshold}},
)
}

View File

@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"log/slog"
"slices"
"time"
"github.com/SigNoz/signoz/pkg/factory"
@@ -16,7 +17,6 @@ import (
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"golang.org/x/exp/slices"
)
const (

View File

@@ -366,7 +366,7 @@ func derefValue(v any) any {
}
val := reflect.ValueOf(v)
for val.Kind() == reflect.Ptr {
for val.Kind() == reflect.Pointer {
if val.IsNil() {
return nil
}

View File

@@ -2303,6 +2303,15 @@ func unionTemporalities(existing, additional []metrictypes.Temporality) []metric
return existing
}
// resolveMetricType applies the non-monotonic-cumulative-sum-as-gauge rule.
// Monotonicity is only meaningful for cumulative sums; delta sums always stay Sum.
func resolveMetricType(metricType metrictypes.Type, isMonotonic bool, temporality metrictypes.Temporality) metrictypes.Type {
if metricType == metrictypes.SumType && !isMonotonic && temporality == metrictypes.Cumulative {
return metrictypes.GaugeType
}
return metricType
}
func (t *telemetryMetaStore) fetchTemporalityTypeForTable(ctx context.Context, tableName string, adjustedStartTs, adjustedEndTs uint64, metricNames []string, extraConds ...string) (map[string][]metrictypes.Temporality, map[string]metrictypes.Type, error) {
temporalities := make(map[string][]metrictypes.Temporality)
types := make(map[string]metrictypes.Type)
@@ -2339,9 +2348,7 @@ func (t *telemetryMetaStore) fetchTemporalityTypeForTable(ctx context.Context, t
if temporality != metrictypes.Unknown {
temporalities[metricName] = append(temporalities[metricName], temporality)
}
if metricType == metrictypes.SumType && !isMonotonic {
metricType = metrictypes.GaugeType
}
metricType = resolveMetricType(metricType, isMonotonic, temporality)
types[metricName] = metricType
}
if err := rows.Err(); err != nil {
@@ -2392,9 +2399,7 @@ func (t *telemetryMetaStore) fetchMeterSourceMetricsTemporalityAndType(ctx conte
if err := rows.Scan(&metricName, &temporality, &metricType, &isMonotonic); err != nil {
return nil, nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to scan temporality result")
}
if metricType == metrictypes.SumType && !isMonotonic {
metricType = metrictypes.GaugeType
}
metricType = resolveMetricType(metricType, isMonotonic, temporality)
temporalities[metricName] = temporality
types[metricName] = metricType
}

Some files were not shown because too many files have changed in this diff Show More