From 4c399873d13d4a201264bb4e3902774f5570ae02 Mon Sep 17 00:00:00 2001 From: simoleo89 Date: Thu, 27 Aug 2026 18:02:23 +0200 Subject: [PATCH] feat(housekeeping): orchestrate partial providers --- .../foundation/providers/orchestrate.test.ts | 111 ++++++++++++++++++ .../foundation/providers/orchestrate.ts | 55 +++++++++ .../foundation/providers/run-provider.test.ts | 78 ++++++++++++ .../foundation/providers/run-provider.ts | 62 ++++++++++ 4 files changed, 306 insertions(+) create mode 100644 src/features/housekeeping/foundation/providers/orchestrate.test.ts create mode 100644 src/features/housekeeping/foundation/providers/orchestrate.ts create mode 100644 src/features/housekeeping/foundation/providers/run-provider.test.ts create mode 100644 src/features/housekeeping/foundation/providers/run-provider.ts diff --git a/src/features/housekeeping/foundation/providers/orchestrate.test.ts b/src/features/housekeeping/foundation/providers/orchestrate.test.ts new file mode 100644 index 00000000..41e9bd53 --- /dev/null +++ b/src/features/housekeeping/foundation/providers/orchestrate.test.ts @@ -0,0 +1,111 @@ +import { describe, expect, it, vi } from "vitest"; +import { + anyCapability, + type HousekeepingCapabilityContext, + ok, +} from "../contracts"; +import { orchestrateProviders, type ProviderPolicy } from "./orchestrate"; +import type { HousekeepingProvider } from "./run-provider"; + +interface Item { + key: string; + label: string; + order: number; +} + +const policy: ProviderPolicy = { + timeoutMs: 2_000, + perProviderLimit: 2, + combinedLimit: 3, + sort: (items) => [...items].sort((left, right) => left.order - right.order), + dedupeKey: (item) => item.key, +}; + +const allowedContext: HousekeepingCapabilityContext = { + actor: { id: 42, username: "operator", rank: 0 }, + isSuperAdmin: false, + has: (slug) => slug === "housekeeping.read", + hasAny: (...slugs) => slugs.includes("housekeeping.read"), + hasAll: (...slugs) => slugs.every((slug) => slug === "housekeeping.read"), +}; + +function provider( + id: string, + items: readonly Item[], +): HousekeepingProvider { + return { + id, + capability: anyCapability("housekeeping.read"), + run: async () => ok(items, `${id}-correlation`), + }; +} + +describe("orchestrateProviders", () => { + it("preserves a successful sibling when another provider times out", async () => { + vi.useFakeTimers(); + const slow: HousekeepingProvider = { + id: "slow", + capability: anyCapability("housekeeping.read"), + run: async (_context, signal) => + await new Promise((_, reject) => { + signal.addEventListener("abort", () => reject(new Error("aborted"))); + }), + }; + + const result = orchestrateProviders( + [slow, provider("fast", [{ key: "fast", label: "Fast", order: 1 }])], + allowedContext, + policy, + ); + await vi.advanceTimersByTimeAsync(2_000); + + expect(await result).toMatchObject({ + items: [{ key: "fast", label: "Fast", order: 1 }], + errors: [{ providerId: "slow", code: "TIMEOUT" }], + }); + vi.useRealTimers(); + }); + + it("keeps the first duplicate before applying deterministic sorting", async () => { + const result = await orchestrateProviders( + [ + provider("first", [ + { key: "shared", label: "First", order: 2 }, + { key: "unique", label: "Unique", order: 1 }, + ]), + provider("second", [{ key: "shared", label: "Second", order: 0 }]), + ], + allowedContext, + policy, + ); + + expect(result.items).toEqual([ + { key: "unique", label: "Unique", order: 1 }, + { key: "shared", label: "First", order: 2 }, + ]); + }); + + it("enforces per-provider and combined caps before sorting", async () => { + const result = await orchestrateProviders( + [ + provider("first", [ + { key: "one", label: "One", order: 3 }, + { key: "two", label: "Two", order: 2 }, + { key: "discarded", label: "Discarded", order: 0 }, + ]), + provider("second", [ + { key: "three", label: "Three", order: 1 }, + { key: "combined-discarded", label: "Combined discarded", order: 0 }, + ]), + ], + allowedContext, + policy, + ); + + expect(result.items).toEqual([ + { key: "three", label: "Three", order: 1 }, + { key: "two", label: "Two", order: 2 }, + { key: "one", label: "One", order: 3 }, + ]); + }); +}); diff --git a/src/features/housekeeping/foundation/providers/orchestrate.ts b/src/features/housekeeping/foundation/providers/orchestrate.ts new file mode 100644 index 00000000..fb52c181 --- /dev/null +++ b/src/features/housekeeping/foundation/providers/orchestrate.ts @@ -0,0 +1,55 @@ +import type { + HousekeepingCapabilityContext, + HousekeepingErrorCode, +} from "../contracts"; +import { createCorrelationId } from "../correlation"; +import { type HousekeepingProvider, runProvider } from "./run-provider"; + +export interface ProviderPolicy { + timeoutMs: number; + perProviderLimit: number; + combinedLimit: number; + sort(items: readonly T[]): readonly T[]; + dedupeKey(item: T): string; +} + +export interface ProviderBatchResult { + items: readonly T[]; + errors: readonly { providerId: string; code: HousekeepingErrorCode }[]; + correlationId: string; +} + +export async function orchestrateProviders( + providers: readonly HousekeepingProvider[], + context: HousekeepingCapabilityContext, + policy: ProviderPolicy, +): Promise> { + const providerResults = await Promise.all( + providers.map((provider) => + runProvider(provider, context, policy.timeoutMs), + ), + ); + const items: T[] = []; + const seen = new Set(); + 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(), + }; +} diff --git a/src/features/housekeeping/foundation/providers/run-provider.test.ts b/src/features/housekeeping/foundation/providers/run-provider.test.ts new file mode 100644 index 00000000..d72c6b9e --- /dev/null +++ b/src/features/housekeeping/foundation/providers/run-provider.test.ts @@ -0,0 +1,78 @@ +import { describe, expect, it, vi } from "vitest"; +import { + anyCapability, + type HousekeepingCapabilityContext, + ok, +} from "../contracts"; +import { type HousekeepingProvider, runProvider } from "./run-provider"; + +function context(granted: readonly string[]): HousekeepingCapabilityContext { + return { + actor: { id: 42, username: "operator", rank: 999 }, + isSuperAdmin: false, + has: (slug) => granted.includes(slug), + hasAny: (...slugs) => slugs.some((slug) => granted.includes(slug)), + hasAll: (...slugs) => slugs.every((slug) => granted.includes(slug)), + }; +} + +describe("runProvider", () => { + it("aborts a slow provider at its 2,000 ms deadline and returns TIMEOUT", async () => { + vi.useFakeTimers(); + let observedAbort = false; + const provider: HousekeepingProvider = { + id: "slow", + capability: anyCapability("housekeeping.read"), + run: async (_context, signal) => + await new Promise((_, reject) => { + signal.addEventListener("abort", () => { + observedAbort = true; + reject(new Error("aborted")); + }); + }), + }; + + const result = runProvider(provider, context(["housekeeping.read"]), 2_000); + await vi.advanceTimersByTimeAsync(2_000); + + expect(await result).toEqual({ + providerId: "slow", + items: [], + error: "TIMEOUT", + }); + expect(observedAbort).toBe(true); + vi.useRealTimers(); + }); + + it("skips a provider whose required capability is absent", async () => { + let invoked = false; + const provider: HousekeepingProvider = { + id: "restricted", + capability: anyCapability("housekeeping.read"), + run: async () => { + invoked = true; + return ok(["private"], "provider-correlation"); + }, + }; + + expect(await runProvider(provider, context([]), 2_000)).toEqual({ + providerId: "restricted", + items: [], + }); + expect(invoked).toBe(false); + }); + + it("maps one thrown provider failure to one INTERNAL partial error", async () => { + const provider: HousekeepingProvider = { + id: "broken", + capability: anyCapability("housekeeping.read"), + run: async () => { + throw new Error("database unavailable"); + }, + }; + + expect( + await runProvider(provider, context(["housekeeping.read"]), 2_000), + ).toEqual({ providerId: "broken", items: [], error: "INTERNAL" }); + }); +}); diff --git a/src/features/housekeeping/foundation/providers/run-provider.ts b/src/features/housekeeping/foundation/providers/run-provider.ts new file mode 100644 index 00000000..f64ae06a --- /dev/null +++ b/src/features/housekeeping/foundation/providers/run-provider.ts @@ -0,0 +1,62 @@ +import { satisfiesCapability } from "../capability-context"; +import type { + CapabilityRequirement, + HousekeepingCapabilityContext, + HousekeepingErrorCode, + HousekeepingResult, +} from "../contracts"; + +export interface HousekeepingProvider { + id: string; + capability: CapabilityRequirement; + run( + context: HousekeepingCapabilityContext, + signal: AbortSignal, + ): Promise>; +} + +export interface ProviderRunResult { + providerId: string; + items: readonly T[]; + error?: HousekeepingErrorCode; +} + +export async function runProvider( + provider: HousekeepingProvider, + context: HousekeepingCapabilityContext, + timeoutMs: number, +): Promise> { + if (!satisfiesCapability(context, provider.capability)) { + return { providerId: provider.id, items: [] }; + } + + const controller = new AbortController(); + let timeout: ReturnType | undefined; + const run = provider + .run(context, controller.signal) + .then( + (result): ProviderRunResult => + result.ok + ? { providerId: provider.id, items: result.data } + : { providerId: provider.id, items: [], error: result.error.code }, + ) + .catch( + (): ProviderRunResult => ({ + providerId: provider.id, + items: [], + error: "INTERNAL", + }), + ); + const timeoutResult = new Promise>((resolve) => { + timeout = setTimeout(() => { + controller.abort(); + resolve({ providerId: provider.id, items: [], error: "TIMEOUT" }); + }, timeoutMs); + }); + + try { + return await Promise.race([run, timeoutResult]); + } finally { + if (timeout !== undefined) clearTimeout(timeout); + } +}