feat: Refactor asyncEvent middleware and add websocket support (#13696)

This commit is contained in:
Rob DiCiuccio
2021-03-23 07:23:23 -07:00
committed by GitHub
parent 5b79f84e1b
commit 452b53092b
9 changed files with 460 additions and 421 deletions

View File

@@ -17,249 +17,231 @@
* under the License.
*/
import fetchMock from 'fetch-mock';
import WS from 'jest-websocket-mock';
import sinon from 'sinon';
import * as featureFlags from 'src/featureFlags';
import initAsyncEvents from 'src/middleware/asyncEvent';
jest.useFakeTimers();
import { parseErrorJson } from 'src/utils/getClientErrorObject';
import * as asyncEvent from 'src/middleware/asyncEvent';
describe('asyncEvent middleware', () => {
const next = sinon.spy();
const state = {
charts: {
123: {
id: 123,
status: 'loading',
asyncJobId: 'foo123',
const asyncPendingEvent = {
status: 'pending',
result_url: null,
job_id: 'foo123',
channel_id: '999',
errors: [],
};
const asyncDoneEvent = {
id: '1518951480106-0',
status: 'done',
result_url: '/api/v1/chart/data/cache-key-1',
job_id: 'foo123',
channel_id: '999',
errors: [],
};
const asyncErrorEvent = {
id: '1518951480107-0',
status: 'error',
result_url: null,
job_id: 'foo123',
channel_id: '999',
errors: [{ message: "Error: relation 'foo' does not exist" }],
};
const chartData = {
result: [
{
cache_key: '199f01f81f99c98693694821e4458111',
cached_dttm: null,
cache_timeout: 86400,
annotation_data: {},
error: null,
is_cached: false,
query:
'SELECT product_line AS product_line,\n sum(sales) AS "(Sales)"\nFROM cleaned_sales_data\nGROUP BY product_line\nLIMIT 50000',
status: 'success',
stacktrace: null,
rowcount: 7,
colnames: ['product_line', '(Sales)'],
coltypes: [1, 0],
data: [
{
product_line: 'Classic Cars',
'(Sales)': 3919615.66,
},
],
applied_filters: [
{
column: '__time_range',
},
],
rejected_filters: [],
},
345: {
id: 345,
status: 'loading',
asyncJobId: 'foo345',
},
},
};
const events = [
{
status: 'done',
result_url: '/api/v1/chart/data/cache-key-1',
job_id: 'foo123',
channel_id: '999',
errors: [],
},
{
status: 'done',
result_url: '/api/v1/chart/data/cache-key-2',
job_id: 'foo345',
channel_id: '999',
errors: [],
},
];
const mockStore = {
getState: () => state,
dispatch: sinon.stub(),
};
const action = {
type: 'GENERIC_ACTION',
],
};
const EVENTS_ENDPOINT = 'glob:*/api/v1/async_event/*';
const CACHED_DATA_ENDPOINT = 'glob:*/api/v1/chart/data/*';
const config = {
GLOBAL_ASYNC_QUERIES_TRANSPORT: 'polling',
GLOBAL_ASYNC_QUERIES_POLLING_DELAY: 500,
};
let featureEnabledStub: any;
function setup() {
const getPendingComponents = sinon.stub();
const successAction = sinon.spy();
const errorAction = sinon.spy();
const testCallback = sinon.stub();
const testCallbackPromise = sinon.stub();
testCallbackPromise.returns(
new Promise(resolve => {
testCallback.callsFake(resolve);
}),
);
return {
getPendingComponents,
successAction,
errorAction,
testCallback,
testCallbackPromise,
};
}
beforeEach(() => {
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: [] },
});
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 200,
body: { result: { some: 'data' } },
});
beforeEach(async () => {
featureEnabledStub = sinon.stub(featureFlags, 'isFeatureEnabled');
featureEnabledStub.withArgs('GLOBAL_ASYNC_QUERIES').returns(true);
});
afterEach(() => {
fetchMock.reset();
next.resetHistory();
featureEnabledStub.restore();
});
afterAll(fetchMock.reset);
it('should initialize and call next', () => {
const { getPendingComponents, successAction, errorAction } = setup();
getPendingComponents.returns([]);
const asyncEventMiddleware = initAsyncEvents({
config,
getPendingComponents,
successAction,
errorAction,
});
asyncEventMiddleware(mockStore)(next)(action);
expect(next.callCount).toBe(1);
});
describe('polling transport', () => {
const config = {
GLOBAL_ASYNC_QUERIES_TRANSPORT: 'polling',
GLOBAL_ASYNC_QUERIES_POLLING_DELAY: 50,
GLOBAL_ASYNC_QUERIES_WEBSOCKET_URL: '',
};
it('should fetch events when there are pending components', () => {
const {
getPendingComponents,
successAction,
errorAction,
testCallback,
testCallbackPromise,
} = setup();
getPendingComponents.returns(Object.values(state.charts));
const asyncEventMiddleware = initAsyncEvents({
config,
getPendingComponents,
successAction,
errorAction,
processEventsCallback: testCallback,
beforeEach(async () => {
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: [asyncDoneEvent] },
});
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 200,
body: { result: chartData },
});
asyncEvent.init(config);
});
asyncEventMiddleware(mockStore)(next)(action);
it('resolves with chart data on event done status', async () => {
await expect(
asyncEvent.waitForAsyncData(asyncPendingEvent),
).resolves.toEqual([chartData]);
return testCallbackPromise().then(() => {
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(1);
});
});
it('should fetch cached when there are successful events', () => {
const {
getPendingComponents,
successAction,
errorAction,
testCallback,
testCallbackPromise,
} = setup();
fetchMock.reset();
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: events },
});
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 200,
body: { result: { some: 'data' } },
});
getPendingComponents.returns(Object.values(state.charts));
const asyncEventMiddleware = initAsyncEvents({
config,
getPendingComponents,
successAction,
errorAction,
processEventsCallback: testCallback,
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(1);
});
asyncEventMiddleware(mockStore)(next)(action);
it('rejects on event error status', async () => {
fetchMock.reset();
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: [asyncErrorEvent] },
});
const errorResponse = await parseErrorJson(asyncErrorEvent);
await expect(
asyncEvent.waitForAsyncData(asyncPendingEvent),
).rejects.toEqual(errorResponse);
return testCallbackPromise().then(() => {
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(1);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(2);
expect(successAction.callCount).toBe(2);
});
});
it('should call errorAction for cache fetch error responses', () => {
const {
getPendingComponents,
successAction,
errorAction,
testCallback,
testCallbackPromise,
} = setup();
fetchMock.reset();
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: events },
});
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 400,
body: { errors: ['error'] },
});
getPendingComponents.returns(Object.values(state.charts));
const asyncEventMiddleware = initAsyncEvents({
config,
getPendingComponents,
successAction,
errorAction,
processEventsCallback: testCallback,
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(0);
});
asyncEventMiddleware(mockStore)(next)(action);
it('rejects on cached data fetch error', async () => {
fetchMock.reset();
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: [asyncDoneEvent] },
});
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 400,
});
const errorResponse = [{ error: 'Bad Request' }];
await expect(
asyncEvent.waitForAsyncData(asyncPendingEvent),
).rejects.toEqual(errorResponse);
return testCallbackPromise().then(() => {
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(1);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(2);
expect(errorAction.callCount).toBe(2);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(1);
});
});
it('should handle event fetching error responses', () => {
const {
getPendingComponents,
successAction,
errorAction,
testCallback,
testCallbackPromise,
} = setup();
fetchMock.reset();
fetchMock.get(EVENTS_ENDPOINT, {
status: 400,
body: { message: 'error' },
});
getPendingComponents.returns(Object.values(state.charts));
const asyncEventMiddleware = initAsyncEvents({
config,
getPendingComponents,
successAction,
errorAction,
processEventsCallback: testCallback,
describe('ws transport', () => {
let wsServer: WS;
const config = {
GLOBAL_ASYNC_QUERIES_TRANSPORT: 'ws',
GLOBAL_ASYNC_QUERIES_POLLING_DELAY: 50,
GLOBAL_ASYNC_QUERIES_WEBSOCKET_URL: 'ws://127.0.0.1:8080/',
};
beforeEach(async () => {
fetchMock.get(EVENTS_ENDPOINT, {
status: 200,
body: { result: [asyncDoneEvent] },
});
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 200,
body: { result: chartData },
});
wsServer = new WS(config.GLOBAL_ASYNC_QUERIES_WEBSOCKET_URL);
asyncEvent.init(config);
});
asyncEventMiddleware(mockStore)(next)(action);
return testCallbackPromise().then(() => {
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(1);
});
});
it('should not fetch events when async queries are disabled', () => {
featureEnabledStub.restore();
featureEnabledStub = sinon.stub(featureFlags, 'isFeatureEnabled');
featureEnabledStub.withArgs('GLOBAL_ASYNC_QUERIES').returns(false);
const { getPendingComponents, successAction, errorAction } = setup();
getPendingComponents.returns(Object.values(state.charts));
const asyncEventMiddleware = initAsyncEvents({
config,
getPendingComponents,
successAction,
errorAction,
afterEach(() => {
WS.clean();
});
asyncEventMiddleware(mockStore)(next)(action);
expect(getPendingComponents.called).toBe(false);
it('resolves with chart data on event done status', async () => {
await wsServer.connected;
const promise = asyncEvent.waitForAsyncData(asyncPendingEvent);
wsServer.send(JSON.stringify(asyncDoneEvent));
await expect(promise).resolves.toEqual([chartData]);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(1);
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(0);
});
it('rejects on event error status', async () => {
await wsServer.connected;
const promise = asyncEvent.waitForAsyncData(asyncPendingEvent);
wsServer.send(JSON.stringify(asyncErrorEvent));
const errorResponse = await parseErrorJson(asyncErrorEvent);
await expect(promise).rejects.toEqual(errorResponse);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(0);
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(0);
});
it('rejects on cached data fetch error', async () => {
fetchMock.reset();
fetchMock.get(CACHED_DATA_ENDPOINT, {
status: 400,
});
await wsServer.connected;
const promise = asyncEvent.waitForAsyncData(asyncPendingEvent);
wsServer.send(JSON.stringify(asyncDoneEvent));
const errorResponse = [{ error: 'Bad Request' }];
await expect(promise).rejects.toEqual(errorResponse);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(1);
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(0);
});
it('resolves when events are received before listener', async () => {
await wsServer.connected;
wsServer.send(JSON.stringify(asyncDoneEvent));
const promise = asyncEvent.waitForAsyncData(asyncPendingEvent);
await expect(promise).resolves.toEqual([chartData]);
expect(fetchMock.calls(CACHED_DATA_ENDPOINT)).toHaveLength(1);
expect(fetchMock.calls(EVENTS_ENDPOINT)).toHaveLength(0);
});
});
});