diff --git a/plugins/scaffolder-backend/package.json b/plugins/scaffolder-backend/package.json index 8c61ee65c4..5d23ea0aff 100644 --- a/plugins/scaffolder-backend/package.json +++ b/plugins/scaffolder-backend/package.json @@ -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", diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts index 639d61d4fd..52fd175dc8 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts @@ -125,7 +125,7 @@ describe('StorageTaskBroker', () => { const logPromise = new Promise(resolve => { const observedEvents = new Array(); - 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(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', diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 71653f503b..0a1fc31bc8 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -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 { diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index a28ccce239..57b92e0222 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -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; vacuumTasks(options: { timeoutS: number }): Promise; - 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; } diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index b20dfa5d24..f8336179eb 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -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; ( - taskBroker.observe as jest.Mocked['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['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; ( - taskBroker.observe as jest.Mocked['observe'] - ).mockImplementation(({ taskId }, callback) => { - setImmediate(() => { - callback(undefined, { - events: [ - { - id: 1, - taskId, - type: 'completion', - createdAt: '', - body: { message: 'Finished!' }, - }, - ], + taskBroker.event$ as jest.Mocked['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; ( - taskBroker.observe as jest.Mocked['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['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; ( - taskBroker.observe as jest.Mocked['observe'] - ).mockImplementation((_, callback) => { - callback(undefined, { events: [] }); - return { unsubscribe }; + taskBroker.event$ as jest.Mocked['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); }); }); }); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 50a1e4bba0..b36670a340 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -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); }); });