feat: CB-1093 cancellable promise implementation

This commit is contained in:
Wroud
2021-07-16 22:43:56 +03:00
parent f65cf434af
commit 799824d581
4 changed files with 214 additions and 51 deletions
@@ -6,7 +6,19 @@
* you may not use this file except in compliance with the License.
*/
export interface ITask<T> {
readonly id: T;
readonly task: Promise<any>;
export interface ITask<TValue> extends Promise<TValue> {
readonly cancelled: boolean;
readonly executing: boolean;
readonly cancellable: boolean;
readonly run: () => this;
readonly cancel: () => Promise<void> | void;
then: <TResult1 = TValue, TResult2 = never>(
onfulfilled?: ((value: TValue) => TResult1 | PromiseLike<TResult1>) | null,
onrejected?: ((reason: any) => TResult2 | PromiseLike<TResult2>) | null
) => ITask<TResult1 | TResult2>;
catch: <TResult = never>(
onrejected?: ((reason: any) => TResult | PromiseLike<TResult>) | null
) => ITask<TValue | TResult>;
finally: (onfinally?: (() => void) | null) => ITask<TValue>;
}
@@ -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<TValue> implements ITask<TValue> {
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<TValue>;
get [Symbol.toStringTag](): string {
return 'Task';
}
constructor(
readonly task: () => Promise<TValue>,
private externalCancel?: () => Promise<void> | 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<TResult1 = TValue, TResult2 = never>(
onfulfilled?: ((value: TValue) => TResult1 | PromiseLike<TResult1>) | null,
onrejected?: ((reason: any) => TResult2 | PromiseLike<TResult2>) | null
): ITask<TResult1 | TResult2> {
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<TResult = never>(
onrejected?: ((reason: any) => TResult | PromiseLike<TResult>) | null
): ITask<TValue | TResult> {
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<TValue> {
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> | void {
this.cancelled = true;
if (!this.executing) {
this.reject(new Error('Task was cancelled'));
return;
}
if (this.externalCancel) {
return this.externalCancel();
}
}
}
@@ -9,8 +9,20 @@
import { computed, observable, makeObservable } from 'mobx';
import type { ITask } from './ITask';
import { Task } from './Task';
interface ITaskContainer<T, TValue> {
readonly id: T;
task: ITask<TValue>;
}
export type BlockedExecution<T> = (active: T, current: T) => boolean;
export interface IScheduleOptions {
cancel?: () => Promise<any> | any;
after?: () => Promise<any> | any;
success?: () => Promise<any> | any;
error?: (exception: Error) => Promise<any> | any;
}
const queueLimit = 100;
@@ -23,7 +35,7 @@ export class TaskScheduler<TIdentifier> {
return this.queue.length > 0;
}
private readonly queue: Array<ITask<TIdentifier>>;
private readonly queue: Array<ITaskContainer<TIdentifier, any>>;
private readonly isBlocked: BlockedExecution<TIdentifier> | null;
@@ -37,65 +49,79 @@ export class TaskScheduler<TIdentifier> {
this.isBlocked = isBlocked;
}
async schedule<T>(
isExecuting(id: TIdentifier): boolean {
if (!this.isBlocked) {
return this.executing;
}
return this.queue.some(active => this.isBlocked!(active.id, id));
}
schedule<T>(
id: TIdentifier,
promise: () => Promise<T>,
after?: () => Promise<any> | any,
success?: () => Promise<any> | any,
error?: (exception: Error) => Promise<any> | any,
): Promise<T> {
const task: ITask<TIdentifier> = {
id,
task: this.scheduler(id, promise),
};
options?: IScheduleOptions,
): ITask<T> {
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<T>(promise, options?.cancel);
const container: ITaskContainer<TIdentifier, T> = { 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<void> {
const containers = this.queue.filter(container => container.id === id && !container.task.cancelled);
for (const container of containers) {
await container.task.cancel();
}
}
async wait(): Promise<void> {
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<T>(
id: TIdentifier,
promise: () => Promise<T>,
) {
try {
if (!this.isBlocked) {
return await promise();
private async execute<T>(container: ITaskContainer<TIdentifier, T>) {
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();
}
}
@@ -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<void> {
@@ -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) {