import { defineBatchStrategyMap, type BatchOperationStrategy, } from './batching-types'; type OpenSosDataScalarPayload = Record & { entity_name: string; state: string; }; type OpenSosDataBulkPayload = Record & { entities: Array<{ entity_name: string; state: string }>; }; type OpenSosDataBulkResultRow = Record & { entity_name?: string; state?: string; success?: boolean; data?: unknown; error?: string | null; cost?: number; }; type OpenSosDataBulkResult = Record & { job_id?: string; status?: string; polling_interrupted?: boolean; polling_timed_out?: boolean; recovery_operation?: string; results?: OpenSosDataBulkResultRow[]; }; function normalizeEntityName(value: unknown): string { return String(value ?? '') .trim() .replace(/\s+/g, ' ') .toLowerCase(); } function normalizeState(value: unknown): string { return String(value ?? '') .trim() .toUpperCase(); } function correlationKey(input: { entity_name?: unknown; state?: unknown; }): string { return JSON.stringify([ normalizeEntityName(input.entity_name), normalizeState(input.state), ]); } function requireResultIdentity(row: OpenSosDataBulkResultRow): string { if ( typeof row.entity_name !== 'string' || !row.entity_name.trim() || typeof row.state !== 'string' || !row.state.trim() ) { throw new Error( 'OpenSOSData bulk result is missing entity_name/state correlation identity.', ); } return correlationKey(row); } function toScalarResult( row: OpenSosDataBulkResultRow, ): Record { const scalarResult: Record = { ...row }; delete scalarResult.entity_name; delete scalarResult.state; if (typeof scalarResult.success !== 'boolean') { scalarResult.success = row.data !== undefined; } return scalarResult; } export const opensosDataBusinessLookupBatchStrategy: BatchOperationStrategy< OpenSosDataScalarPayload, OpenSosDataBulkPayload, OpenSosDataBulkResult, Record, OpenSosDataBulkResultRow > = { sourceOperation: 'opensosdata_business_lookup', batchOperation: 'opensosdata_bulk_lookup', kind: 'identifier_batch', maxBatchSize: 256, canBatchWith() { return true; }, toBucketKey() { return 'opensosdata_bulk_lookup'; }, toItemKey(payload) { return correlationKey(payload); }, compile(payloads) { return { batchOperation: 'opensosdata_bulk_lookup', batchPayload: { entities: payloads.map((payload) => ({ entity_name: payload.entity_name, state: payload.state, })), }, items: payloads.map((payload) => ({ itemKey: correlationKey(payload), payload, })), }; }, splitResult(fullResult, compiled) { if ( typeof fullResult.job_id === 'string' && (fullResult.status === 'queued' || fullResult.status === 'processing') ) { return compiled.items.map((item) => ({ itemKey: item.itemKey, result: { job_id: fullResult.job_id, status: fullResult.status, polling_interrupted: fullResult.polling_interrupted === true, polling_timed_out: fullResult.polling_timed_out === true, recovery_operation: fullResult.recovery_operation ?? 'opensosdata_get_bulk_result', }, rawResult: fullResult, })); } if (!Array.isArray(fullResult.results)) { throw new Error( 'OpenSOSData completed bulk result is missing its results array.', ); } const remainingByKey = new Map(); for (const item of compiled.items) { remainingByKey.set( item.itemKey, (remainingByKey.get(item.itemKey) ?? 0) + 1, ); } const rowsByKey = new Map(); for (const row of fullResult.results) { const key = requireResultIdentity(row); const remaining = remainingByKey.get(key) ?? 0; if (remaining < 1) { throw new Error( `OpenSOSData bulk result has unmatched correlation identity ${key}.`, ); } remainingByKey.set(key, remaining - 1); const rows = rowsByKey.get(key); if (rows) rows.push(row); else rowsByKey.set(key, [row]); } return compiled.items.map((item) => { const row = rowsByKey.get(item.itemKey)?.shift(); if (!row) { throw new Error( `OpenSOSData bulk result is missing result identity ${item.itemKey}.`, ); } return { itemKey: item.itemKey, result: toScalarResult(row), rawResult: row, }; }); }, }; export const opensosdataBatchStrategies = defineBatchStrategyMap({ opensosdata_business_lookup: opensosDataBusinessLookupBatchStrategy, });