Загрузка данных


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$