Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 15 additions & 34 deletions src/domain/entities/generic/Future.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,10 @@
import {
FutureWithAccumulation,
ParallelAccumulatedData,
ParallelWithAccumulationOptions,
SequentialAccumulatedData,
SequentialWithAccumulationOptions,
} from "./FutureWithAccumulation";
import * as rcpromise from "real-cancellable-promise";

/**
Expand Down Expand Up @@ -169,38 +176,16 @@ export class Future<E, D> {

static sequentialWithAccumulation<E, D>(
futures: Array<Future<E, D>>,
options: { stopOnError?: boolean } = {}
options: SequentialWithAccumulationOptions = {}
): Future<never, SequentialAccumulatedData<E, D>> {
const { stopOnError = false } = options;
const processSequentially = (
futures: Array<Future<E, D>>,
accumulatedData: D[] = []
): Future<never, SequentialAccumulatedData<E, D>> => {
const [firstFuture, ...remainingFutures] = futures;

if (!firstFuture) {
return Future.success({ type: "success", data: accumulatedData });
}

return firstFuture
.flatMap(resultData => {
return processSequentially(remainingFutures, [...accumulatedData, resultData]);
})
.flatMapError((error: E) => {
if (stopOnError) {
const accumulatedDataWithError: SequentialAccumulatedData<E, D> = {
type: "error",
error: error,
data: accumulatedData,
};
return Future.success(accumulatedDataWithError);
} else {
return processSequentially(remainingFutures, accumulatedData);
}
});
};
return FutureWithAccumulation.sequential(futures, options);
}

return processSequentially(futures);
static parallelWithAccumulation<E, D>(
futures: Array<Future<E, D>>,
options: ParallelWithAccumulationOptions = {}
): Future<never, ParallelAccumulatedData<E, D>> {
return FutureWithAccumulation.parallel(futures, options);
}

static fromPromise<Data>(promise: Promise<Data>): FutureData<Data> {
Expand All @@ -211,10 +196,6 @@ export class Future<E, D> {
}
}

export type SequentialAccumulatedData<E, D> =
| { type: "success"; data: D[] }
| { type: "error"; error: E; data: D[] };

export type Cancel = (() => void) | undefined;

interface CaptureAsync<E> {
Expand Down
110 changes: 110 additions & 0 deletions src/domain/entities/generic/FutureWithAccumulation.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
import { Future } from "./Future";

export class FutureWithAccumulation {
static sequential<E, D>(
futures: Array<Future<E, D>>,
options: SequentialWithAccumulationOptions = {}
): Future<never, SequentialAccumulatedData<E, D>> {
const { stopOnError = false } = options;
const processSequentially = (
futures: Array<Future<E, D>>,
accumulatedData: D[] = []
): Future<never, SequentialAccumulatedData<E, D>> => {
const [firstFuture, ...remainingFutures] = futures;

if (!firstFuture) {
return Future.success({ type: "success", data: accumulatedData });
}

return firstFuture
.flatMap(resultData => {
return processSequentially(remainingFutures, [...accumulatedData, resultData]);
})
.flatMapError((error: E) => {
if (stopOnError) {
const accumulatedDataWithError: SequentialAccumulatedData<E, D> = {
type: "error",
error: error,
data: accumulatedData,
};
return Future.success(accumulatedDataWithError);
} else {
return processSequentially(remainingFutures, accumulatedData);
}
});
};

return processSequentially(futures);
}

static parallel<E, D>(
futures: Array<Future<E, D>>,
options: ParallelWithAccumulationOptions = {}
): Future<never, ParallelAccumulatedData<E, D>> {
const { concurrency = 10, stopOnError = true } = options;

const toParallelResult = (future: Future<E, D>): Future<never, ParallelResult<E, D>> => {
return future
.map<ParallelResult<E, D>>(data => ({ type: "success", data }))
.mapError<ParallelResult<E, D>>(error => ({ type: "error", error }))
.flatMapError(errorResult =>
Future.success<never, ParallelResult<E, D>>(errorResult)
);
};

const processInParallel = (
pendingFutures: Array<Future<E, D>>,
accumulatedData: D[] = []
): Future<never, ParallelAccumulatedData<E, D>> => {
if (pendingFutures.length === 0) {
return Future.success({ type: "success", data: accumulatedData });
}

const currentBatch = pendingFutures.slice(0, concurrency);
const remainingFutures = pendingFutures.slice(concurrency);

return Future.parallel(currentBatch.map(toParallelResult), {
concurrency: concurrency,
}).flatMap(batchResults => {
const successfulData = batchResults.flatMap(result =>
result.type === "success" ? [result.data] : []
);

const batchErrors = batchResults.flatMap(result =>
result.type === "error" ? [result.error] : []
);

const nextAccumulatedData = [...accumulatedData, ...successfulData];

if (batchErrors.length > 0 && stopOnError) {
return Future.success({
type: "error",
errors: batchErrors,
data: nextAccumulatedData,
});
}

return processInParallel(remainingFutures, nextAccumulatedData);
});
};

return processInParallel(futures);
}
}

export type ParallelWithAccumulationOptions = {
concurrency?: number;
stopOnError?: boolean;
};

export type ParallelAccumulatedData<E, D> =
| { type: "success"; data: D[] }
| { type: "error"; errors: E[]; data: D[] };

type ParallelResult<E, D> = { type: "success"; data: D } | { type: "error"; error: E };

export type SequentialWithAccumulationOptions = { stopOnError?: boolean };

export type SequentialAccumulatedData<E, D> =
| { type: "success"; data: D[] }
| { type: "error"; error: E; data: D[] };
Loading
Loading