diff --git a/server/modules/arenaImport/service/arenaImport/jobs/filesImportJob.js b/server/modules/arenaImport/service/arenaImport/jobs/filesImportJob.js index 97579b3b17..90be311c24 100644 --- a/server/modules/arenaImport/service/arenaImport/jobs/filesImportJob.js +++ b/server/modules/arenaImport/service/arenaImport/jobs/filesImportJob.js @@ -8,6 +8,8 @@ import * as FileService from '@server/modules/record/service/fileService' import * as ArenaSurveyFileZip from '../model/arenaSurveyFileZip' +const detailedLogEnabled = false + export default class FilesImportJob extends Job { constructor(params) { super('FilesImportJob', params) @@ -76,31 +78,56 @@ export default class FilesImportJob extends Job { const { surveyId, dryRun } = context const fileUuid = RecordFile.getUuid(file) const fileProps = RecordFile.getProps(file) - this.logDebug(`persisting file ${fileUuid}`) + if (detailedLogEnabled) { + this.logDebug(`persisting file ${fileUuid}`) + } const existingFileSummary = await FileService.fetchFileSummaryByUuid(surveyId, fileUuid, this.tx) - if (existingFileSummary) { + if (!existingFileSummary) { + await this.insertNewFile({ dryRun, file, fileProps, fileUuid, surveyId, tx }) + return + } + + if (detailedLogEnabled) { this.logDebug(`file already existing`) - if (RecordFile.isDeleted(existingFileSummary)) { - this.logDebug(`file previously marked as deleted: delete permanently and insert a new one`) - if (!dryRun) { - await FileService.deleteFileByUuid({ surveyId, fileUuid }, tx) - await FileService.insertFile(surveyId, file, tx) - } - this.insertedFileUuids.push(fileUuid) - } else { - this.logDebug('updating props') - if (!dryRun) { - await FileService.updateFileProps(surveyId, fileUuid, fileProps, tx) - } - this.updatedFileUuids.push(fileUuid) - } - } else { + } + + if (RecordFile.isDeleted(existingFileSummary)) { + await this.replaceDeletedFile({ dryRun, file, fileUuid, surveyId, tx }) + return + } + + await this.updateExistingFileProps({ dryRun, fileProps, fileUuid, surveyId, tx }) + } + + async insertNewFile({ dryRun, file, fileProps, fileUuid, surveyId, tx }) { + if (detailedLogEnabled) { this.logDebug(`file not existing: inserting new file`, fileProps) - if (!dryRun) { - await FileService.insertFile(surveyId, file, tx) - } - this.insertedFileUuids.push(fileUuid) } + if (!dryRun) { + await FileService.insertFile(surveyId, file, tx) + } + this.insertedFileUuids.push(fileUuid) + } + + async replaceDeletedFile({ dryRun, file, fileUuid, surveyId, tx }) { + if (detailedLogEnabled) { + this.logDebug(`file previously marked as deleted: delete permanently and insert a new one`) + } + if (!dryRun) { + await FileService.deleteFileByUuid({ surveyId, fileUuid }, tx) + await FileService.insertFile(surveyId, file, tx) + } + this.insertedFileUuids.push(fileUuid) + } + + async updateExistingFileProps({ dryRun, fileProps, fileUuid, surveyId, tx }) { + if (detailedLogEnabled) { + this.logDebug('updating props') + } + if (!dryRun) { + await FileService.updateFileProps(surveyId, fileUuid, fileProps, tx) + } + this.updatedFileUuids.push(fileUuid) } async checkFilesNotExceedingAvailableQuota(filesSummaries) { @@ -125,8 +152,9 @@ export default class FilesImportJob extends Job { } const filesUuids = filesSummaries.map(RecordFile.getUuid) - this.logDebug(`file uuids to be imported: ${filesUuids}`) - + if (detailedLogEnabled) { + this.logDebug(`file uuids to be imported: ${filesUuids}`) + } const missingRecordFileUuidsInFiles = recordsFileUuids.filter( (recordFileUuid) => !filesUuids.includes(recordFileUuid) ) diff --git a/server/modules/dataImport/service/DataImportJob/DataImportBaseJob.js b/server/modules/dataImport/service/DataImportJob/DataImportBaseJob.js index f3350c92c4..977ea22cb5 100644 --- a/server/modules/dataImport/service/DataImportJob/DataImportBaseJob.js +++ b/server/modules/dataImport/service/DataImportJob/DataImportBaseJob.js @@ -63,6 +63,9 @@ export default class DataImportBaseJob extends Job { if (nodesArray.length === 0) return for (const node of nodesArray) { + if (!node[Node.keys.recordUuid]) { + node[Node.keys.recordUuid] = recordUuid + } if (Node.isDeleted(node)) { await this.nodesDeleteBatchPersister.addItem(node) } else if (Node.isCreated(node)) { @@ -72,7 +75,7 @@ export default class DataImportBaseJob extends Job { } } await RecordManager.assocRefDataToNodes({ survey, nodes: nodesArray }, tx) - const { record: recordUpdated, rdbUpdates } = RecordManager.generateRdbUpates({ survey, record, nodesArray }, tx) + const { record: recordUpdated, rdbUpdates } = RecordManager.generateRdbUpates({ survey, record, nodesArray }) await this.rdbUpdatesBatchPersister.addItem(rdbUpdates) this.currentRecord = recordUpdated diff --git a/server/modules/mobile/service/arenaMobileDataImport/arenaMobileDataImportJob.js b/server/modules/mobile/service/arenaMobileDataImport/arenaMobileDataImportJob.js index ccb1c0af83..6db2a0893a 100644 --- a/server/modules/mobile/service/arenaMobileDataImport/arenaMobileDataImportJob.js +++ b/server/modules/mobile/service/arenaMobileDataImport/arenaMobileDataImportJob.js @@ -1,19 +1,21 @@ +import { Surveys, SystemError } from '@openforis/arena-core' + import * as Survey from '@core/survey/survey' +import FilesImportJob from '@server/modules/arenaImport/service/arenaImport/jobs/filesImportJob' +import { RecordsUpdateThreadService } from '@server/modules/record/service/update/surveyRecordsThreadService' +import { RecordsUpdateThreadMessageTypes } from '@server/modules/record/service/update/thread/recordsThreadMessageTypes' +import * as SurveyService from '@server/modules/survey/service/surveyService' +import RecordCheckJob from '@server/modules/survey/service/recordCheckJob' + import Job from '@server/job/job' import FileZip from '@server/utils/file/fileZip' import RecordsImportJob from './jobs/recordsImportJob' -import FilesImportJob from '../../../arenaImport/service/arenaImport/jobs/filesImportJob' -import { RecordsUpdateThreadService } from '@server/modules/record/service/update/surveyRecordsThreadService' -import { RecordsUpdateThreadMessageTypes } from '@server/modules/record/service/update/thread/recordsThreadMessageTypes' -import * as SurveyService from '@server/modules/survey/service/surveyService' -import { Surveys, SystemError } from '@openforis/arena-core' export default class ArenaMobileDataImportJob extends Job { /** * Creates a new data import job to import survey records and files in Arena format. - * * @param {!object} params - The import parameters. * @param {!object} [params.user] - The user performing the import. * @param {!number} [params.surveyId] - The id of the survey in which data will be imported. @@ -22,7 +24,7 @@ export default class ArenaMobileDataImportJob extends Job { * @returns {ArenaMobileDataImportJob} - The import job. */ constructor(params) { - super(ArenaMobileDataImportJob.type, params, [new RecordsImportJob(), new FilesImportJob()]) + super(ArenaMobileDataImportJob.type, params, [new RecordsImportJob(), new FilesImportJob(), new RecordCheckJob()]) } async onStart() { diff --git a/server/modules/mobile/service/arenaMobileDataImport/jobs/recordsImportJob.js b/server/modules/mobile/service/arenaMobileDataImport/jobs/recordsImportJob.js index 803b4a18f3..6ea6e081a1 100644 --- a/server/modules/mobile/service/arenaMobileDataImport/jobs/recordsImportJob.js +++ b/server/modules/mobile/service/arenaMobileDataImport/jobs/recordsImportJob.js @@ -1,4 +1,4 @@ -import { Dates, Objects, Records, Surveys } from '@openforis/arena-core' +import { Dates, Objects, RecordFixer, Records, Surveys } from '@openforis/arena-core' import { ConflictResolutionStrategy } from '@common/dataImport' @@ -25,27 +25,6 @@ const resultKeys = { const categoryItemProvider = CategoryItemProviderDefault const taxonProvider = TaxonProviderDefault -const checkNodeIsValid = ({ nodes, node, nodeDef }) => { - if (!nodeDef) { - return { valid: false, error: 'refers a missing node definition' } - } - const parentUuid = Node.getParentUuid(node) - if ((!parentUuid && !NodeDef.isRoot(nodeDef)) || (parentUuid && !nodes[parentUuid])) { - return { valid: false, error: `has missing or invalid parent_uuid` } - } - if (NodeDef.isMultipleAttribute(nodeDef) && Node.isValueBlank(node)) { - return { valid: false, error: `is multiple and has an empty value` } - } - const nodeHierarchy = Node.getHierarchy(node) - if ( - nodeHierarchy.length !== NodeDef.getMetaHierarchy(nodeDef)?.length || - nodeHierarchy.some((ancestorUuid) => !nodes[ancestorUuid]) - ) { - return { valid: false, error: `has an invalid meta hierarchy` } - } - return { valid: true } -} - const getRecordFormattedKeyValues = ({ survey, record }) => { const rootDef = Surveys.getNodeDefRoot({ survey }) const recordRootEntity = Records.getRoot(record) @@ -57,6 +36,9 @@ const getRecordFormattedKeyValues = ({ survey, record }) => { }) } +const nodeHierarchyLengthComparator = (nodeA, nodeB) => + Node.getHierarchy(nodeA).length - Node.getHierarchy(nodeB).length + export default class RecordsImportJob extends DataImportBaseJob { constructor(params) { super(RecordsImportJob.type, params) @@ -117,19 +99,18 @@ export default class RecordsImportJob extends DataImportBaseJob { trackFileUuids({ nodes }) { // keep track of file uuids found in record attribute values const { survey } = this.context - Object.values(nodes).forEach((node) => { + for (const node of Object.values(nodes)) { const nodeDef = Survey.getNodeDefByUuid(Node.getNodeDefUuid(node))(survey) if (NodeDef.isFile(nodeDef)) { this.trackFileUuid({ node }) } - }) + } } async cleanupCurrentRecord() { const { context, currentRecord: record, user, tx } = this const { survey } = context - const recordUuid = Record.getUuid(record) // check owner uuid: if user not defined, use the job user as owner const ownerUuidSource = Record.getOwnerUuid(record) const ownerSource = await UserService.fetchUserByUuid(ownerUuidSource, tx) @@ -137,25 +118,10 @@ export default class RecordsImportJob extends DataImportBaseJob { // remove invalid nodes and build index from scratch delete record['_nodesIndex'] - const nodes = Record.getNodes(record) - for (const [nodeUuid, node] of Object.entries(nodes)) { - const nodeDefUuid = Node.getNodeDefUuid(node) - const nodeDef = Survey.getNodeDefByUuid(nodeDefUuid)(survey) - const { valid, error } = checkNodeIsValid({ nodes, node, nodeDef }) - if (valid) { - // ensure recordUuid is set in node - node[Node.keys.recordUuid] = recordUuid - Node.removeFlags({ sideEffect: true })(node) - } else { - const messagePrefix = `record ${Record.getUuid(record)}: node with uuid ${Node.getUuid(node)} and node def ${NodeDef.getName(nodeDef)} (uuid ${nodeDefUuid})` - const messageSuffix = `: skipping it` - this.logWarn(`${messagePrefix} ${error} ${messageSuffix}`) - delete nodes[nodeUuid] - } - } - // assoc nodes and build index from scratch - this.currentRecord = Record.assocNodes({ nodes, sideEffect: true })(record) + // fix record (e.g. insert missing nodes, remove status flags) + // do side effect to avoid creating new objects and use less memory + RecordFixer.fixRecord({ survey, record, sideEffect: true }) } findExistingRecordSummaryWithSameKeys() { @@ -272,26 +238,25 @@ export default class RecordsImportJob extends DataImportBaseJob { await RecordManager.insertRecord(user, surveyId, record, true, tx) // insert nodes (add them to batch persister) - const nodesIndexedByUuid = Record.getNodesArray(record) - .sort((nodeA, nodeB) => Node.getHierarchy(nodeA).length - Node.getHierarchy(nodeB).length) - .reduce((acc, node) => { - const nodeUuid = Node.getUuid(node) - const nodeDefUuid = Node.getNodeDefUuid(node) - // check that the node definition associated to the node has not been deleted from the survey - const nodeDef = Survey.getNodeDefByUuid(nodeDefUuid)(survey) - if (nodeDef) { - node[Node.keys.created] = true // do side effect to avoid creating new objects - acc[nodeUuid] = node - if (NodeDef.isFile(nodeDef)) { - this.trackFileUuid({ node }) - } - } else { - this.logDebug( - `Record ${recordUuid}: missing node def with uuid ${nodeDefUuid} in node ${nodeUuid}; skipping it` - ) + const nodesIndexedByUuid = {} + const nodesArraySorted = Record.getNodesArray(record).sort(nodeHierarchyLengthComparator) + for (const node of nodesArraySorted) { + const nodeUuid = Node.getUuid(node) + const nodeDefUuid = Node.getNodeDefUuid(node) + // check that the node definition associated to the node has not been deleted from the survey + const nodeDef = Survey.getNodeDefByUuid(nodeDefUuid)(survey) + if (nodeDef) { + node[Node.keys.created] = true // do side effect to avoid creating new objects + nodesIndexedByUuid[nodeUuid] = node + if (NodeDef.isFile(nodeDef)) { + this.trackFileUuid({ node }) } - return acc - }, {}) + } else { + this.logDebug( + `Record ${recordUuid}: missing node def with uuid ${nodeDefUuid} in node ${nodeUuid}; skipping it` + ) + } + } if (!Record.getDateModified(record)) { this.logDebug(`Empty date modified for record ${Record.getUuid(record)}`) @@ -305,12 +270,17 @@ export default class RecordsImportJob extends DataImportBaseJob { async beforeSuccess() { await super.beforeSuccess() + const { insertedRecordsUuids, updatedRecordsUuids } = this const recordsFileUuidsArray = Array.from(this.recordsFileUuids) - const recordsFilesCount = recordsFileUuidsArray.length - if (recordsFilesCount > 0) { - this.logDebug(`found ${recordsFilesCount} files:`, recordsFileUuidsArray) - } - this.setContext({ recordsFileUuids: recordsFileUuidsArray }) + + this.setContext({ + recordsFileUuids: recordsFileUuidsArray, + // flag to cleanup records in RecordsCheckJob + cleanupRecords: true, + updateRdb: true, + // record UUIDs to consider in RecordCheckJob + recordUuidsIncluded: [...insertedRecordsUuids, ...updatedRecordsUuids], + }) } generateResult() { diff --git a/server/modules/survey/service/recordCheckJob.js b/server/modules/survey/service/recordCheckJob.js index 69555635ef..fbbacc9d06 100644 --- a/server/modules/survey/service/recordCheckJob.js +++ b/server/modules/survey/service/recordCheckJob.js @@ -11,20 +11,30 @@ import * as Validation from '@core/validation/validation' import BatchPersister from '@server/db/batchPersister' import Job from '@server/job/job' +import { RdbUpdatesBatchPersister } from '@server/modules/record/manager/RdbUpdatesBatchPersister' + import * as SurveyManager from '../manager/surveyManager' import * as RecordManager from '../../record/manager/recordManager' export default class RecordCheckJob extends Job { constructor(params) { super(RecordCheckJob.type, params) + } + + async onStart() { + await super.onStart() + const { surveyId, tx, user } = this this.surveyAndNodeDefsByCycle = {} // Cache of surveys and updated node defs by cycle this.nodesBatchInserter = new BatchPersister(this.nodesBatchInsertHandler.bind(this), 2500) this.nodesBatchUpdater = new BatchPersister(this.nodesBatchUpdateHandler.bind(this), 2500) + this.rdbUpdatesBatchPersister = new RdbUpdatesBatchPersister({ user, surveyId, tx }) } async execute() { - const recordsUuidAndCycle = await RecordManager.fetchRecordsUuidAndCycle({ surveyId: this.surveyId }, this.tx) + const { surveyId, tx, context } = this + const { recordUuidsIncluded = null } = context + const recordsUuidAndCycle = await RecordManager.fetchRecordsUuidAndCycle({ surveyId, recordUuidsIncluded }, tx) this.total = R.length(recordsUuidAndCycle) @@ -118,7 +128,7 @@ export default class RecordCheckJob extends Job { const { context, surveyId, user, tx } = this const { survey, nodeDefAddedUuids, nodeDefUpdatedUuids, nodeDefDeletedUuids, allNotDeletedNodeDefUuids } = surveyAndNodeDefs - const { cleanupRecords } = context + const { cleanupRecords, updateRdb } = context // this.logDebug(`checking record ${recordUuid}`) @@ -167,11 +177,16 @@ export default class RecordCheckJob extends Job { nodeDefAddedOrUpdatedUuidsUnique.add(Node.getNodeDefUuid(nodeInserted)) } const nodeDefAddedOrUpdatedUuids = Array.from(nodeDefAddedOrUpdatedUuidsUnique) - if (nodeDefAddedOrUpdatedUuids.length > 0) { + + const nodeDefToCheckForDefaultValuesAndApplicabilityUuids = cleanupRecords + ? allNotDeletedNodeDefUuids + : nodeDefAddedOrUpdatedUuids + + if (nodeDefToCheckForDefaultValuesAndApplicabilityUuids.length > 0) { // this.logDebug('applying default values') const { record: recordUpdate, nodes: nodesUpdatedDefaultValues = {} } = await _applyDefaultValuesAndApplicability( survey, - nodeDefAddedOrUpdatedUuids, + nodeDefToCheckForDefaultValuesAndApplicabilityUuids, record, nodesInsertedByUuid, tx @@ -208,6 +223,13 @@ export default class RecordCheckJob extends Job { this.tx ) } + + if (updateRdb) { + await RecordManager.assocRefDataToNodes({ survey, nodes: allUpdatedNodesArray }, tx) + const { rdbUpdates } = RecordManager.generateRdbUpates({ survey, record, nodesArray: allUpdatedNodesArray }) + await this.rdbUpdatesBatchPersister.addItem(rdbUpdates) + } + // this.logDebug('record check complete') } @@ -251,6 +273,7 @@ export default class RecordCheckJob extends Job { super.beforeSuccess() await this.nodesBatchInserter.flush(this.tx) await this.nodesBatchUpdater.flush(this.tx) + await this.rdbUpdatesBatchPersister.flush() } }