Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 39 additions & 1 deletion grafana-datasource-plugin/pkg/plugin/resources.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@ package plugin
import (
"database/sql"
"encoding/json"
"fmt"
"net/http"
"time"

"github.com/lib/pq"
)
Expand Down Expand Up @@ -47,7 +49,43 @@ type channelEntry struct {
}

func (d *Datasource) handleGetTelemetryChannels(w http.ResponseWriter, r *http.Request) {
rows, err := d.db.QueryContext(r.Context(), "SELECT component, name FROM telemetryDefs ORDER BY component, name;")
query := "SELECT component, name FROM telemetryDefs ORDER BY component, name;"
var args []any
if sources := r.URL.Query()["sources"]; len(sources) > 0 {
fromRaw := r.URL.Query().Get("from")
toRaw := r.URL.Query().Get("to")
if fromRaw == "" || toRaw == "" {
http.Error(w, "source-filtered channel queries require from and to", http.StatusBadRequest)
return
}
from, err := time.Parse(time.RFC3339Nano, fromRaw)
if err != nil {
http.Error(w, "invalid from time", http.StatusBadRequest)
return
}
to, err := time.Parse(time.RFC3339Nano, toRaw)
if err != nil || to.Before(from) {
http.Error(w, "invalid to time", http.StatusBadRequest)
return
}
timeField := r.URL.Query().Get("timeField")
if timeField == "" {
timeField = "ert"
}
if timeField != "time" && timeField != "ert" {
http.Error(w, "invalid time field", http.StatusBadRequest)
return
}
query = fmt.Sprintf(`SELECT d.component, d.name FROM telemetryDefs d
WHERE EXISTS (
SELECT 1 FROM telemetry t
WHERE t.telemetryDefId = d.id AND t.source = ANY($1)
AND t.%s >= $2 AND t.%s <= $3
)
ORDER BY d.component, d.name;`, timeField, timeField)
args = append(args, pq.Array(sources), from, to)
}
rows, err := d.db.QueryContext(r.Context(), query, args...)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
Expand Down
128 changes: 128 additions & 0 deletions grafana-datasource-plugin/pkg/plugin/resources_sources_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
package plugin

import (
"errors"
"fmt"
"net/http"
"net/http/httptest"
"net/url"
"regexp"
"testing"
"time"

"github.com/DATA-DOG/go-sqlmock"
"github.com/lib/pq"
)

func telemetryChannelsFilteredQuery(timeField string) string {
return fmt.Sprintf(`SELECT d.component, d.name FROM telemetryDefs d
WHERE EXISTS (
SELECT 1 FROM telemetry t
WHERE t.telemetryDefId = d.id AND t.source = ANY($1)
AND t.%s >= $2 AND t.%s <= $3
)
ORDER BY d.component, d.name;`, timeField, timeField)
}

func TestTelemetryChannelsSourceFilter(t *testing.T) {
const allQuery = "SELECT component, name FROM telemetryDefs ORDER BY component, name;"
from := time.Date(2026, 9, 14, 18, 0, 0, 0, time.UTC)
to := from.Add(time.Hour)
for _, tt := range []struct {
name string
sources []string
timeField string
empty bool
}{
{name: "all sources"},
{name: "one source", sources: []string{"FSW-A"}},
{name: "multiple sources by onboard time", sources: []string{"FSW-A", "FSW-B"}, timeField: "time"},
{name: "quoted source", sources: []string{"FSW's \"test\",\\source"}},
{name: "empty source identifier", sources: []string{""}},
{name: "no matching source", sources: []string{"unknown"}, empty: true},
} {
t.Run(tt.name, func(t *testing.T) {
db, mock, err := sqlmock.New()
if err != nil {
t.Fatal(err)
}
defer func() { _ = db.Close() }()
values := url.Values{"sources": tt.sources}
var expected *sqlmock.ExpectedQuery
if len(tt.sources) > 0 {
field := tt.timeField
if field == "" {
field = "ert"
}
query := telemetryChannelsFilteredQuery(field)
expected = mock.ExpectQuery(regexp.QuoteMeta(query)).WithArgs(pq.Array(tt.sources), from, to)
values.Set("from", from.Format(time.RFC3339Nano))
values.Set("to", to.Format(time.RFC3339Nano))
values.Set("timeField", field)
} else {
expected = mock.ExpectQuery(regexp.QuoteMeta(allQuery)).WithoutArgs()
}
rows := sqlmock.NewRows([]string{"component", "name"})
body := "[]"
if !tt.empty {
rows.AddRow("CDH", "Temperature")
body = `[{"component":"CDH","name":"Temperature"}]`
}
expected.WillReturnRows(rows).RowsWillBeClosed()
request := httptest.NewRequest(http.MethodGet, "/telemetry/channels?"+values.Encode(), nil)
response := httptest.NewRecorder()
(&Datasource{db: db}).handleGetTelemetryChannels(response, request)
if response.Code != http.StatusOK || response.Body.String() != body {
t.Fatalf("got %d %s, want 200 %s", response.Code, response.Body.String(), body)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Fatal(err)
}
})
}
}

func TestTelemetryChannelsSourceFilterRejectsUnboundedOrInvalidRanges(t *testing.T) {
for _, rawQuery := range []string{
"sources=FSW-A",
"sources=FSW-A&from=2026-09-14T18%3A00%3A00Z&to=2026-09-14T19%3A00%3A00Z&timeField=bogus",
"sources=FSW-A&from=2026-09-14T20%3A00%3A00Z&to=2026-09-14T19%3A00%3A00Z&timeField=ert",
} {
db, _, err := sqlmock.New()
if err != nil {
t.Fatal(err)
}
response := httptest.NewRecorder()
(&Datasource{db: db}).handleGetTelemetryChannels(response, httptest.NewRequest(http.MethodGet, "/telemetry/channels?"+rawQuery, nil))
_ = db.Close()
if response.Code != http.StatusBadRequest {
t.Fatalf("query %q got status %d, want 400", rawQuery, response.Code)
}
}
}

func TestTelemetryChannelsFilteredQueryError(t *testing.T) {
db, mock, err := sqlmock.New()
if err != nil {
t.Fatal(err)
}
defer func() { _ = db.Close() }()
from := time.Date(2026, 9, 14, 18, 0, 0, 0, time.UTC)
to := from.Add(time.Hour)
mock.ExpectQuery("SELECT d.component").WithArgs(pq.Array([]string{"FSW-A"}), from, to).
WillReturnError(errors.New("database unavailable"))
values := url.Values{
"sources": []string{"FSW-A"},
"from": []string{from.Format(time.RFC3339Nano)},
"to": []string{to.Format(time.RFC3339Nano)},
"timeField": []string{"ert"},
}
response := httptest.NewRecorder()
(&Datasource{db: db}).handleGetTelemetryChannels(response, httptest.NewRequest(http.MethodGet, "/telemetry/channels?"+values.Encode(), nil))
if response.Code != http.StatusInternalServerError {
t.Fatalf("got status %d, want 500", response.Code)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Fatal(err)
}
}
6 changes: 4 additions & 2 deletions grafana-datasource-plugin/src/components/BuilderEditor.tsx
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import React from 'react';
import { CollapsableSection, DateTimePicker, InlineField, RadioButtonGroup } from '@grafana/ui';
import { dateTime, DateTime, SelectableValue } from '@grafana/data';
import { dateTime, DateTime, SelectableValue, TimeRange } from '@grafana/data';
import { DataSource } from '../datasource';
import { MyQuery, QueryType, TimeField } from '../types';
import { TelemetryFields } from './TelemetryFields';
Expand All @@ -11,6 +11,7 @@ interface BuilderEditorProps {
onChange: (query: MyQuery) => void;
onRunQuery: () => void;
datasource: DataSource;
range?: TimeRange;
}

const QUERY_TYPE_OPTIONS: Array<SelectableValue<QueryType>> = [
Expand All @@ -24,7 +25,7 @@ const TIME_FIELD_OPTIONS: Array<SelectableValue<TimeField>> = [
{ label: 'On-board Time', value: 'time' },
];

export function BuilderEditor({ query, onChange, onRunQuery, datasource }: BuilderEditorProps) {
export function BuilderEditor({ query, onChange, onRunQuery, datasource, range }: BuilderEditorProps) {
const queryType = query.queryType ?? 'telemetry';

const onQueryTypeChange = (value: QueryType) => {
Expand Down Expand Up @@ -86,6 +87,7 @@ export function BuilderEditor({ query, onChange, onRunQuery, datasource }: Build
onChange={onChange}
onRunQuery={onRunQuery}
datasource={datasource}
range={range}
sharedOptions={sharedOptions}
/>
) : (
Expand Down
1 change: 1 addition & 0 deletions grafana-datasource-plugin/src/components/QueryEditor.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ export function QueryEditor({ query, onChange, onRunQuery, datasource, range }:
onChange={onChange}
onRunQuery={onRunQuery}
datasource={datasource}
range={range}
/>
)}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
import React from 'react';
import { act, render, waitFor } from '@testing-library/react';
import { MultiCombobox } from '@grafana/ui';
import { dateTime, TimeRange } from '@grafana/data';
import { TelemetryFields } from './TelemetryFields';
import { DataSource } from '../datasource';
import { ChannelRef, MyQuery } from '../types';

let sourceVariable = 'FSW-A';
jest.mock('@grafana/runtime', () => ({
getTemplateSrv: () => ({
replace: (value: string) => value.replace('$source', sourceVariable),
getVariables: () => [],
}),
}));
jest.mock('@grafana/ui', () => ({
InlineField: ({ children }: { children: React.ReactNode }) => children,
Combobox: () => null,
MultiCombobox: jest.fn(() => null),
}));
jest.mock('./TransformFields', () => ({ TransformFields: () => null }));

function channelProps() {
const calls = (MultiCombobox as unknown as jest.Mock).mock.calls;
return [...calls].reverse().find(([props]) => props['data-testid'] === 'query-editor-channel')![0];
}

function deferred() {
let resolve!: (channels: ChannelRef[]) => void;
let reject!: (error: Error) => void;
const promise = new Promise<ChannelRef[]>((res, rej) => { resolve = res; reject = rej; });
return { promise, resolve, reject };
}

function setup(sources: string[], getChannels = jest.fn().mockResolvedValue([]), range?: TimeRange) {
const datasource = {
getChannels,
getSources: jest.fn().mockResolvedValue([]),
getKeys: jest.fn().mockResolvedValue([]),
} as unknown as DataSource;
const query: MyQuery = {
refId: 'A', queryType: 'telemetry', aggregation: 'avg', sources,
channels: [{ component: 'CDH', name: 'Temperature' }], keys: [],
};
const onChange = jest.fn();
const onRunQuery = jest.fn();
const view = (nextSources: string[]) => (
<TelemetryFields query={{ ...query, sources: nextSources }} datasource={datasource} range={range}
onChange={onChange} onRunQuery={onRunQuery} />
);
const rendered = render(view(sources));
return { getChannels, onChange, ...rendered, setSources: (next: string[]) => rendered.rerender(view(next)) };
}

beforeEach(() => {
jest.clearAllMocks();
sourceVariable = 'FSW-A';
});


it('bounds source-filtered channel lookup to the active panel range', async () => {
const range: TimeRange = {
from: dateTime('2026-09-14T18:00:00Z'),
to: dateTime('2026-09-14T19:00:00Z'),
raw: { from: 'now-1h', to: 'now' },
};
const getChannels = jest.fn().mockResolvedValue([]);
setup(['FSW-A'], getChannels, range);
await waitFor(() => expect(getChannels).toHaveBeenCalledWith(['FSW-A'], {
from: '2026-09-14T18:00:00.000Z',
to: '2026-09-14T19:00:00.000Z',
timeField: 'ert',
}));
});

it('reloads channels for selected sources and restores all channels when cleared', async () => {
const { getChannels, setSources, onChange } = setup(['FSW-A']);
await waitFor(() => expect(getChannels).toHaveBeenLastCalledWith(['FSW-A']));
setSources(['FSW-A', 'FSW-B']);
await waitFor(() => expect(getChannels).toHaveBeenLastCalledWith(['FSW-A', 'FSW-B']));
setSources([]);
await waitFor(() => expect(getChannels).toHaveBeenLastCalledWith([]));
expect(onChange).not.toHaveBeenCalled();
});

it('does not reload for a new array containing the same sources', async () => {
const { getChannels, setSources } = setup(['FSW-A']);
await waitFor(() => expect(channelProps().loading).toBe(false));
setSources(['FSW-A']);
expect(getChannels).toHaveBeenCalledTimes(1);
});

it('reloads when a source template variable changes', async () => {
const { getChannels, setSources } = setup(['$source']);
await waitFor(() => expect(getChannels).toHaveBeenLastCalledWith(['FSW-A']));
sourceVariable = 'FSW-B';
setSources(['$source']);
await waitFor(() => expect(getChannels).toHaveBeenLastCalledWith(['FSW-B']));
});

it.each(['resolve', 'reject'] as const)('ignores an obsolete request that later %ss', async (settle) => {
const old = deferred();
const current = deferred();
const getChannels = jest.fn().mockReturnValueOnce(old.promise).mockReturnValueOnce(current.promise);
const { setSources } = setup(['FSW-A'], getChannels);
setSources(['FSW-B']);
const channels = [{ component: 'PWR', name: 'Current' }];
await act(async () => { current.resolve(channels); });
await act(async () => {
if (settle === 'resolve') {
old.resolve([{ component: 'CDH', name: 'Temperature' }]);
} else {
old.reject(new Error('obsolete request'));
}
});
expect(await channelProps().options('')).toEqual([
expect.objectContaining({ label: 'PWR.Current', value: JSON.stringify(channels[0]) }),
]);
expect(channelProps().loading).toBe(false);
});

it('keeps the loading state while a newer request is pending', async () => {
const old = deferred();
const current = deferred();
const getChannels = jest.fn().mockReturnValueOnce(old.promise).mockReturnValueOnce(current.promise);
const { setSources } = setup(['FSW-A'], getChannels);
setSources(['FSW-B']);
await act(async () => { old.resolve([]); });
expect(channelProps().loading).toBe(true);
await act(async () => { current.resolve([]); });
expect(channelProps().loading).toBe(false);
});

it('clears stale options and finishes loading when the current request fails', async () => {
const current = deferred();
const getChannels = jest.fn()
.mockResolvedValueOnce([{ component: 'CDH', name: 'Temperature' }])
.mockReturnValueOnce(current.promise);
const { setSources } = setup(['FSW-A'], getChannels);
await waitFor(() => expect(channelProps().loading).toBe(false));
setSources(['FSW-B']);
expect(await channelProps().options('')).toEqual([]);
await act(async () => { current.reject(new Error('request failed')); });
expect(await channelProps().options('')).toEqual([]);
expect(channelProps().loading).toBe(false);
});
Loading
Loading