Загрузка данных
dmitriev-aal@VDI-Dmitriev-A:~/Desktop/rshbintech/it-corp/rprul/ecp/acquiring/backend/ms-acquiring-manage-b2c-limits$ sed -n '1,120p' src/infrastructure/database/database.module.ts
import { TypeOrmModule, TypeOrmModuleOptions } from '@nestjs/typeorm';
import { ConfigService } from '@nestjs/config';
import { Module, Global } from '@nestjs/common';
import { DatabaseConfig } from '../config/database.config';
import { Person } from '../../modules/person/person.entity';
import { ClientPerson } from '../../modules/client-person/client-person.entity';
import { Phone } from '../../modules/phone/phone.entity';
import { Limit } from '../../modules/limit/limit.entity';
import { LimitHistory } from '../../modules/limit-history/limit-history.entity';
import { Operation } from '../../modules/operation/operation.entity';
import { PreAuthorization } from '../../modules/pre-authorization/pre-authorization.entity';
@Global()
@Module({
imports: [
TypeOrmModule.forRootAsync({
inject: [ConfigService],
useFactory: (config: ConfigService): TypeOrmModuleOptions => {
const db = config.get<DatabaseConfig>('database')!;
return {
type: 'postgres',
host: db.host,
port: db.port,
username: db.username,
password: db.password,
database: db.database,
schema: db.schema,
entities: [
Person,
ClientPerson,
Phone,
Limit,
LimitHistory,
Operation,
PreAuthorization,
],
synchronize: db.synchronize,
logging: db.logging,
};
},
}),
],
})
export class DatabaseModule {}
dmitriev-aal@VDI-Dmitriev-A:~/Desktop/rshbintech/it-corp/rprul/ecp/acquiring/backend/ms-acquiring-manage-b2c-limits$ sed -n '1,180p' src/infrastructure/kafka/kafka.service.ts
import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common';
import { randomBytes } from 'crypto';
import { Consumer, Kafka, KafkaMessage, Producer } from 'kafkajs';
import { DefaultLogger } from '@rshbintech.rprul.ecp/core-nodejs-starter-lib';
import { SchedulerRegistry } from '@nestjs/schedule';
import { ConfigService } from '@nestjs/config';
import { EventEmitter2, OnEvent } from '@nestjs/event-emitter';
import {
EVENT_KAFKA_RECEIVED_MESSAGE,
EVENT_KAFKA_SEND_MESSAGE,
} from '../../common/constants/constans';
import { KafkaMessageReceivedEvent } from '../../common/events/kafka-message-received.event';
import { KafkaSendMessageEvent } from '../../common/events/kafka-send-message.event';
import { KafkaSettings } from '../../infrastructure/setting/kafka/kafka.settings';
interface ConsumerTopicPair {
consumer: Consumer | null;
topic: string;
}
@Injectable()
export class KafkaService implements OnModuleInit, OnModuleDestroy {
public kafka: Kafka;
public consumers: ConsumerTopicPair[] = [];
public producer: Producer | null;
private readonly clientId: string;
private readonly bootstrapUrl: string;
private readonly groupId: string;
private readonly offset: boolean;
private readonly consumerTopics: string[];
private readonly producerTopics: string[];
private isConnected = false;
private isConnecting = false;
private isShuttingDown = false;
private readonly INTERVAL_NAME = 'kafkaConnectionCheckAcquiringB2cLimits';
private readonly context = KafkaService.name;
constructor(
private readonly configService: ConfigService,
private readonly logger: DefaultLogger,
private readonly schedulerRegistry: SchedulerRegistry,
private readonly eventEmitter: EventEmitter2,
) {
const s = this.configService.get<KafkaSettings>('kafkaSettings', {
infer: true,
});
this.clientId = s.APP_KAFKA_CLIENT_ID;
this.bootstrapUrl = s.APP_PLATFORM_KAFKA_BOOTSTRAP_SERVER;
this.groupId = s.APP_KAFKA_BOOTSTRAP_GROUP_ID;
this.offset = s.APP_KAFKA_CLIENT_OFFSET;
this.consumerTopics = [
s.APP_KAFKA_CONSUMER_1_TOPIC,
s.APP_KAFKA_CONSUMER_2_TOPIC,
s.APP_KAFKA_CONSUMER_3_TOPIC,
];
this.producerTopics = [
s.APP_KAFKA_PRODUCER_1_TOPIC,
s.APP_KAFKA_PRODUCER_2_TOPIC,
s.APP_KAFKA_PRODUCER_3_TOPIC,
];
}
/**
* На старте модуля запускает проверку подключения к Kafka
*/
onModuleInit() {
const s = this.configService.get<KafkaSettings>('kafkaSettings', {
infer: true,
});
const interval = s.APP_KAFKA_RETRY_INTERVAL_MS;
this.schedulerRegistry.addInterval(
this.INTERVAL_NAME,
setInterval(() => {
this.tryToConnectToKafka().catch((error) => {
this.logger.error(
'Ошибка подключения к Kafka',
String(error),
this.context,
);
});
}, interval),
);
}
/**
* Попытки подключения к Kafka
*/
async tryToConnectToKafka() {
if (this.isConnected || this.isShuttingDown || this.isConnecting) {
return;
}
if (!this.groupId || !this.bootstrapUrl || !this.clientId) {
this.logger.warn(
'Пропуск подключения к Kafka: не заданы обязательные переменные окружения (APP_KAFKA_BOOTSTRAP_GROUP_ID, APP_PLATFORM_KAFKA_BOOTSTRAP_SERVER, APP_KAFKA_CLIENT_ID)',
this.context,
);
return;
}
await this.initializeKafka();
}
/**
* Инициализация Kafka-клиента с настройками подключения, создание consumer и producer, подписка на топики и запуск обработки сообщений
*/
private async initializeKafka() {
if (this.isConnecting || this.isShuttingDown) {
return;
}
this.isConnecting = true;
try {
if (!this.kafka) {
this.kafka = new Kafka({
clientId: this.clientId,
brokers: [this.bootstrapUrl],
retry: { retries: 1 },
});
}
// Создаём consumer для каждого топика
for (let i = 0; i < this.consumerTopics.length; i++) {
if (!this.consumers[i]?.consumer) {
this.consumers[i] = {
consumer: this.kafka.consumer({
groupId: i === 0 ? this.groupId : `${this.groupId}-${i + 1}`,
}),
topic: this.consumerTopics[i],
};
}
}
if (!this.producer) {
this.producer = this.kafka.producer();
await this.producer.connect();
}
this.isConnected = true;
this.setupEventHandlers();
for (const pair of this.consumers) {
await this.runSingleConsumer(pair);
}
this.logger.log('Подключение к Kafka выполнено успешно!', this.context);
} catch (error) {
this.logger.warn(
`Ошибка при подключении к Kafka: ${error.message}`,
this.context,
);
this.logger.warn(
'Будет повторяться в следующем цикле cron',
this.context,
);
this.isConnected = false;
this.consumers = [];
this.producer = null;
} finally {
this.isConnecting = false;
}
}
/**
* Запуск одного consumer - подключение, подписка на топик, обработка входящих сообщений
*/
private async runSingleConsumer(pair: ConsumerTopicPair) {
const { consumer, topic } = pair;
if (!consumer) {
return;
}
try {
await consumer.connect();
await consumer.subscribe({ topic, fromBeginning: this.offset });
dmitriev-aal@VDI-Dmitriev-A:~/Desktop/rshbintech/it-corp/rprul/ecp/acquiring/backend/ms-acquiring-manage-b2c-limits$ sed -n '1,140p' src/infrastructure/dictionaries/dictionaries.service.ts
import { Inject, Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { CACHE_MANAGER } from '@nestjs/cache-manager';
import { Cache } from 'cache-manager';
import { DefaultLogger } from '@rshbintech.rprul.ecp/core-nodejs-starter-lib';
import { HttpService } from '../http/http.service';
import { TokenProvider } from '../auth/token.provider';
import { DepartmentItem } from './types/dictionary.types';
@Injectable()
export class DictionariesService implements OnModuleInit {
private readonly maxRetries = 10;
private readonly retryBaseDelayMs = 2000;
private readonly cacheKey = 'dict:departments';
private readonly refreshIntervalMs = 24 * 60 * 60 * 1000;
private readonly dictionariesBaseUrl: string;
constructor(
private readonly logger: DefaultLogger,
private readonly httpService: HttpService,
private readonly configService: ConfigService,
private readonly tokenProvider: TokenProvider,
@Inject(CACHE_MANAGER) private readonly cacheManager: Cache,
) {
this.logger.setContext(DictionariesService.name);
this.dictionariesBaseUrl = this.configService.get<string>(
'APP_DICTIONARIES_BASE_URL',
'',
);
}
/** При старте модуля: первичная загрузка кэша и запуск периодического обновления */
async onModuleInit(): Promise<void> {
await this.refreshCache();
setInterval(() => this.refreshCache(), this.refreshIntervalMs);
}
/**
* Метод для получения департамента по номеру из закэшированных данных
*
* @param number - номер департамента (например, "3500")
* @returns DepartmentItem или null если департамент не найден в кэше
*/
async getDepartmentByNumber(number: string): Promise<DepartmentItem | null> {
const departments = await this.getCachedDepartments();
return departments.find((d) => d.number === number) ?? null;
}
/**
* Метод для получения закэшированного списка департаментов
*
* @returns массив DepartmentItem из кэша, при промахе - загружает из API
*/
private async getCachedDepartments(): Promise<DepartmentItem[]> {
try {
const cached = await this.cacheManager.get<DepartmentItem[]>(
this.cacheKey,
);
if (cached && cached.length > 0) {
return cached;
}
} catch (error) {
this.logger.error('Ошибка чтения кэша департаментов', error);
}
return this.fetchAndCache();
}
/**
* Метод для принудительного обновления кэша департаментов
*
* Повторяет попытки с экспоненциальной задержкой при ошибках
*
* @returns void
*/
private async refreshCache(): Promise<void> {
this.logger.log('Плановое обновление кэша департаментов');
for (let attempt = 1; attempt <= this.maxRetries; attempt++) {
try {
await this.fetchAndCache();
this.logger.log('Кэш департаментов обновлён');
return;
} catch (error) {
this.logger.error(
`Ошибка обновления (попытка ${attempt}/${this.maxRetries})`,
error,
);
if (attempt < this.maxRetries) {
await new Promise((r) =>
setTimeout(r, this.retryBaseDelayMs * attempt),
);
}
}
}
this.logger.error(
`Не удалось обновить кэш после ${this.maxRetries} попыток - оставлены старые данные`,
);
}
/**
* Метод для загрузки списка департаментов из API и сохранения в кэш
*
* @returns массив DepartmentItem, загруженный из ms-core-dictionaries
*/
private async fetchAndCache(): Promise<DepartmentItem[]> {
const url = `${this.dictionariesBaseUrl}/api/dictionaries/departments/public?pageSize=1000`;
const token = await this.tokenProvider.getToken();
const res = await this.httpService.get<DepartmentItem[]>(url, {
headers: {
'Content-Type': 'application/json',
Authorization: `Bearer ${token.access_token}`,
},
});
const raw = res.data as any;
const items: DepartmentItem[] = Array.isArray(raw)
? raw
: raw?.values
? Object.values(raw.values)
: [];
await this.cacheManager.set(this.cacheKey, items, 0);
this.logger.log(`Закэшировано ${items.length} департаментов`);
return items;
}
}
dmitriev-aal@VDI-Dmitriev-A:~/Desktop/rshbintech/it-corp/rprul/ecp/acquiring/backend/ms-acquiring-manage-b2c-limits$