diff --git a/core/in_queue.js b/core/in_queue.js index 22820e7..1dce96a 100644 --- a/core/in_queue.js +++ b/core/in_queue.js @@ -24,7 +24,8 @@ const { deepClone, getKafkaConnectionSettings, getMQTTConnectionSettings, - getURLProtocol + getURLProtocol, + extractQueuePrcExecMsg } = require("./utils"); //Вспомогательные функции const { NINC_EXEC_CNT_YES, NIS_ORIGINAL_NO, NIS_ORIGINAL_YES } = require("../models/prms_db_connector"); //Схемы валидации параметров функций модуля взаимодействия с БД const objInQueueSchema = require("../models/obj_in_queue"); //Схема валидации сообщений обмена с бработчиком очереди входящих сообщений @@ -123,6 +124,8 @@ class InQueue extends EventEmitter { let optionsResp = {}; //Флаг прекращения обработки сообщения let bStopPropagation = false; + //Сообщение обработчика БД + let sDbExecMsg = null; //Нормализация сообщения const toMessageBuffer = body => { if (body === undefined || body === null) return null; @@ -258,6 +261,8 @@ class InQueue extends EventEmitter { //Если результат - ошибка аутентификации, то и её пробрасываем, но с правильным кодом if (prcRes.sResult == objQueueSchema.SPRC_RESP_RESULT_UNAUTH) throw new ServerError(SERR_UNAUTH, prcRes.sMsg || "Нет аутентификации"); + //Запомним сообщение обработчика БД + sDbExecMsg = extractQueuePrcExecMsg(prcRes); //Выставим статус сообщению очереди - исполнено обработчиком БД q = await this.dbConn.setQueueState({ nQueueId: q.nId, @@ -360,6 +365,7 @@ class InQueue extends EventEmitter { //Фиксируем успех обработки - в статусе сообщения q = await this.dbConn.setQueueState({ nQueueId: q.nId, + sExecMsg: sDbExecMsg, nIncExecCnt: NINC_EXEC_CNT_YES, nExecState: objQueueSchema.NQUEUE_EXEC_STATE_OK }); @@ -420,6 +426,8 @@ class InQueue extends EventEmitter { let options = {}; //Флаг прекращения обработки сообщения let bStopPropagation = false; + //Сообщение обработчика БД + let sDbExecMsg = null; //Получим тело сообщения blMsg = message.value ? message.value : null; //Определимся с параметрами сообщения полученными от внешней системы @@ -503,6 +511,8 @@ class InQueue extends EventEmitter { let prcRes = await this.dbConn.execQueueDBPrc({ nQueueId: q.nId }); //Если результат - ошибка пробрасываем её if (prcRes.sResult == objQueueSchema.SPRC_RESP_RESULT_ERR) throw new ServerError(SERR_DB_SERVER, prcRes.sMsg); + //Запомним сообщение обработчика БД + sDbExecMsg = extractQueuePrcExecMsg(prcRes); //Выставим статус сообщению очереди - исполнено обработчиком БД q = await this.dbConn.setQueueState({ nQueueId: q.nId, @@ -514,6 +524,7 @@ class InQueue extends EventEmitter { //Фиксируем успех обработки - в статусе сообщения q = await this.dbConn.setQueueState({ nQueueId: q.nId, + sExecMsg: sDbExecMsg, nIncExecCnt: NINC_EXEC_CNT_YES, nExecState: objQueueSchema.NQUEUE_EXEC_STATE_OK }); diff --git a/core/out_queue_processor.js b/core/out_queue_processor.js index 3ac4260..4c6ac47 100644 --- a/core/out_queue_processor.js +++ b/core/out_queue_processor.js @@ -24,7 +24,8 @@ const { getKafkaBroker, getKafkaAuth, getURLProtocol, - wrapPromiseTimeout + wrapPromiseTimeout, + extractQueuePrcExecMsg } = require("./utils"); //Вспомогательные функции const { ServerError } = require("./server_errors"); //Типовая ошибка const objOutQueueProcessorSchema = require("../models/obj_out_queue_processor"); //Схема валидации сообщений обмена с бработчиком очереди исходящих сообщений @@ -536,6 +537,8 @@ const dbProcess = async prms => { }, попытка исполнения - ${prms.queue.nExecCnt + 1}`, { nQueueId: prms.queue.nId } ); + //Сообщение обработчика БД + let sDbExecMsg = null; //Если обработчик со стороны БД указан if (prms.function.sPrcResp) { //Вызываем его @@ -544,6 +547,8 @@ const dbProcess = async prms => { if (prcRes.sResult == objQueueSchema.SPRC_RESP_RESULT_ERR) throw new ServerError(SERR_DB_SERVER, prcRes.sMsg); //Если результат - ошибка аутентификации, то и её пробрасываем, но с правильным кодом if (prcRes.sResult == objQueueSchema.SPRC_RESP_RESULT_UNAUTH) throw new ServerError(SERR_UNAUTH, prcRes.sMsg || "Нет аутентификации"); + //Запомним сообщение обработчика БД + sDbExecMsg = extractQueuePrcExecMsg(prcRes); } //Фиксируем успешное исполнение сервером БД - в протоколе работы сервиса await logger.info(`Исходящее сообщение ${prms.queue.nId} успешно отработано сервером БД`, { @@ -552,6 +557,7 @@ const dbProcess = async prms => { //Фиксируем успешное исполнение (полное - дальше обработки нет) - в статусе сообщения res = await dbConn.setQueueState({ nQueueId: prms.queue.nId, + sExecMsg: sDbExecMsg, nIncExecCnt: prms.queue.nExecCnt == 0 ? NINC_EXEC_CNT_YES : NINC_EXEC_CNT_NO, nExecState: objQueueSchema.NQUEUE_EXEC_STATE_OK }); diff --git a/core/utils.js b/core/utils.js index 87bffe4..ff335b5 100644 --- a/core/utils.js +++ b/core/utils.js @@ -444,6 +444,14 @@ const deepCopyObject = obj => JSON.parse(JSON.stringify(obj)); //Проверка на undefined const isUndefined = value => value === undefined; +//Извлечение сообщения обработчика БД +const extractQueuePrcExecMsg = prcRes => { + //Нет результата обработчика - сообщения нет + if (!prcRes) return null; + //Вернём сообщение обработчика (null, если не передано) + return isUndefined(prcRes.sMsg) ? null : prcRes.sMsg; +}; + //Считывание параметров подключения для сервиса обмена (при service === "" считывание подключения "По умолчанию", settingsArray - массив объектов [{sService: "", ...},...]) const getConnectionSettings = (service, settingsArray) => { //Считываем параметры и возвращаем @@ -565,6 +573,7 @@ exports.deepMerge = deepMerge; exports.deepClone = deepClone; exports.deepCopyObject = deepCopyObject; exports.isUndefined = isUndefined; +exports.extractQueuePrcExecMsg = extractQueuePrcExecMsg; exports.getKafkaConnectionSettings = getKafkaConnectionSettings; exports.getMQTTConnectionSettings = getMQTTConnectionSettings; exports.getKafkaBroker = getKafkaBroker;