import { jsonFingerprint } from "./json.js"; import type { ControllerStore } from "./store.js"; import type { ControllerEffects, ControllerResource, EffectApplication, EffectDefinition, EffectObservation, EffectRecord, JsonObject, } from "./types.js"; export class ControllerEffectService implements ControllerEffects { private used = false; constructor( private readonly store: ControllerStore, private readonly resource: ControllerResource, private readonly signal: AbortSignal, ) {} async ensure(definition: EffectDefinition): Promise { if (this.used) { throw new Error("A reconciliation pass may ensure only one external effect"); } this.used = true; const requestFingerprint = jsonFingerprint(definition.request); const reservation = this.store.reserveEffect({ key: definition.key, resourceUid: this.resource.metadata.uid, generation: this.resource.metadata.generation, kind: definition.kind, requestFingerprint, }); if (reservation.record.state === "applied" || reservation.record.state === "rejected") { return reservation.record; } if (!reservation.created) { const observed = await this.observe(definition); const recovered = this.applyObservation(definition.key, observed); if (recovered !== undefined) { return recovered; } } return await this.apply(definition); } private async observe( definition: EffectDefinition, ): Promise { try { return await definition.observe(this.signal); } catch (error) { this.recordEvent("effect_observe_failed", definition.key, { error: boundedError(error), }); throw error; } } private applyObservation(key: string, observation: EffectObservation): EffectRecord | undefined { if (observation.state === "not_applied") { return undefined; } const record = this.store.updateEffect({ resourceUid: this.resource.metadata.uid, key, state: observation.state === "applied" ? "applied" : "indeterminate", ...("externalRef" in observation && observation.externalRef !== undefined ? { externalRef: observation.externalRef } : {}), }); this.recordEvent( observation.state === "applied" ? "effect_recovered" : "effect_indeterminate", key, ); return record; } private async apply(definition: EffectDefinition): Promise { let result: EffectApplication; try { result = await definition.apply(this.signal); } catch (error) { const message = boundedError(error); const record = this.store.updateEffect({ resourceUid: this.resource.metadata.uid, key: definition.key, state: "indeterminate", error: message, }); this.recordEvent("effect_indeterminate", definition.key, { error: message }); return record; } const record = this.store.updateEffect({ resourceUid: this.resource.metadata.uid, key: definition.key, state: result.state, ...(result.state === "applied" && result.externalRef !== undefined ? { externalRef: result.externalRef } : {}), ...(result.state !== "applied" && result.error !== undefined ? { error: result.error } : {}), }); this.recordEvent(`effect_${result.state}`, definition.key); return record; } private recordEvent(type: string, effectKey: string, extra: JsonObject = {}): void { this.store.recordEvent({ controller: this.resource.metadata.controller, key: this.resource.metadata.key, type, payload: { effectKey, ...extra }, }); } } function boundedError(error: unknown): string { const message = error instanceof Error ? error.message : String(error); return message.length <= 8_192 ? message : `${message.slice(0, 8_192)}…`; }