ЦИТК-1079. Добавление возможности сохранения сообщения обработки в очереди обмена при успешном выполнении функции сервиса обмена

This commit is contained in:
boa604 2026-07-27 12:44:35 +03:00
parent a0177a2d4c
commit 5690a79d1a
3 changed files with 28 additions and 2 deletions

View File

@ -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
});

View File

@ -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
});

View File

@ -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;