From 799824d581f182c129f81d030bd024cfb9f76cc2 Mon Sep 17 00:00:00 2001 From: Wroud Date: Fri, 16 Jul 2021 22:43:56 +0300 Subject: [PATCH] feat: CB-1093 cancellable promise implementation --- .../core-executor/src/TaskScheduler/ITask.ts | 18 ++- .../core-executor/src/TaskScheduler/Task.ts | 121 ++++++++++++++++++ .../src/TaskScheduler/TaskScheduler.ts | 110 ++++++++++------ .../core-sdk/src/Resource/CachedResource.ts | 16 ++- 4 files changed, 214 insertions(+), 51 deletions(-) create mode 100644 webapp/packages/core-executor/src/TaskScheduler/Task.ts diff --git a/webapp/packages/core-executor/src/TaskScheduler/ITask.ts b/webapp/packages/core-executor/src/TaskScheduler/ITask.ts index 19a3d9f1e1..208876f932 100644 --- a/webapp/packages/core-executor/src/TaskScheduler/ITask.ts +++ b/webapp/packages/core-executor/src/TaskScheduler/ITask.ts @@ -6,7 +6,19 @@ * you may not use this file except in compliance with the License. */ -export interface ITask { - readonly id: T; - readonly task: Promise; +export interface ITask extends Promise { + readonly cancelled: boolean; + readonly executing: boolean; + readonly cancellable: boolean; + readonly run: () => this; + readonly cancel: () => Promise | void; + + then: ( + onfulfilled?: ((value: TValue) => TResult1 | PromiseLike) | null, + onrejected?: ((reason: any) => TResult2 | PromiseLike) | null + ) => ITask; + catch: ( + onrejected?: ((reason: any) => TResult | PromiseLike) | null + ) => ITask; + finally: (onfinally?: (() => void) | null) => ITask; } diff --git a/webapp/packages/core-executor/src/TaskScheduler/Task.ts b/webapp/packages/core-executor/src/TaskScheduler/Task.ts new file mode 100644 index 0000000000..080019bde6 --- /dev/null +++ b/webapp/packages/core-executor/src/TaskScheduler/Task.ts @@ -0,0 +1,121 @@ +/* + * CloudBeaver - Cloud Database Manager + * Copyright (C) 2020-2021 DBeaver Corp and others + * + * Licensed under the Apache License, Version 2.0. + * you may not use this file except in compliance with the License. + */ + +import { makeObservable, observable } from 'mobx'; + +import type { ITask } from './ITask'; + +export class Task implements ITask { + cancelled: boolean; + executing: boolean; + + get cancellable(): boolean { + return !this.cancelled && (!this.executing || !!this.externalCancel); + } + + private resolve!: (value: TValue) => void; + private reject!: (reason?: any) => void; + private innerPromise: Promise; + + get [Symbol.toStringTag](): string { + return 'Task'; + } + + constructor( + readonly task: () => Promise, + private externalCancel?: () => Promise | void + ) { + this.innerPromise = new Promise((resolve, reject) => { + this.reject = reject; + this.resolve = resolve; + }); + this.cancelled = false; + this.executing = false; + + makeObservable(this, { + cancelled: observable, + executing: observable, + }); + } + + then( + onfulfilled?: ((value: TValue) => TResult1 | PromiseLike) | null, + onrejected?: ((reason: any) => TResult2 | PromiseLike) | null + ): ITask { + return new Task(async () => { + try { + const value = await this.innerPromise; + return await onfulfilled?.(value) as TResult1; + } catch (e) { + if (onrejected) { + return await onrejected(e); + } + throw e; + } + }, () => this.cancel()).run(); + } + + catch( + onrejected?: ((reason: any) => TResult | PromiseLike) | null + ): ITask { + return new Task(async () => { + try { + return await this.innerPromise; + } catch (exception) { + if (onrejected) { + return await onrejected(exception); + } + throw exception; + } + }, () => this.cancel()).run(); + } + + finally(onfinally?: (() => void) | null): ITask { + return new Task(async () => { + try { + return await this.innerPromise; + } finally { + onfinally?.(); + } + }, () => this.cancel()).run(); + } + + run(): this { + if (this.cancelled) { + return this; + } + + if (this.executing) { + throw new Error('Task already executing'); + } + + this.executing = true; + + this.task() + .then(value => this.resolve(value)) + .catch(reason => this.reject(reason)) + .finally(() => { + this.executing = false; + }); + + return this; + } + + cancel(): Promise | void { + this.cancelled = true; + + if (!this.executing) { + this.reject(new Error('Task was cancelled')); + return; + } + + if (this.externalCancel) { + return this.externalCancel(); + } + } +} diff --git a/webapp/packages/core-executor/src/TaskScheduler/TaskScheduler.ts b/webapp/packages/core-executor/src/TaskScheduler/TaskScheduler.ts index a9b860db49..ae3b5f0e5d 100644 --- a/webapp/packages/core-executor/src/TaskScheduler/TaskScheduler.ts +++ b/webapp/packages/core-executor/src/TaskScheduler/TaskScheduler.ts @@ -9,8 +9,20 @@ import { computed, observable, makeObservable } from 'mobx'; import type { ITask } from './ITask'; +import { Task } from './Task'; + +interface ITaskContainer { + readonly id: T; + task: ITask; +} export type BlockedExecution = (active: T, current: T) => boolean; +export interface IScheduleOptions { + cancel?: () => Promise | any; + after?: () => Promise | any; + success?: () => Promise | any; + error?: (exception: Error) => Promise | any; +} const queueLimit = 100; @@ -23,7 +35,7 @@ export class TaskScheduler { return this.queue.length > 0; } - private readonly queue: Array>; + private readonly queue: Array>; private readonly isBlocked: BlockedExecution | null; @@ -37,65 +49,79 @@ export class TaskScheduler { this.isBlocked = isBlocked; } - async schedule( + isExecuting(id: TIdentifier): boolean { + if (!this.isBlocked) { + return this.executing; + } + return this.queue.some(active => this.isBlocked!(active.id, id)); + } + + schedule( id: TIdentifier, promise: () => Promise, - after?: () => Promise | any, - success?: () => Promise | any, - error?: (exception: Error) => Promise | any, - ): Promise { - const task: ITask = { - id, - task: this.scheduler(id, promise), - }; - + options?: IScheduleOptions, + ): ITask { if (this.queue.length > queueLimit) { throw new Error('Execution queue limit is reached'); } - this.queue.push(task); - try { - const value = await task.task; - await success?.(); - return value; - } catch (exception) { - await error?.(exception); - throw exception; - } finally { - await after?.(); + const task = new Task(promise, options?.cancel); + const container: ITaskContainer = { id, task }; + this.queue.push(container); + + this.execute(container); + + return task + .then(async value => { + await options?.success?.(); + return value; + }) + .catch(async exception => { + await options?.error?.(exception); + throw exception; + }) + .finally(() => options?.after?.()) + .finally(() => this.queue.splice(this.queue.indexOf(container), 1)); + } + + async cancel(id: TIdentifier): Promise { + const containers = this.queue.filter(container => container.id === id && !container.task.cancelled); + + for (const container of containers) { + await container.task.cancel(); } } async wait(): Promise { const queueList = this.queue.slice(); - for (const task of queueList) { + for (const container of queueList) { try { - await task.task; + await container.task; } catch {} } } - private async scheduler( - id: TIdentifier, - promise: () => Promise, - ) { - try { - if (!this.isBlocked) { - return await promise(); + private async execute(container: ITaskContainer) { + if (this.isBlocked) { + const queueList = this.queue.filter( + active => + active !== container + && this.isBlocked!(active.id, container.id) + ); + + for (const _container of queueList) { + if (container.task.cancelled) { + throw new Error('Task was cancelled'); + } + if (!_container.task.cancelled) { + try { + await _container.task; + } catch {} + } } - - const queueList = this.queue.filter(active => this.isBlocked!(active.id, id)); - - for (const task of queueList) { - try { - await task.task; - } catch {} - } - - return await promise(); - } finally { - this.queue.splice(this.queue.findIndex(task => task.id === id), 1); } + + container.task.run(); } } diff --git a/webapp/packages/core-sdk/src/Resource/CachedResource.ts b/webapp/packages/core-sdk/src/Resource/CachedResource.ts index c5f50758cd..d461e92054 100644 --- a/webapp/packages/core-sdk/src/Resource/CachedResource.ts +++ b/webapp/packages/core-sdk/src/Resource/CachedResource.ts @@ -258,9 +258,11 @@ export abstract class CachedResource< return await this.taskWrapper(param, context, update); }, - () => this.markDataLoaded(param, context), - () => this.onDataUpdate.execute(param), - exception => this.markDataError(exception, param, context)); + { + after: () => this.markDataLoaded(param, context), + success: () => this.onDataUpdate.execute(param), + error: exception => this.markDataError(exception, param, context), + }); } protected async loadData(param: TParam, refresh: boolean, context: TContext): Promise { @@ -286,9 +288,11 @@ export abstract class CachedResource< await this.taskWrapper(param, context, this.loadingTask); }, - () => this.markDataLoaded(param, context), - () => this.onDataUpdate.execute(param), - exception => this.markDataError(exception, param, context)); + { + after: () => this.markDataLoaded(param, context), + success: () => this.onDataUpdate.execute(param), + error: exception => this.markDataError(exception, param, context), + }); } private async loadingTask(param: TParam, context: TContext) {