Слой базы данных
Хост Postgres
Назначение: Обеспечивает работу с PostgreSQL через Oracle-совместимые интерфейсы. Предоставляет загрузку метаданных функций, пулы соединений, короткие и длинные транзакции, нотификации и утилиты генерации обёрток. Включает адаптеры, скрывающие различия между Oracle и PostgreSQL.
Пакет предназначен для работы в составе MorphCluster‑приложения (например, Supervisor) и использует общие сервисы: DbServices, RegistryHelper, GlobalServices и т.д.
Точка входа (src/index.mjs)
export default {
name: "Postgres",
config,
clients: {},
services: {
DbServices, PgFunctions, OraFunctions, PgPool, OraQueries,
PgTransactions, PgTools, PgNotify, OraLongTransactions,
OraQueriesPool, OLTTest, OraAdpTest
}
}
Все сервисы регистрируются в ServiceHost. Они общаются через механизм зависимостей (requirements) и прямые вызовы методов друг друга.
Основные сервисы
1. PgFunctions
Файл: src/pg-queries/pg-functions/index.mjs
Загружает метаданные о всех функциях/процедурах из системных таблиц PostgreSQL и кэширует их. Используется другими сервисами для поиска функций по имени и пространству имён.
Зависимости: нет
Ключевые методы:
-
reload()– перечитывает функции из БД (вызывается при старте). -
find({namespace, name})– поиск функции по схеме и имени. -
findByName({name})– поиск только по имени (менее точно).
2. OraFunctions
Файл: src/ora-adapter/ora-functions/index.mjs
Представляет функции PostgreSQL в «оракул-подобном» виде: пакеты и процедуры. Использует PgFunctions для получения сырых метаданных и конвертирует их в формат, привычный для Oracle-клиентов.
Зависимости: PgFunctions
Ключевые методы:
-
list()– возвращает список всех схем (пакетов). -
getPackage({packageName})– список функций в схеме. -
getFunc({packageName, funcName})– детальное описание функции (параметры, типы, направления). -
reload()– обновляет данные, перезагружаяPgFunctions.
3. PgPool
Файл: src/pg-queries/pg-pool/index.mjs
Пул «коротких» соединений с PostgreSQL. Каждый запрос выполняется в отдельной транзакции (BEGIN → COMMIT/ROLLBACK). Используется сервисами OraQueries для быстрых запросов.
Зависимости: PgFunctions, RegistryHelper
Управление через CommonRegistry:
-
pgQueryPool.connectionCount– размер пула (по умолчанию 0 = отключен).
Методы:
-
execSql({sql, values, session})– выполнить SQL. -
execFunc({namespace, funcName, params, session})– выполнить функцию. -
reconnectAll()– пересоздать все соединения. -
state()– текущий статус соединений (дляOraQueriesPool). -
reload()– инициализация ссылок на служебные функции (set_user_id,write_error_log).
Особенности:
- Устанавливает контекст пользователя через
accounts.set_user_id. - Логирует ошибки в таблицу
logs.write_error_log. - Автоматически переподключается при обрывах (таймер).
4. PgTransactions
Файл: src/pg-queries/pg-transactions/index.mjs
Управляет длинными транзакциями. В отличие от PgPool, соединение удерживается открытым между вызовами, поддерживается ручное управление (BEGIN, выполнение нескольких запросов, COMMIT/ROLLBACK).
Зависимости: RegistryHelper, DbServices, PgFunctions
Управление через CommonRegistry:
-
PgTransactions.idleTimeout– таймаут бездействия транзакции (мс). -
PgTransactions.cleanTimeout– через сколько очищать закрытые транзакции (мс).
Внутренний класс PGTPool содержит массив экземпляров PGTransaction. Каждая транзакция имеет жизненный цикл: created → connecting → ready → executing → executed → fetching → closing → closed.
Методы (внешние через OraLongTransactions):
-
create({session})– создать новую транзакцию, возвращаетtrxId. -
execFunc({trxId, namespace, funcName, params, session})– выполнить функцию в транзакции. -
execSql({trxId, sql, values, session})– выполнить SQL. -
getExecResult({trxId})– получить результат после выполнения (блокирующий). -
fetch({trxId, count})– дочитать строки курсора. -
commit({trxId})/rollback({trxId})– завершить транзакцию. -
poolStatus()– состояние всех транзакций.
5. PgTools
Файл: src/pg-queries/pg-tools/index.mjs
Вспомогательный сервис, создающий функцию, возвращающую TABLE(...), на основе существующей функции, возвращающей REFCURSOR. Упрощает миграцию с Oracle.
Методы:
-
funcTableFromRefcursor({funcFrom, funcTo})– генерирует SQL-обёртку и создаёт её в БД.
6. PgNotify
Файл: src/pg-queries/pg-notify/index.mjs
Реализует механизм асинхронных уведомлений PostgreSQL (LISTEN/NOTIFY). Поддерживает подписку на каналы, переподключение и единое событие onNotify.
События:
-
onNotify– генерируется при получении уведомления (полезная нагрузка парсится из JSON).
Методы:
-
subscribe({channelName})/unsubscribe({channelName}). -
isConnected()– проверка соединения. -
getSubscriptions()– список активных подписок.
7. OraQueries
Файл: src/ora-adapter/ora-queries/index.mjs
Основной адаптер для выполнения «коротких» Oracle-подобных запросов. Принимает вызовы вида package.func с именованными параметрами, преобразует их через DbServices в реальные вызовы PostgreSQL и исполняет через PgPool.
Зависимости: PgPool, DbServices
Методы:
-
exec({queryName, params, count, offset, session})– универсальный вызов. -
execFunc(...)– явное указание пакета и функции. -
execSql({sql, params, session})– выполнение SQL с заменой именованных параметров:paramна позиционные$1. -
resetUser()/resetAllUsers()– сброс сессионного контекста (не реализованы).
8. OraLongTransactions
Файл: src/ora-adapter/ora-long-transactions/index.mjs
Аналог OraQueries, но для длинных транзакций. Использует PgTransactions для удержания соединения. Поддерживает многошаговые сценарии: создать транзакцию, выполнить несколько запросов, получить результаты, зафиксировать.
Зависимости: PgTransactions, DbServices
Методы:
-
create({session})– создать транзакцию (возвращаетoltId). -
execFunc({oltId, package, func, params, session})/execSql(...). -
getExecResult({oltId})– получить результат выполненного запроса. -
fetch({oltId, count})– получить строки из открытого курсора. -
commit({oltId})/rollback({oltId}). -
poolStatus()– состояние пула длинных транзакций.
9. OraQueriesPool
Файл: src/ora-adapter/ora-queries-pool/index.mjs
Административный интерфейс для просмотра состояния пула PgPool.
Зависимости: PgPool, DbServices
Методы:
-
state()– возвращает детализацию по соединениям (количество, статусы).
Вспомогательные классы и утилиты
ArgumentParser
Распарсивает строку аргументов функции PostgreSQL (из pg_get_function_arguments) в структурированный вид: { name, type, isOut, hasDefault }.
PgFormats
Отвечает за:
- Конвертацию типов полей при чтении из БД (числа, даты).
- Формирование SQL-вызова функции с именованными параметрами (
"arg" => $1). - Проверку кодов ошибок PostgreSQL (отделение логических ошибок от проблем соединения).
PgConnection
Управляет одним подключением для «коротких» запросов (PgPool). Выполняет SQL в транзакции (BEGIN → запрос → COMMIT/ROLLBACK). Обрабатывает refcursor, загружая данные из курсора.
PGTransaction
Класс одной длинной транзакции. Хранит состояние, управляет жизненным циклом, поддерживает выполнение SQL и функций, фетч курсора. Используется внутри PGTPool.
PGTPool
Менеджер пула PGTransaction. Создаёт транзакции по требованию, отслеживает таймауты, автоматически подчищает закрытые соединения, пишет ошибки в лог БД.
PgNotifyConnection
Подключение для LISTEN/NOTIFY. Инкапсулирует логику переподключения и подписки.
Конфигурация
Файл config.mjs содержит параметры по умолчанию:
-
NatsConnection.queue="pgqueries",timeout= 60000 -
NatsPublisher.interval= 60000
Остальные настройки берутся из config.common и CommonRegistry:
| Параметр | Описание |
|---|---|
common.postgresUri |
Строка подключения к PostgreSQL |
common.timeZone |
Часовой пояс для сессий БД |
pgQueryPool.connectionCount |
Размер пула коротких запросов (в CommonRegistry) |
PgTransactions.idleTimeout |
Таймаут простоя транзакции (мс) |
PgTransactions.cleanTimeout |
Интервал очистки закрытых транзакций (мс) |
Все сервисы получают конфигурацию через стандартный механизм Config и могут динамически обновляться через подписку на CommonRegistry.
Схема взаимодействия
-
Загрузка метаданных
PgFunctionsпри старте загружает список всех функций PostgreSQL.OraFunctionsна его основе строит «оракул-подобное» представление. -
Регистрация в DbServices Сервис
DbServices(из внешнего пакета) хранит каталог всех доступных вызовов (жёсткие и мягкие).OraQueriesиOraLongTransactionsиспользуютDbServices.getOraSql()для генерации PL/SQL‑блока с подстановкой параметров. -
Выполнение запросов
-
Короткие запросы:
OraQueries.exec→DbServices.getOraSql→ формирование вызова →PgPool.execFunc→PgConnection. -
Длинные транзакции:
OraLongTransactions.create→PgTransactions.create(создаётсяPGTransaction) → далееexecFunc/execSql→ по окончанииgetExecResult+commit/rollback.
-
Короткие запросы:
-
Нотификации
PgNotifyслушает каналы PostgreSQL и генерирует событиеonNotify, на которое могут подписываться другие сервисы. -
Администрирование
OraQueriesPoolи методpoolStatusу длинных транзакций дают мониторинг соединений.
Примечания
- Для корректной работы необходима настройка
common.postgresUriи наличие в БД схемыcarabimetaс процедурамиset_user_id,write_error_log.
Пакет DbServices
Обзор
@morphcluster/db-services — это библиотека для работы с сервисами базы данных. Она предоставляет:
-
Клиентские прокси для удалённого вызова сервисов
DbServices,OraQueriesиOraLongTransactionsчерез NATS. -
Конструкторы SQL-запросов для программного построения сложных
SELECT-выражений (SqlQuery,WhereCondition,SqlModBuilder). -
Хелперы для удобного выполнения запросов к БД из любого сервиса, включая управление длинными транзакциями (
QueriesHelper,OLTHelper). -
Инструменты развёртывания для загрузки описаний сервисов из файлов и автоматической установки хранимых функций в PostgreSQL (
DbServicesLoader,DbFunctionsInstaller).
Все компоненты ориентированы на миграцию с Oracle на PostgreSQL и обеспечивают совместимость с ожидаемыми интерфейсами Oracle-клиентов.
Состав библиотеки
Основные экспорты из index.mjs:
export { default as DbServices } from './clients/DbServices.mjs'
export { default as OraLongTransactions } from './clients/OraLongTransactions.mjs'
export { default as OraQueries } from './clients/OraQueries.mjs'
export { SqlEscape, SqlQuery, WhereCondition } from './SqlQuery.mjs'
export { default as SqlModBuilder } from './SqlModBuilder.mjs'
export { default as DbServicesLoader } from './DbServicesLoader.mjs'
export { default as QueriesHelper } from './QueriesHelper/index.mjs'
export { default as OLTHelper } from './OLTHelper.mjs'
export { default as DbFunctionsInstaller } from './DbFunctionsInstaller/index.mjs'
Клиентские прокси (ServiceNats)
Клиенты DbServices, OraQueries и OraLongTransactions являются наследниками ServiceNats (устаревшего, но всё ещё используемого) и предназначены для прозрачного обращения к соответствующим сервисам, опубликованным в NATS.
SQL-построители
SqlQuery
Класс для построения SQL-запроса SELECT с поддержкой всех основных секций: SELECT, FROM, JOIN, WHERE, ORDER BY, LIMIT/OFFSET. Предназначен для генерации читаемого и оптимизированного SQL.
import { SqlQuery } from '@morphcluster/db-services';
Методы:
-
select(expression | [expressions])— добавить колонки в SELECT. -
from(tableOrSubquery, alias?)— добавить таблицу или подзапрос в FROM. -
innerJoin(table, alias?, condition?),leftJoin(...),addJoin(type, table, alias?, condition?)— добавить JOIN. -
where— публичное свойство, экземплярWhereCondition(по умолчанию пустая AND-группа). Все условия добавляются в него. -
orderBy(column | [columns])— добавить сортировку. -
limit(limit, offset?)— установить LIMIT и OFFSET. -
build(options?)→string— собрать готовый SQL. -
buildPretty(options?)— собрать SQL с отступами (pretty print).
Параметры build:
{ pretty, indentSize, isSubquery, indent } (обычно используется внутри для рекурсивного построения подзапросов).
Пример:
const q = new SqlQuery()
.select(['u.id', 'u.name'])
.from('users', 'u')
.leftJoin('orders', 'o', new WhereCondition('simple', 'u.id = o.user_id'));
q.where.addSimple("u.active = 1");
q.where.addOr([
new WhereCondition('simple', "u.role = 'admin'"),
new WhereCondition('simple', "u.role = 'manager'")
]);
q.orderBy('u.name');
q.limit(10, 20);
console.log(q.build({ pretty: true }));
WhereCondition
Класс для построения деревьев условий WHERE. Поддерживает:
- простые выражения (
simple) -
EXISTSс подзапросом - логические группы
AND/OR - отрицание
NOT - метод
optimize()для упрощения дерева (удаление избыточных скобок,1=1/1=0).
import { WhereCondition } from '@morphcluster/db-services';
Типы условий (поле type):
-
'simple'— строка SQL (например,"a = b"). -
'exists'— подзапрос (экземплярSqlQuery). При построении превращается вEXISTS (subquery). -
'and','or'— группа условий, хранятся в массивеconditions. -
'not'— отрицание другого условия.
Конструктор: new WhereCondition(type, param), где param зависит от типа.
Основные методы:
-
add(condition)— добавить условие в текущую группу (только дляand/or). -
addSimple(expr)— добавить простое условие. -
addExists(subquery)— добавитьEXISTS. -
addNotExists(subquery)— добавитьNOT EXISTS. -
addNot(condition)— добавитьNOT. -
addOr(conditions)— создать и добавить OR-группу. -
addAnd(conditions)— создать и добавить AND-группу. -
optimize()— вернуть оптимизированную копию условия. -
build(options?)→string— собрать SQL-представление. -
buildPretty(options?)— то же с форматированием.
Оптимизация удаляет лишние уровни вложенности, упрощает константные выражения (1=1, 1=0), применяет законы де Моргана для NOT.
SqlModBuilder
Простой класс для формирования запросов INSERT и UPDATE на основе описания колонок.
import SqlModBuilder from '@morphcluster/db-services';
Конструктор: new SqlModBuilder(table) — указывает имя таблицы.
Методы:
-
set(column, type, value)— добавить значение для установки (для INSERT или UPDATE). -
setId(column, type, value)— установить идентификатор (первичный ключ) для WHERE в UPDATE. -
getInsertQuery()→{ sql, params }— возвращаетINSERT ... VALUES (... :param ...) RETURNING idи массив объектов{ name, type, value }. -
getUpdateQuery()→{ sql, params }—UPDATE ... SET ... WHERE id=:ID.
Параметры для Oracle-нотации формируются в виде именованных параметров (:VARNAME).
SqlEscape
Утилитарная функция SqlEscape(str) — экранирует строку для безопасной вставки в SQL: удваивает одинарные кавычки.
import { SqlEscape } from '@morphcluster/db-services';
DbServicesLoader
Сервис, который загружает описания DB-сервисов (hard-services) из JSON-файлов на диске и отправляет их в DbServices при старте и по событию onResendHardServices.
import { DbServicesLoader } from '@morphcluster/db-services';
Конструктор: new DbServicesLoader(host, config)
- Помечается как
local = true, т.е. не публикуется глобально. - Требует зависимость
DbServices. - Загружает файлы из директории
{Config.packageRoot}/db-services/с расширением.json. - При старте вызывает
reload()(чтение файлов) иsendDbServices()(отправка каждого сервиса черезDbServices.setHardService). - Подписывается на
DbServices.onResendHardServices, чтобы повторно отправлять сервисы по требованию.
Структура JSON-файла сервиса должна соответствовать тому, что ожидает setHardService (полное описание сервиса с запросами и т.д.). В примере кода в этих файлах добавляется поле hostName.
QueriesHelper
Универсальный хелпер для выполнения запросов к базе данных из любого сервиса. Скрывает детали работы с OraQueries и OraLongTransactions, предоставляя простые методы query, select, querySql и др.
import { QueriesHelper } from '@morphcluster/db-services';
Конструктор: new QueriesHelper(host)
- Помечается
local = true. - Зависимости:
OraLongTransactions,RegistryHelper. Опционально ожидаетOraQueries, но при необходимости создаёт его сам. - В
start()создаёт подписку на схему реестраregistrySchema, которая содержит полеqueries.type(например,"ora"или"pg").
Основные методы:
| Метод | Назначение |
|---|---|
queryRaw(queryName, params, count, offset, options) |
Выполнить вызов функции (PKG.FUNC) и вернуть сырой результат (массив объектов {paramName, type, value}) |
query(queryName, params, count, offset, options) |
То же, но преобразует курсоры в массивы объектов (колонки → ключи) и возвращает объект с ключами paramName |
select(queryName, params, count, offset, options) |
Взять первое значение из результата query |
selectRow(queryName, params, options) |
Вызвать select с count=1 и вернуть первый элемент массива, если есть |
querySql(log, SQL, rawParams, options) |
Выполнить сырой SQL через OraQueries.execSql (или через транзакцию, если передан trxId) |
querySqlCursor(log, SQL, rawParams, options) |
Как querySql, но ожидает, что первый результат — курсор, и преобразует его через convertCursor |
transactionWrap({ log, callback, trxId?, session? }) |
Выполнить функцию в транзакции: создать транзакцию, вызвать callback(trxId), при успехе — commit, при ошибке — rollback |
convertCursor(cursor) |
Преобразует объект курсора { columns, list } в массив объектов, где ключи — имена колонок |
escapeId(str) |
Экранирует идентификатор для PostgreSQL |
Параметры options для методов:
-
log— логгер (обязательно). -
trxId— идентификатор существующей транзакции. Если передан, запросы выполняются в ней, иначе — в простом пуле. -
session— объект сессии (обычно содержитuserId). -
userId— алиас дляsession.userId.
queryRaw и query автоматически определяют, нужно ли использовать транзакцию (если передан trxId), и в этом случае вызывают OraLongTransactions.execFunc + getExecResult + извлечение курсора. Если без транзакции — вызывают OraQueries.execFuncInner (прямое выполнение).
Внутренний метод _queryExec выбирает пул в зависимости от this.DBType, но в текущей версии поддерживается только "ora", который приводит к вызову OraQueries.execFuncInner.
OLTHelper
Помощник для работы с одной длинной транзакцией (Oracle Long Transaction). Упрощает создание транзакции, выполнение функций, извлечение результатов и фетчинг курсоров.
import { OLTHelper } from '@morphcluster/db-services';
Конструктор: new OLTHelper(olt, ws, oltId?, log, options?)
-
olt— экземпляр клиентаOraLongTransactions. -
oltId— существующий идентификатор транзакции (если есть). -
log— логгер (обязателен). -
options.userId/options.noUser— параметры сессии.
Методы:
| Метод | Описание |
|---|---|
getTrxId() |
Вернуть oltId, если нет — создать транзакцию и вернуть |
create() |
Создать новую транзакцию и сохранить oltId |
getResult() |
Получить getExecResult с повторными попытками при OltTimeout |
fetch(OutParams, count, keepColIndex?) |
Извлечь count строк из курсора, присутствующего в OutParams. Возвращает массив строк (уже преобразованных, если keepColIndex=false) |
execFuncMulti(queryName, params) |
Выполнить функцию и вернуть объект всех выходных параметров (ключ = paramName) |
execFunc(queryName, params) |
Выполнить функцию и вернуть значение первого выходного параметра |
selectFunc(queryName, params, count) |
Выполнить функцию и выбрать count строк из курсора |
selectRowFunc(queryName, params) |
selectFunc с count=1, вернуть первый ряд или null |
execSql(sql, params) |
Выполнить сырой SQL |
selectSql(sql, params, count) |
Выполнить SQL и выбрать строки из курсора |
selectRowSql(sql, params) |
Аналогично, одна строка |
commit() |
Зафиксировать транзакцию |
rollback() |
Откатить транзакцию |
Все методы обрабатывают ошибку OltTimeout автоматическими повторами. execFunc и execSql обновляют this.result, который затем используется для фетча. convertRow преобразует числовые поля из строк в числа.
DbFunctionsInstaller
Сервис для автоматической установки хранимых функций в PostgreSQL. Сравнивает хеши функций, хранящиеся в таблице db_installer.db_functions_1, с переданными, и при несовпадении выполняет их пересоздание.
import { DbFunctionsInstaller } from '@morphcluster/db-services';
Конструктор: new DbFunctionsInstaller(host, config, dbFunctions)
-
config.disabled— отключить автоустановку. -
dbFunctions— массив объектов схем и функций:[ { name: 'schema_name', functions: [ { name: 'func_name', source: 'CREATE OR REPLACE FUNCTION ...', hash: 'sha256...' } ] } ] - Зависит от
QueriesHelper(для выполнения SQL).
Запросы:
-
install({})— основной метод: в транзакции создаёт служебную таблицу, получает текущие хеши, находит функции, требующие обновления, и переустанавливает их. При установке используетсяSET check_function_bodies = offдля игнорирования ошибок в теле функции. После установки всех функций записывает/обновляет запись вdb_installer.db_functions_1.
Схема сервиса определена в service-schema.mjs: запрос install, требует needAdmin: true.
Примечания
-
Обработка ошибок:
QueriesHelperиOLTHelperавтоматически обрабатываютOltTimeoutповторными попытками. В случае других ошибок они пробрасываются выше, при этом вQueriesHelper.queryRawк ошибке добавляются поляqueryиqueryParams. -
Безопасность: Функция
SqlEscapeэкранирует только одинарные кавычки. Для предотвращения SQL-инъекций всегда используйте параметризованные запросы, передавая значения черезparams, а не встраивая их в SQL.
Сервис QueriesHelper
QueriesHelper — локальный сервис-помощник для выполнения запросов к базе данных в сервисах, предоставляя единый интерфейс для вызова хранимых функций, выполнения SQL и работы с транзакциями.
1. Подключение в вашем сервисе
Ваш сервис должен наследоваться от ServiceRequire и указать QueriesHelper в списке зависимостей.
import { ServiceRequire } from '@morphcluster/core'
import { QueriesHelper } from '@morphcluster/carabi'
export default class MyService extends ServiceRequire {
constructor(host, config) {
super(host, config)
this.requirements = ['QueriesHelper']
}
async start(log) {
await super.start(log)
// this.QueriesHelper уже доступен
}
}
После вызова super.start(log) свойство this.QueriesHelper будет содержать готовый к использованию экземпляр.
2. Конфигурация
QueriesHelper получает тип подключения из реестра через RegistryHelper. Ключ конфигурации: queries.type.
Возможные значения:
-
"ora"— работа черезOraQueries(Oracle-совместимый адаптер поверх PostgreSQL). -
"pg"— прямое подключение черезPgQueryPool(на данный момент не используется; при попытке использования в_queryExecвыбрасывается ошибка).
3. Основные методы запросов
3.1. queryRaw(queryName, params, count, offset, options)
Базовый метод выполнения хранимой функции или процедуры.
Формат имени: 'PKG.FUN' (пакет/схема и имя функции через точку). Для Postgres пакет заменен схемой.
Параметры:
-
queryName(string) – полное имя функции, например'csp_mainmenu.get_dashboard'. -
params(object) – объект с параметрами, ключи соответствуют именам параметров функции (без префикса:). -
count(number, по умолчанию 1) – количество строк для извлечения из курсора (если возвращается курсор). -
offset(number, по умолчанию 0) – смещение для курсора. -
options(object) – дополнительные настройки (см. раздел «Опции запросов»).
Возвращает: Promise<Array<{ paramName: string, type: string, value: any }>>
Массив объектов, описывающих все выходные параметры функции. Для курсоров value будет содержать структуру { columns: Array<[string, string]>, list: Array<Array<any>> }.
Пример:
const result = await this.QueriesHelper.queryRaw(
'csp_mainmenu.get_items',
{ user_id: 42, role_id: 1 },
10,
0,
{ log, session }
);
// result[0] может быть { paramName: 'RESULT', type: 'CURSOR', value: {...} }
3.2. query(queryName, params, count, offset, options)
Обёртка над queryRaw, которая:
- Преобразует курсоры из сырого формата в массив объектов (через
convertCursor). - Упаковывает все выходные параметры в объект, где ключ —
paramName, а значение — преобразованное значение.
Возвращает: Promise<Object>
Объект вида { PARAM1: value1, PARAM2: value2 }. Курсоры превращаются в массив объектов, где ключи — имена колонок.
Пример:
const data = await this.QueriesHelper.query(
'csp_mainmenu.get_user_info',
{ user_id: 100 },
1, 0,
{ log, session }
);
// data.USER_INFO = [ { NAME: 'John', AGE: 30 } ]
// data.STATUS = 'OK'
3.3. select(queryName, params, count, offset, options)
Упрощённый метод, когда ожидается ровно одно выходное значение. Фактически возвращает первое свойство из результата query.
Возвращает: значение первого выходного параметра (после преобразования курсора).
Пример:
const itemCount = await this.QueriesHelper.select(
'csp_mainmenu.count_items',
{ category: 'books' },
1, 0,
{ log, session }
);
// itemCount = 12 (если единственный out-параметр — число)
3.4. selectRow(queryName, params, options)
Используется, когда ожидается ровно одна строка из курсора. Возвращает первый элемент массива-результата или null.
Параметры:
-
queryName,params,options(count и offset не принимаются – под капотом используетсяcount=1, offset=0).
Возвращает: объект строки или null.
Пример:
const user = await this.QueriesHelper.selectRow(
'csp_mainmenu.get_user_by_id',
{ user_id: 55 },
{ log, session }
);
// user = { NAME: 'Alice', EMAIL: 'alice@example.com' } или null, если не найден
3.5. querySql(log, SQL, rawParams, options)
Выполняет произвольный SQL-запрос (не хранимую функцию). Поддерживает как простые запросы, так и выполнение в транзакции.
Параметры:
-
log(Logger) – обязательно. -
SQL(string) – текст запроса с плейсхолдерами в нотации Oracle (:param_name). -
rawParams(Array) – массив объектов параметров:{ name: 'PARAM', type: 'number', value: 123 }. -
options(object) –{ userId, session, trxId, count, offset }.
Возвращает: Promise<Array<{ paramName, type, value }>> (как queryRaw, но только для SQL).
Пример:
const result = await this.QueriesHelper.querySql(
log,
`UPDATE users SET name = :NEW_NAME WHERE id = :ID`,
[
{ name: 'NEW_NAME', type: 'varchar2', value: 'Bob' },
{ name: 'ID', type: 'number', value: 100 }
],
{ session, trxId: currentTrx } // если нужно в транзакции
);
3.6. querySqlCursor(log, SQL, rawParams, options)
То же, что querySql, но ожидает, что единственный выходной параметр — курсор. Сразу преобразует его через convertCursor и возвращает массив объектов (строк).
Возвращает: массив объектов (строк курсора).
Пример:
const rows = await this.QueriesHelper.querySqlCursor(
log,
`SELECT id, name FROM users WHERE role = :ROLE`,
[{ name: 'ROLE', type: 'varchar2', value: 'admin' }],
{ log, session }
);
// rows = [ { ID: 1, NAME: 'Alice' }, { ID: 2, NAME: 'Bob' } ]
3.7. convertCursor(cursor)
Вспомогательный метод, который вы можете использовать отдельно, если получили сырой курсор. Преобразует колонки и значения в массив объектов, попутно парся числа.
Сигнатура:
convertCursor(rawCursor: { columns: Array<[string, string]>, list: Array<Array<any>> }): Array<Object>
4. Опции запросов
Почти все методы принимают объект options со следующими необязательными полями:
-
log— экземплярLogger(обязателен для методовquerySql,querySqlCursor). -
session— объект сессии текущего пользователя (обычно{ userId, token, isAdmin }). -
userId— числовой ID пользователя; если передан, он будет добавлен вsession.userId. -
trxId— ID открытой транзакции, если запрос нужно выполнить в её контексте. -
count,offset— для курсоров, если не переданы отдельными аргументами (вquerySqlони берутся из options).
5. Работа с транзакциями
5.1. transactionWrap({ log, callback, trxId, session })
Выполняет переданную функцию в рамках транзакции. Если trxId не указан, создаёт новую транзакцию, после успешного выполнения коммитит, при ошибке — откатывает.
Параметры:
-
callback(async function) — принимаетtrxIdи выполняет запросы. -
trxId— можно передать существующую транзакцию, тогда коммит/откат не управляется автоматически. -
session— сессия для создания транзакции. -
log— логгер.
Пример:
await this.QueriesHelper.transactionWrap({
log,
session,
callback: async (trxId) => {
// Передаём trxId в опции всех запросов
await this.QueriesHelper.querySql(log,
`INSERT INTO audit (user_id, action) VALUES (:UID, :ACT)`,
[
{ name: 'UID', type: 'number', value: session.userId },
{ name: 'ACT', type: 'varchar2', value: 'update' }
],
{ trxId, session }
);
// ещё запросы...
}
});
Если не передать trxId, транзакция будет создана и автоматически завершена.
Если вы передали внешний trxId (например, из длинной транзакции), обёртка не выполняет commit/rollback — управление остаётся за вами.
5.2. Прямое управление транзакциями
Вы можете самостоятельно создавать, коммитить и откатывать транзакции через сервис OraLongTransactions (доступен как this.QueriesHelper.OraLongTransactions). Однако обычно удобнее использовать transactionWrap.
6. Параметры запросов
Именованные параметры передаются в виде массива объектов:
interface QueryParam {
name: string; // имя параметра (без двоеточия)
type: string; // тип: 'varchar2', 'number', 'numeric', 'json' и т.д.
value: any; // значение
}
В тексте SQL плейсхолдеры указываются с двоеточием: :USER_ID. QueriesHelper сам преобразует их в позиционные ($1, $2, ...) в зависимости от типа БД.
Для вызова хранимых функций (query, queryRaw) параметры передаются объектом, где ключи соответствуют именам параметров функции (без префикса p_ или других соглашений – смотрите документацию к вашим функциям). Внутренний механизм сам сопоставит их с аргументами функции.
7. Типы данных и преобразования
-
NUMBER — PostgreSQL-значения числовых типов автоматически преобразуются в
floatпри конвертации курсора (convertCursor). - DATE / TIMESTAMP — возвращаются как строки в формате ISO (зависит от адаптера OraQueries).
- VARCHAR2 — строки.
- CURSOR — преобразуется в массив объектов.
Если вы используете метод queryRaw, вы получаете сырые значения без дополнительных преобразований чисел.
8. Обработка ошибок
При ошибке выполнения запроса выбрасывается исключение. В него добавляются поля query и queryParams для упрощения отладки. Также ошибка автоматически логируется через log.writeExceptionOnly(e).
Рекомендуется оборачивать вызовы в try/catch и при необходимости выбрасывать ComplexError для клиентов.
9. Полный пример сервиса
import { ServiceRequire } from '@morphcluster/core'
import { QueriesHelper } from '@morphcluster/carabi'
export default class ReportService extends ServiceRequire {
constructor(host, config) {
super(host, config)
this.requirements = ['QueriesHelper']
}
async start(log) {
await super.start(log)
// Можно выполнить начальную загрузку
}
async getDailyReport({ date, session }, ws, log) {
const rows = await this.QueriesHelper.querySqlCursor(log,
`SELECT product, sum(amount) as total
FROM sales
WHERE sale_date = :SALE_DATE
GROUP BY product`,
[{ name: 'SALE_DATE', type: 'date', value: date }],
{ session, count: 1000 }
);
return rows;
}
async updateProductStock({ productId, quantity }, ws, log) {
await this.QueriesHelper.transactionWrap({
log,
session: ws.session,
callback: async (trxId) => {
await this.QueriesHelper.querySql(log,
`UPDATE products SET stock = stock - :QTY WHERE id = :ID`,
[
{ name: 'QTY', type: 'number', value: quantity },
{ name: 'ID', type: 'number', value: productId }
],
{ trxId }
);
// Дополнительные действия...
}
});
return { success: true };
}
}
Сервис SysProcesses
Сервис для асинхронного выполнения фоновых процессов бизнес-логики. Процессы могут быть инициированы через API или автоматически путём вставки записей в служебную таблицу sysprocesses.sysprocesses.
Конфигурация:
SysProcesses.disableCheck – если true, автоматическая проверка очереди процессов отключается (полезно для отладки).
Основные возможности:
- Автоматический запуск процессов, ожидающих выполнения в таблице
sysprocesses.sysprocesses. - Поддержка очередей (
QUEUE): внутри одной очереди процессы выполняются строго последовательно. - Отложенный запуск по полю
START_AT. - Ограничение количества одновременно выполняющихся процессов (
maxProcessses, по умолчанию1). - Сохранение результатов или ошибок в БД.
- API для добавления, удаления и просмотра процессов.
- Оповещение о завершении/ошибке через события
onProcessCompletedиonProcessError.
Сервис поддерживает три типа процессов:
-
csp- вызов сервисов -
db- выполнение хранимых процедур в БД с транзакцией () -
sql- выполнение произвольного SQL-кода ()
Методы
addProcess
Добавляет новый процесс в очередь.
Параметры:
-
session– объект сессии пользователя (обязательно, для записиuser_id). -
type– тип процесса ('csp','db','sql'). -
name– имя процесса:- для
csp– строка вида'ServiceName.requestName'. - для
db– полное имя хранимой функции (например,'PKG_VOCAB.INSERT_VALUE'). - для
sql– произвольное описание (сам код передаётся вsourceCode).
- для
-
params– объект с параметрами (дляcspиdb). При вызове сервиса (csp) параметры будут переданы в запрос, а также будет добавлено полеsessionсuserId. -
sourceCode– SQL-код для типаsql. -
queue– имя очереди (строковый идентификатор). Процессы в одной очереди выполняются последовательно. -
startAt– время отложенного запуска в формате ISO. -
keep– еслиtrue, запись о процессе не удаляется автоматически после завершения. Полезно для отслеживания выполнения процесса.
Возвращает: идентификатор созданного процесса (id).
Пример:
const id = await gsvc.sendRequest('addProcess', {
session: { userId: 123, isAdmin: true },
type: 'csp',
name: 'ReportService.generate',
params: { date: '2026-01-01' },
queue: 'reports',
keep: true
}, log);
delProcess
Удаляет процесс по идентификатору. Запрещено удалять выполняющийся в данный момент процесс.
Нужно только для процессов соданых с keep: true
Параметры:
-
session– сессия (обязательно). -
id– ID процесса.
Права доступа: администратор или владелец процесса (совпадение user_id).
Ошибки:
Пример:
await gsvc.sendRequest('delProcess', {
session: { userId: 123 },
id: 42
}, log);
list
Возвращает список процессов с информацией о статусе.
Параметры:
-
session– обязателен.
Права: администратор видит все процессы, обычный пользователь – только свои (user_id совпадает с session.userId).
Возвращаемый формат:
[
{
"ID": 1,
"TYPE": "csp",
"NAME": "ReportService.generate",
"USER_ID": 123,
"QUEUE": "reports",
"START_AT": null,
"STARTED": "2026-06-22T10:00:00",
"COMPLETED": null,
"ERROR_TEXT": null,
"STATUS": "running",
"INFO": { ... }
}
]
Поле STATUS может принимать значения:
-
'waiting'– ожидает запуска. -
'running'– выполняется в данный момент. -
'completed'– успешно завершён. -
'failed'– завершён с ошибкой.
Для выполняющихся процессов в INFO попадают данные, возвращаемые методом getInfo() конкретного обработчика (например, идентификатор транзакции).
checkList
Инструмент отладки. Принудительно проверяет таблицу процессов и запускает все ожидающие (с учётом очередей и лимитов).
Требует прав администратора.
Пример:
await gsvc.sendRequest('checkList', {}, log);
События
-
onProcessCompleted– возникает при успешном завершении процесса. Параметры обработчика:{ process, result }. -
onProcessError– возникает при ошибке. Параметры:{ process, errorText }.
Внутреннее устройство
Сервис использует два механизма для обнаружения новых процессов:
-
Таймер
checkTimer– периодически (каждые 60 секунд) вызываетcheckList(). -
Подписка на уведомления PostgreSQL через клиент
PgNotify. При вставке новой записи в таблицуsysprocesses.sysprocesses(через триггер в БД) отправляется уведомлениеnew_sysprocess, которое мгновенно активирует проверку очереди.
При запуске сервиса (start()) выполняется первоначальная проверка очереди, затем запускается таймер и оформляется подписка на канал new_sysprocess. При остановке (stop()) подписка снимается.
Примечания
- Максимальное число одновременно выполняющихся процессов задаётся жёстко (
maxProcessses = 1). При необходимости изменения можно унаследовать сервис и переопределить это свойство. - Автоматическое удаление записи после завершения можно отключить флагом
keep: trueпри добавлении процесса.