/* eslint-disable max-len */ /* eslint-disable @typescript-eslint/no-explicit-any */ /* eslint-disable @typescript-eslint/no-non-null-assertion */ import { BaseDependentTypeImportSpec, TypeImportSpec, WebApiSpec } from '../schema/WebApiSchema.js'; import { Config } from '../types/Config.js'; import { PagingState } from '../types/PagingState.js'; import { Plugin } from '../types/Plugin.js'; import { PluginContext } from '../types/PluginContext.js'; import { SqupImportContext } from '../types/SqupContext.js'; import type { } from '../types/fetch.js'; // NodeJS v18 has fetch in core, but no types :( import { CallApi } from './ApiFetcher.js'; import { CreateVertex } from './CreateVertex.js'; import { getArray } from './GetArray.js'; type Edge = { label: string, inV: string, outV: string }; type Job = { importSpec: BaseDependentTypeImportSpec, parentItems: Record[] }; export async function ImportObjects( plugin: Plugin, webapiSpec: WebApiSpec, config: Config, squpContext: SqupImportContext ): Promise { const imports = webapiSpec.import; const apiDefinitions = (Array.isArray(webapiSpec.api) ? webapiSpec.api : []).reduce( (acc, val) => acc.set(val.name, val), new Map() ); const pluginContext: PluginContext = { plugin, config, apiDefinitions, lookUps: new Map(), lookUpCache: new Map(), dataTransforms: webapiSpec.dataTransforms ?? [] }; squpContext.log.debug('Import 1', { plugin, imports, config } ); if (!Array.isArray(imports?.types) || imports.types.length <= 0) { // No import required return; } if (!squpContext.pageApi.get('initialized')) { squpContext.pageApi.set('firstGraphImport', true); squpContext.pageApi.set('importState', []); squpContext.pageApi.set('complete', false); squpContext.pageApi.set('initialized', true); } // Main import loop let first = true; while (!squpContext.pageApi.get('complete') && (first || (squpContext.canContinue ? squpContext.canContinue() : true))) { first = false; // First we check for queued jobs from earlier... const jobQueue: Job[] = squpContext.pageApi.get('jobQueue'); if (Array.isArray(jobQueue) && jobQueue.length > 0) { const job = jobQueue.shift(); if (job) { // There's a job to run, so run it now squpContext.pageApi.set('jobQueue', jobQueue); await RunQueuedJob(pluginContext, job.importSpec, squpContext, job.parentItems); // Return to while loop check continue; } } // Get the current paging state object let typeNum = squpContext.pageApi.get('importState').length - 1; let pagingState: PagingState = typeNum === -1 ? { complete: true, recordCount: 0 } : squpContext.pageApi.get('importState')[typeNum]; if (pagingState.complete) { // Move to importing next object type (if there are more) if (++typeNum >= imports.types.length) { squpContext.pageApi.set('complete', true); } else { // Create new paging state for the next object type pagingState = { complete: false, recordCount: 0 }; squpContext.pageApi.set('importState', [...squpContext.pageApi.get('importState'), pagingState]); } } if (!pagingState.complete) { await PerformSingleImport( pluginContext, imports.types[typeNum], pagingState, squpContext ); } } if (squpContext.pageApi.get('complete')) { await squpContext.graphImport( config, { vertices: [], edges: [] }, squpContext.pageApi.get('firstGraphImport'), squpContext.pageApi.get('complete')); squpContext.pageApi.clear(); } } async function handleDependencies ( pluginContext: PluginContext, dependentObjects: BaseDependentTypeImportSpec[], item: any, squpContext: SqupImportContext, parentItems: Record[] = [] ) { for (const dependentObject of dependentObjects) { if (dependentObject.api) { if (!dependentObject.api.isScalar) { throw new Error('API for dependent objects must be scalar'); } // We can't do this now as we may blow our caller's elapsed time or payload size limit, so queue it up... const jobQueue: Job[] = squpContext.pageApi.get('jobQueue') ?? []; // It is naive pushing on the whole of the parent array every time here. We should hash the objects and // use hashes here to avoid the risk of making pagingContext too large. jobQueue.push({ importSpec: dependentObject, parentItems }); squpContext.pageApi.set('jobQueue', jobQueue); } else { const pagingState = { complete: false, recordCount: 0 }; if (dependentObject.forEach) { if (!dependentObject.forEach.startsWith('@')) { throw new Error(`Invalid forEach value (${dependentObject.forEach}) - must start with "@"`); } const arr = getArray(item, dependentObject.forEach.substring(1)); if (Array.isArray(arr)) { for (const el of arr) { await PerformSingleImport(pluginContext, dependentObject, pagingState, squpContext, [el, item, ...parentItems]); } } } else { await PerformSingleImport(pluginContext, dependentObject, pagingState, squpContext, [item, ...parentItems]); } } } }; async function PerformSingleImport( pluginContext: PluginContext, importTypeSpec: TypeImportSpec, pagingState: PagingState, squpContext: SqupImportContext, parentItems: Record[] = [] ) { const vertices: Record[] = []; let edges: Edge[] = []; if (importTypeSpec.api) { const { data } = await CallApi(pluginContext, importTypeSpec.api, pagingState, squpContext); if (importTypeSpec.api.isScalar) { if (Array.isArray(data) && data.length !== 1) { throw new Error('Got array when expecting scalar'); } const { vertex, edges: vertexEdges } = await CreateVertex(pluginContext, importTypeSpec, squpContext, data, parentItems); if (vertex) { vertices.push(vertex); } edges = [...edges, ...vertexEdges]; if (Array.isArray(importTypeSpec.dependentObjects) && importTypeSpec.dependentObjects.length > 0) { await handleDependencies(pluginContext, importTypeSpec.dependentObjects, data, squpContext, [data, ...parentItems]); } } else { const dataArray = Array.isArray(data) ? data : [data]; for (const item of dataArray) { if (importTypeSpec.fields.length > 0) { const { vertex, edges: vertexEdges } = await CreateVertex(pluginContext, importTypeSpec, squpContext, item, parentItems); if (vertex) { vertices.push(vertex); } edges = [...edges, ...vertexEdges]; } if (Array.isArray(importTypeSpec.dependentObjects) && importTypeSpec.dependentObjects.length > 0) { await handleDependencies(pluginContext, importTypeSpec.dependentObjects, item, squpContext, [item, ...parentItems]); } } } } else { // Hard-coded node const { vertex, edges: vertexEdges } = await CreateVertex(pluginContext, importTypeSpec, squpContext, {}, parentItems); if (vertex) { vertices.push(vertex); } edges = [...edges, ...vertexEdges]; } if (vertices.length > 0 || edges.length > 0) { const graphImportPayload = { vertices, edges }; await squpContext.graphImport( pluginContext.config, graphImportPayload, squpContext.pageApi.get('firstGraphImport') ); squpContext.pageApi.set('firstGraphImport', false); } if (!importTypeSpec.api) { pagingState.complete = true; } } async function RunQueuedJob( pluginContext: PluginContext, importTypeSpec: TypeImportSpec, squpContext: SqupImportContext, parentItems: Record[] = [] ) { const vertices: Record[] = []; let edges: Edge[] = []; if (!importTypeSpec.api || !importTypeSpec.api.isScalar) { throw new Error('Queued jobs should be for scalar APIs only'); } const baseReplacements: Record = {}; if (parentItems.length > 0) { baseReplacements.parent = parentItems[0]; baseReplacements.parents = parentItems; } const pagingState = { complete: false, recordCount: 0 }; const { data } = await CallApi(pluginContext, importTypeSpec.api, pagingState, squpContext, baseReplacements); if (Array.isArray(data) && data.length !== 1) { throw new Error('Got array when expecting scalar'); } const { vertex, edges: vertexEdges } = await CreateVertex(pluginContext, importTypeSpec, squpContext, data, parentItems); if (vertex) { vertices.push(vertex); } edges = [...edges, ...vertexEdges]; if (Array.isArray(importTypeSpec.dependentObjects)) { await handleDependencies(pluginContext, importTypeSpec.dependentObjects, data, squpContext, [data, ...parentItems]); } if (vertices.length > 0 || edges.length > 0) { const graphImportPayload = { vertices, edges }; await squpContext.graphImport( pluginContext.config, graphImportPayload, squpContext.pageApi.get('firstGraphImport') ); squpContext.pageApi.set('firstGraphImport', false); } }