chore: remove the observe method in favour of an observable

Signed-off-by: blam <ben@blam.sh>
This commit is contained in:
blam
2022-02-22 14:31:20 +01:00
parent 9d9b2bab47
commit 67771d2b21
6 changed files with 155 additions and 175 deletions
+3 -1
View File
@@ -72,7 +72,8 @@
"uuid": "^8.2.0",
"winston": "^3.2.1",
"yaml": "^1.10.0",
"vm2": "^3.9.6"
"vm2": "^3.9.6",
"zen-observable": "^0.8.15"
},
"devDependencies": {
"@backstage/cli": "^0.14.0",
@@ -83,6 +84,7 @@
"@types/mock-fs": "^4.13.0",
"@types/nunjucks": "^3.1.4",
"@types/supertest": "^2.0.8",
"@types/zen-observable": "^0.8.0",
"esbuild": "^0.14.1",
"jest-when": "^3.1.0",
"mock-fs": "^5.1.0",
@@ -125,7 +125,7 @@ describe('StorageTaskBroker', () => {
const logPromise = new Promise<SerializedTaskEvent[]>(resolve => {
const observedEvents = new Array<SerializedTaskEvent>();
broker2.observe({ taskId, after: undefined }, (_err, { events }) => {
broker2.event$({ taskId, after: undefined }).subscribe(({ events }) => {
observedEvents.push(...events);
if (events.some(e => e.type === 'completion')) {
resolve(observedEvents);
@@ -147,9 +147,11 @@ describe('StorageTaskBroker', () => {
]);
const afterLogs = await new Promise<string[]>(resolve => {
broker2.observe({ taskId, after: logs[1].id }, (_err, { events }) =>
resolve(events.map(e => e.body.message as string)),
);
broker2
.event$({ taskId, after: logs[1].id })
.subscribe(({ events }) =>
resolve(events.map(e => e.body.message as string)),
);
});
expect(afterLogs).toEqual([
'log 3',
@@ -13,7 +13,8 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
import { JsonObject } from '@backstage/types';
import { JsonObject, Observable } from '@backstage/types';
import ObservableImpl from 'zen-observable';
import { assertError } from '@backstage/errors';
import { TaskSpec } from '@backstage/plugin-scaffolder-common';
import { Logger } from 'winston';
@@ -176,43 +177,33 @@ export class StorageTaskBroker implements TaskBroker {
return this.storage.getTask(taskId);
}
observe(
options: {
taskId: string;
after: number | undefined;
},
callback: (
error: Error | undefined,
result: { events: SerializedTaskEvent[] },
) => void,
): { unsubscribe: () => void } {
const { taskId } = options;
event$(options: {
taskId: string;
after?: number;
}): Observable<{ events: SerializedTaskEvent[] }> {
return new ObservableImpl(observer => {
const { taskId } = options;
let cancelled = false;
const unsubscribe = () => {
cancelled = true;
};
(async () => {
let after = options.after;
while (!cancelled) {
const result = await this.storage.listEvents({ taskId, after: after });
const { events } = result;
if (events.length) {
after = events[events.length - 1].id;
try {
callback(undefined, result);
} catch (error) {
assertError(error);
callback(error, { events: [] });
let cancelled = false;
(async () => {
while (!cancelled) {
const result = await this.storage.listEvents({ taskId, after });
const { events } = result;
if (events.length) {
after = events[events.length - 1].id;
observer.next(result);
}
await new Promise(resolve => setTimeout(resolve, 1000));
}
})();
await new Promise(resolve => setTimeout(resolve, 1000));
}
})();
return { unsubscribe };
return () => {
cancelled = true;
};
});
}
async vacuumTasks(options: { timeoutS: number }): Promise<void> {
@@ -14,7 +14,7 @@
* limitations under the License.
*/
import { JsonValue, JsonObject } from '@backstage/types';
import { JsonValue, JsonObject, Observable } from '@backstage/types';
import { TaskSpec } from '@backstage/plugin-scaffolder-common';
/**
@@ -148,16 +148,10 @@ export interface TaskBroker {
options: TaskBrokerDispatchOptions,
): Promise<TaskBrokerDispatchResult>;
vacuumTasks(options: { timeoutS: number }): Promise<void>;
observe(
options: {
taskId: string;
after: number | undefined;
},
callback: (
error: Error | undefined,
result: { events: SerializedTaskEvent[] },
) => void,
): { unsubscribe: () => void };
event$(options: {
taskId: string;
after: number | undefined;
}): Observable<{ events: SerializedTaskEvent[] }>;
get(taskId: string): Promise<SerializedTask>;
}
@@ -38,6 +38,7 @@ import {
import { CatalogApi } from '@backstage/catalog-client';
import { TemplateEntityV1beta2 } from '@backstage/plugin-scaffolder-common';
import { ConfigReader } from '@backstage/config';
import ObservableImpl from 'zen-observable';
import express from 'express';
import request from 'supertest';
/**
@@ -113,7 +114,7 @@ describe('createRouter', () => {
jest.spyOn(taskBroker, 'dispatch');
jest.spyOn(taskBroker, 'get');
jest.spyOn(taskBroker, 'observe');
jest.spyOn(taskBroker, 'event$');
const router = await createRouter({
logger: getVoidLogger(),
@@ -194,37 +195,38 @@ describe('createRouter', () => {
describe('GET /v2/tasks/:taskId/eventstream', () => {
it('should return log messages', async () => {
const unsubscribe = jest.fn();
let subscriber: ZenObservable.SubscriptionObserver<any>;
(
taskBroker.observe as jest.Mocked<TaskBroker>['observe']
).mockImplementation(({ taskId }, callback) => {
// emit after this function returned
setImmediate(() => {
callback(undefined, {
events: [
{
id: 0,
taskId,
type: 'log',
createdAt: '',
body: { message: 'My log message' },
},
],
});
callback(undefined, {
events: [
{
id: 1,
taskId,
type: 'completion',
createdAt: '',
body: { message: 'Finished!' },
},
],
taskBroker.event$ as jest.Mocked<TaskBroker>['event$']
).mockImplementation(({ taskId }) => {
return new ObservableImpl(observer => {
subscriber = observer;
setImmediate(() => {
observer.next({
events: [
{
id: 0,
taskId,
type: 'log',
createdAt: '',
body: { message: 'My log message' },
},
],
});
observer.next({
events: [
{
id: 1,
taskId,
type: 'completion',
createdAt: '',
body: { message: 'Finished!' },
},
],
});
});
});
return { unsubscribe };
// emit after this function returned
});
let statusCode: any = undefined;
@@ -264,34 +266,32 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{
`);
expect(taskBroker.observe).toBeCalledTimes(1);
expect(taskBroker.observe).toBeCalledWith(
{ taskId: 'a-random-id' },
expect.any(Function),
);
expect(unsubscribe).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledWith({ taskId: 'a-random-id' });
expect(subscriber!.closed).toBe(true);
});
it('should return log messages with after query', async () => {
const unsubscribe = jest.fn();
let subscriber: ZenObservable.SubscriptionObserver<any>;
(
taskBroker.observe as jest.Mocked<TaskBroker>['observe']
).mockImplementation(({ taskId }, callback) => {
setImmediate(() => {
callback(undefined, {
events: [
{
id: 1,
taskId,
type: 'completion',
createdAt: '',
body: { message: 'Finished!' },
},
],
taskBroker.event$ as jest.Mocked<TaskBroker>['event$']
).mockImplementation(({ taskId }) => {
return new ObservableImpl(observer => {
subscriber = observer;
setImmediate(() => {
observer.next({
events: [
{
id: 1,
taskId,
type: 'completion',
createdAt: '',
body: { message: 'Finished!' },
},
],
});
});
});
return { unsubscribe };
});
let statusCode: any = undefined;
@@ -318,41 +318,43 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{
expect(statusCode).toBe(200);
expect(headers['content-type']).toBe('text/event-stream');
expect(taskBroker.observe).toBeCalledTimes(1);
expect(taskBroker.observe).toBeCalledWith(
{ taskId: 'a-random-id', after: 10 },
expect.any(Function),
);
expect(taskBroker.event$).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledWith({
taskId: 'a-random-id',
after: 10,
});
expect(unsubscribe).toBeCalledTimes(1);
expect(subscriber!.closed).toBe(true);
});
});
describe('GET /v2/tasks/:taskId/events', () => {
it('should return log messages', async () => {
const unsubscribe = jest.fn();
let subscriber: ZenObservable.SubscriptionObserver<any>;
(
taskBroker.observe as jest.Mocked<TaskBroker>['observe']
).mockImplementation(({ taskId }, callback) => {
callback(undefined, {
events: [
{
id: 0,
taskId,
type: 'log',
createdAt: '',
body: { message: 'My log message' },
},
{
id: 1,
taskId,
type: 'completion',
createdAt: '',
body: { message: 'Finished!' },
},
],
taskBroker.event$ as jest.Mocked<TaskBroker>['event$']
).mockImplementation(({ taskId }) => {
return new ObservableImpl(observer => {
subscriber = observer;
observer.next({
events: [
{
id: 0,
taskId,
type: 'log',
createdAt: '',
body: { message: 'My log message' },
},
{
id: 1,
taskId,
type: 'completion',
createdAt: '',
body: { message: 'Finished!' },
},
],
});
});
return { unsubscribe };
});
const response = await request(app).get('/v2/tasks/a-random-id/events');
@@ -375,21 +377,20 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{
},
]);
expect(taskBroker.observe).toBeCalledTimes(1);
expect(taskBroker.observe).toBeCalledWith(
{ taskId: 'a-random-id' },
expect.any(Function),
);
expect(unsubscribe).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledWith({ taskId: 'a-random-id' });
expect(subscriber!.closed).toBe(true);
});
it('should return log messages with after query', async () => {
const unsubscribe = jest.fn();
let subscriber: ZenObservable.SubscriptionObserver<any>;
(
taskBroker.observe as jest.Mocked<TaskBroker>['observe']
).mockImplementation((_, callback) => {
callback(undefined, { events: [] });
return { unsubscribe };
taskBroker.event$ as jest.Mocked<TaskBroker>['event$']
).mockImplementation(() => {
return new ObservableImpl(observer => {
subscriber = observer;
observer.next({ events: [] });
});
});
const response = await request(app)
@@ -399,12 +400,12 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{
expect(response.status).toEqual(200);
expect(response.body).toEqual([]);
expect(taskBroker.observe).toBeCalledTimes(1);
expect(taskBroker.observe).toBeCalledWith(
{ taskId: 'a-random-id', after: 10 },
expect.any(Function),
);
expect(unsubscribe).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledTimes(1);
expect(taskBroker.event$).toBeCalledWith({
taskId: 'a-random-id',
after: 10,
});
expect(subscriber!.closed).toBe(true);
});
});
});
@@ -301,15 +301,13 @@ export async function createRouter(
});
// After client opens connection send all events as string
const { unsubscribe } = taskBroker.observe(
{ taskId, after },
(error, { events }) => {
if (error) {
logger.error(
`Received error from event stream when observing taskId '${taskId}', ${error}`,
);
}
const subscription = taskBroker.event$({ taskId, after }).subscribe({
error: error => {
logger.error(
`Received error from event stream when observing taskId '${taskId}', ${error}`,
);
},
next: ({ events }) => {
let shouldUnsubscribe = false;
for (const event of events) {
res.write(
@@ -317,19 +315,18 @@ export async function createRouter(
);
if (event.type === 'completion') {
shouldUnsubscribe = true;
// Closing the event stream here would cause the frontend
// to automatically reconnect because it lost connection.
}
}
// res.flush() is only available with the compression middleware
res.flush?.();
if (shouldUnsubscribe) unsubscribe();
if (shouldUnsubscribe) subscription.unsubscribe();
},
);
});
// When client closes connection we update the clients list
// avoiding the disconnected one
req.on('close', () => {
unsubscribe();
subscription.unsubscribe();
logger.debug(`Event stream observing taskId '${taskId}' closed`);
});
})
@@ -337,36 +334,29 @@ export async function createRouter(
const { taskId } = req.params;
const after = Number(req.query.after) || undefined;
let unsubscribe = () => {};
// cancel the request after 30 seconds. this aligns with the recommendations of RFC 6202.
const timeout = setTimeout(() => {
unsubscribe();
res.json([]);
}, 30_000);
// Get all known events after an id (always includes the completion event) and return the first callback
({ unsubscribe } = taskBroker.observe(
{ taskId, after },
(error, { events }) => {
// stop the timeout
const subscription = taskBroker.event$({ taskId, after }).subscribe({
error: error => {
logger.error(
`Received error from event stream when observing taskId '${taskId}', ${error}`,
);
},
next: ({ events }) => {
clearTimeout(timeout);
unsubscribe();
if (error) {
logger.error(
`Received error from log when observing taskId '${taskId}', ${error}`,
);
}
subscription.unsubscribe();
res.json(events);
},
));
});
// When client closes connection we update the clients list
// avoiding the disconnected one
req.on('close', () => {
unsubscribe();
subscription.unsubscribe();
clearTimeout(timeout);
});
});