import type { BoundedOtelOperationResult, OtelOperationSettlement, } from "../otel/shutdown.ts"; type OtelOperation = BoundedOtelOperationResult["operation"]; type OtelOperationRetry = () => Promise; type OtelOperationSettlementObserver = (settlement: OtelOperationSettlement) => void; type OwnedOperationState = "settling" | "retryable"; interface OwnedOtelOperation { readonly operation: OtelOperation; readonly retry: OtelOperationRetry; readonly observeSettlement: OtelOperationSettlementObserver; state: OwnedOperationState; pendingSettlement?: Promise; } export interface OtelOperationOwnership { readonly hasUnresolvedOperations: boolean; retain: ( result: BoundedOtelOperationResult, retry: OtelOperationRetry, observeSettlement: OtelOperationSettlementObserver, ) => void; resolveBeforeStart: () => Promise; takeStartupDiagnostic: () => boolean; } const PROCESS_OWNERSHIP_KEY = Symbol.for("@senad-d/observme.otel-operation-ownership.v1"); const operationRetryOrder: readonly OtelOperation[] = ["flush", "shutdown"]; type ObservMeProcessGlobal = typeof globalThis & { [PROCESS_OWNERSHIP_KEY]?: OtelOperationOwnership; }; export function createOtelOperationOwnership(): OtelOperationOwnership { return new ProcessOtelOperationOwnership(); } export function getProcessOtelOperationOwnership(): OtelOperationOwnership { const processGlobal = globalThis as ObservMeProcessGlobal; const existing = processGlobal[PROCESS_OWNERSHIP_KEY]; if (existing) return existing; const ownership = createOtelOperationOwnership(); Object.defineProperty(processGlobal, PROCESS_OWNERSHIP_KEY, { configurable: false, enumerable: false, value: ownership, writable: false, }); return ownership; } class ProcessOtelOperationOwnership implements OtelOperationOwnership { readonly #operations = new Set(); #diagnosticEmitted = false; #resolution?: Promise; get hasUnresolvedOperations(): boolean { return this.#operations.size > 0; } retain( result: BoundedOtelOperationResult, retry: OtelOperationRetry, observeSettlement: OtelOperationSettlementObserver, ): void { if (operationCompleted(result)) return; const ownedOperation: OwnedOtelOperation = { operation: result.operation, retry, observeSettlement, state: operationCanRetry(result) ? "retryable" : "settling", }; this.#operations.add(ownedOperation); this.observePendingSettlement(ownedOperation, result.settlement); } async resolveBeforeStart(): Promise { if (this.#resolution) return this.#resolution; const resolution = this.retryUnresolvedOperations(); this.#resolution = resolution; return this.finishResolution(resolution); } takeStartupDiagnostic(): boolean { if (!this.hasUnresolvedOperations || this.#diagnosticEmitted) return false; this.#diagnosticEmitted = true; return true; } private async finishResolution(resolution: Promise): Promise { try { return await resolution; } finally { if (this.#resolution === resolution) this.#resolution = undefined; } } private async retryUnresolvedOperations(): Promise { for (const operation of operationRetryOrder) { for (const ownedOperation of this.#operations) { if (ownedOperation.operation !== operation) continue; if (ownedOperation.state === "settling") return false; const result = await ownedOperation.retry(); if (operationCompleted(result)) { this.release(ownedOperation); continue; } ownedOperation.state = operationCanRetry(result) ? "retryable" : "settling"; this.observePendingSettlement(ownedOperation, result.settlement); return false; } } return !this.hasUnresolvedOperations; } private observePendingSettlement( ownedOperation: OwnedOtelOperation, settlement: Promise | undefined, ): void { if (!settlement) return; ownedOperation.state = "settling"; ownedOperation.pendingSettlement = settlement; void this.applyPendingSettlement(ownedOperation, settlement); } private async applyPendingSettlement( ownedOperation: OwnedOtelOperation, settlementPromise: Promise, ): Promise { const settlement = await resolveSettlement(ownedOperation.operation, settlementPromise); if (!this.#operations.has(ownedOperation) || ownedOperation.pendingSettlement !== settlementPromise) return; ownedOperation.pendingSettlement = undefined; notifySettlementObserver(ownedOperation.observeSettlement, settlement); if (operationCompleted(settlement)) { this.release(ownedOperation); return; } ownedOperation.state = "retryable"; } private release(ownedOperation: OwnedOtelOperation): void { if (!this.#operations.delete(ownedOperation)) return; ownedOperation.pendingSettlement = undefined; if (this.#operations.size === 0) this.#diagnosticEmitted = false; } } async function resolveSettlement( operation: OtelOperation, settlementPromise: Promise, ): Promise { try { return await settlementPromise; } catch (error) { return { operation, completed: false, timedOut: false, error }; } } function notifySettlementObserver( observer: OtelOperationSettlementObserver, settlement: OtelOperationSettlement, ): void { try { observer(settlement); } catch { return; } } function operationCompleted( result: Pick, ): boolean { return result.completed && !result.timedOut && !result.error; } function operationCanRetry(result: BoundedOtelOperationResult): boolean { return !result.timedOut && !result.settlement; }