56 lines
1.5 KiB
TypeScript
56 lines
1.5 KiB
TypeScript
import type {
|
|
HousekeepingCapabilityContext,
|
|
HousekeepingErrorCode,
|
|
} from "../contracts";
|
|
import { createCorrelationId } from "../correlation";
|
|
import { type HousekeepingProvider, runProvider } from "./run-provider";
|
|
|
|
export interface ProviderPolicy<T> {
|
|
timeoutMs: number;
|
|
perProviderLimit: number;
|
|
combinedLimit: number;
|
|
sort(items: readonly T[]): readonly T[];
|
|
dedupeKey(item: T): string;
|
|
}
|
|
|
|
export interface ProviderBatchResult<T> {
|
|
items: readonly T[];
|
|
errors: readonly { providerId: string; code: HousekeepingErrorCode }[];
|
|
correlationId: string;
|
|
}
|
|
|
|
export async function orchestrateProviders<T>(
|
|
providers: readonly HousekeepingProvider<T>[],
|
|
context: HousekeepingCapabilityContext,
|
|
policy: ProviderPolicy<T>,
|
|
): Promise<ProviderBatchResult<T>> {
|
|
const providerResults = await Promise.all(
|
|
providers.map((provider) =>
|
|
runProvider(provider, context, policy.timeoutMs),
|
|
),
|
|
);
|
|
const items: T[] = [];
|
|
const seen = new Set<string>();
|
|
const errors: { providerId: string; code: HousekeepingErrorCode }[] = [];
|
|
|
|
for (const result of providerResults) {
|
|
if (result.error !== undefined) {
|
|
errors.push({ providerId: result.providerId, code: result.error });
|
|
}
|
|
|
|
for (const item of result.items.slice(0, policy.perProviderLimit)) {
|
|
if (items.length >= policy.combinedLimit) break;
|
|
const key = policy.dedupeKey(item);
|
|
if (seen.has(key)) continue;
|
|
seen.add(key);
|
|
items.push(item);
|
|
}
|
|
}
|
|
|
|
return {
|
|
items: policy.sort(items),
|
|
errors,
|
|
correlationId: createCorrelationId(),
|
|
};
|
|
}
|