npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

pp-event-bus

v2.0.0

Published

Distributed event bus library with pluggable transports (Redis Streams / DragonflyDB)

Readme

pp-event-bus

Rozproszona magistrala zdarzeń dla mikroserwisów PolskiePolisy. Architektura portów i adapterów zgodna z firmowym przewodnikiem architecture-ddd-srp-guide.md: trzy warstwy DDD z granicami egzekwowanymi przez lintera, jeden port transportu, dwie jego implementacje (Redis Streams / DragonflyDB oraz InMemory do testów). Dostarczanie at-least-once z przeciwciśnieniem, dead-letter queue, deduplikacją i transakcyjnym outboxem. Integracja z NestJS pod osobnym wejściem pakietu. Jedyna zależność produkcyjna: ioredis.

Historia zmian: CHANGELOG.md. Migracja z wersji 1.x: MIGRATION.md.


Spis treści

  1. Główne funkcje
  2. Architektura
  3. Instalacja
  4. Szybki start
  5. Definiowanie zdarzeń
  6. Publikacja i subskrypcja
  7. Wydajność
  8. Gwarancje dostarczenia
  9. Struktura projektu
  10. Zmienne środowiskowe
  11. Integracja z NestJS
  12. Obserwowalność i kontrola zdrowia
  13. Trwałość i konfiguracja DragonflyDB
  14. Bezpieczeństwo
  15. Testowanie
  16. Rozwój
  17. Ograniczenia
  18. Dokumentacja

1. Główne funkcje

  • Wymienny transport — cała komunikacja z brokerem przechodzi przez jeden interfejs EventTransportPort. Zmiana brokera nie dotyka logiki magistrali.
  • Gwarancje at-least-once — potwierdzenie następuje wyłącznie po pomyślnym wykonaniu handlera albo po trwałym zapisie do DLQ. Nie ma ścieżki, w której zdarzenie znika bez śladu.
  • Współbieżność z kolejnością — konfigurowalna równoległość, a kolejność zachowana w obrębie partitionKey (np. per orderId).
  • Przeciwciśnienie — pętla odczytu nie pobiera więcej, niż handlery są w stanie przetworzyć.
  • Dead-letter queue — zdarzenie po wyczerpaniu prób trafia do osobnego strumienia wraz z przyczyną, zamiast być ponawiane w nieskończoność.
  • Deduplikacja — dwufazowe oznaczanie odporne na awarię procesu w trakcie przetwarzania.
  • Transakcyjny outbox — port i relay zamykające lukę „commit do bazy przeszedł, publikacja padła".
  • Retencja czasowa — XTRIM MINID zamiast MAXLEN; zdarzenia żyją ustaloną liczbę dni niezależnie od natężenia ruchu.
  • Integracja z NestJS — EventBusModule, dekorator @EventHandler, wskaźnik zdrowia dla Terminusa. Osobne wejście pakietu, NestJS jako opcjonalna zależność równorzędna.
  • Testowalność — transport pamięciowy pozwala testować handlery bez brokera; identyfikatory zdarzeń i czas da się uczynić deterministycznymi.

2. Architektura

Trzy warstwy DDD, zależności prowadzą wyłącznie do wewnątrz.

2.1. Warstwy

┌──────────────────────────────────────────────────────────────────────┐
│                          INFRASTRUCTURE                              │
│   ┌────────────────┐  ┌────────────────┐  ┌───────────────────────┐  │
│   │ Redis Streams  │  │ InMemory       │  │ Logger, IdGenerator,  │  │
│   │ (ioredis)      │  │ (testy)        │  │ ConsumerIdentity, env │  │
│   └───────┬────────┘  └───────┬────────┘  └──────────┬────────────┘  │
└───────────┼───────────────────┼──────────────────────┼───────────────┘
            │                   │                      │
            ▼                   ▼                      ▼
┌──────────────────────────────────────────────────────────────────────┐
│                            APPLICATION                               │
│   EventBus · MessageDispatcher · ConsumerSupervisor                   │
│   PartitionedQueue · AckBatcher · OutboxRelay                        │
│   SubscriptionRegistry · SubscriptionPolicy                          │
│   Zna wyłącznie porty - ani jednego odwołania do ioredis.            │
└───────────────────────────────┬──────────────────────────────────────┘
                                ▼
┌──────────────────────────────────────────────────────────────────────┐
│                              DOMAIN                                  │
│   Event · EventEnvelope · RetryPolicy · błędy                        │
│   Porty: EventTransport, Outbox, DeduplicationStore,                 │
│          Metrics, Logger, Clock, IdGenerator                          │
│   Zero zależności technicznych - także bez modułów Node.             │
└──────────────────────────────────────────────────────────────────────┘

Granice są egzekwowane regułami ESLint (no-restricted-imports), a nie tylko opisane. Warstwa domenowa nie może zaimportować żadnego modułu Node — łącznie z node:crypto i node:os. Potrzebne prymitywy (generowanie identyfikatorów, czas) dostaje przez porty konfigurowane w punkcie kompozycji. Reguła jest zweryfikowana testem negatywnym: sonda importująca node:crypto w src/domain musi wywalić lintera.

2.2. Port transportu

export interface EventTransportPort {
  publish(records: readonly PublishRecord[]): Promise<readonly PublishResult[]>;
  subscribe(spec: SubscriptionSpec): Promise<TransportSubscription>;
  close(): Promise<void>;
}

Sprawy specyficzne dla brokera — retencja, tworzenie grup, format kluczy, pula połączeń — są wewnętrzną odpowiedzialnością implementacji i nie pojawiają się w kontrakcie.

Poprawność abstrakcji jest weryfikowana wspólnym zestawem testów kontraktowych (src/infrastructure/__tests__/transport-contract.suite.ts), który przechodzą zarówno Redis Streams, jak i transport pamięciowy. Gdyby kontrakt przeciekał szczegółami Redisa, wersja pamięciowa nie mogłaby go spełnić.

Aby dodać nowy transport, wystarczy zaimplementować port:

const bus = await EventBus.create({ transport: new MojTransport(config) });

3. Instalacja

npm install pp-event-bus

Zależności równorzędne (@nestjs/common, @nestjs/core, reflect-metadata) są opcjonalne — potrzebne wyłącznie przy korzystaniu z wejścia pp-event-bus/nestjs.


4. Szybki start

import { createRedisEventBus, Event, loggerFromEnv } from 'pp-event-bus';

interface OrderPayload {
  orderId: string;
  customerName: string;
  amount: number;
}

class OrderCreatedEvent extends Event<OrderPayload> {
  public static override readonly version = '1.0.0';
}

const bus = await createRedisEventBus({
  url: process.env.REDIS_URL,
  logger: loggerFromEnv(),
});

// Konsument
await bus.subscribe(
  OrderCreatedEvent,
  async (event) => {
    // event.payload jest w pełni otypowany - bez powtarzania definicji
    await billing.charge(event.payload.orderId, event.payload.amount);
  },
  { group: 'rozliczenia' },
);

// Producent
await bus.emit(
  new OrderCreatedEvent({
    orderId: 'ZAM-2024-001',
    customerName: 'Jan Kowalski',
    amount: 249.99,
  }),
);

// Zamknięcie - dokańcza zdarzenia w trakcie przetwarzania
await bus.close();

5. Definiowanie zdarzeń

Zdarzenie to klasa dziedzicząca po Event<TPayload>. Nazwą zdarzenia jest nazwa klasy, a wersja trafia do koperty.

class OrderCreatedEvent extends Event<OrderPayload> {
  public static override readonly version = '1.2.0';
}

const event = new OrderCreatedEvent(payload)
  .setMetadata('correlationId', ctx.correlationId)
  .setMetadata('source', 'zamowienia-api');

event.id; // UUID v4
event.name; // 'OrderCreatedEvent'
event.version; // '1.2.0'
event.occurredAt; // znacznik czasu (ms)
event.payload; // dane biznesowe

5.1. Wersjonowanie

Wersja nie wpływa na routing — wszystkie wersje danego zdarzenia trafiają do tego samego strumienia. Konsument decyduje, co zrobić z wersją, której nie zna:

await bus.subscribe(
  OrderCreatedEvent,
  async (event) => {
    if (!event.version.startsWith('1.')) {
      logger.warn('Nieobsługiwana wersja zdarzenia', { version: event.version });
      return;
    }
    // ...
  },
  { group: 'rozliczenia' },
);

5.2. Minifikacja a nazwa strumienia

Nazwą zdarzenia jest nazwa klasy (constructor.name). Budowanie produkcyjne z minifikacją bez keep_classnames (esbuild, terser, webpack) zamienia ją w jednoliterowy identyfikator: producent publikuje wtedy na strumień a, a konsument z innego builda słucha b. Zero błędów, zero zdarzeń.

Biblioteka ostrzega, gdy rozwiązana nazwa wygląda na zmanglowaną, ale heurystyka nie jest niezawodna. Dla zdarzeń wymienianych między serwisami ustawiaj nazwę jawnie:

class OrderCreatedEvent extends Event<OrderPayload> {
  public static override readonly eventName = 'OrderCreatedEvent';
}

5.3. Zmiana nazwy klasy

Nazwa klasy jest kontraktem między serwisami. Jeśli klasa musi zostać przemianowana, zachowaj dotychczasową nazwę kontraktu:

class OrderPlacedEvent extends Event<OrderPayload> {
  // Konsumenci nadal czytają strumień pod starą nazwą
  public static override readonly eventName = 'OrderCreatedEvent';
}

6. Publikacja i subskrypcja

6.1. Publikacja

// Pojedyncze zdarzenie
await bus.emit(new OrderCreatedEvent(payload));

// Partia - jeden pipeline zamiast rundy na każde zdarzenie
await bus.emitMany([event1, event2, event3]);

6.2. Subskrypcja

const subscription = await bus.subscribe(OrderCreatedEvent, handler, {
  group: 'rozliczenia', // wymagane - grupa konsumencka
  consumer: 'replika-1', // opcjonalne - domyślnie stabilna nazwa procesu
  concurrency: 16, // ile zdarzeń równolegle
  partitionKey: (e) => e.payload.orderId,
  retry: { maxAttempts: 5, initialDelayMs: 1000, jitter: true },
  deadLetter: true,
  deduplication: { ttlMs: 86_400_000 },
  startFrom: 'beginning',
  batchSize: 64,
  onError: (error, event) => sentry.captureException(error, { extra: { event } }),
  onFatalError: (error) => alerting.page('Pętla konsumpcji przerwana', error),
  sanitizeError: (error) => error.message.replace(/\S+@\S+/g, '***'),
});

await subscription.unsubscribe();

6.3. Grupy konsumenckie

Każda grupa dostaje własną kopię zdarzenia (rozgłaszanie 1:N). W obrębie jednej grupy zdarzenia są rozdzielane między repliki (rozkładanie obciążenia).

Dwie subskrypcje tego samego zdarzenia w tej samej grupie i w tym samym procesie są odrzucane błędem ConfigurationError. Dzieliłyby tożsamość konsumenta, przez co cykl odzyskiwania jednej mógłby przejąć wiadomość przetwarzaną przez drugą — i ten sam handler wykonałby się dwukrotnie.

  • Chcesz, aby oba handlery dostały każde zdarzenie → użyj różnych nazw grup.
  • Chcesz świadomie rozłożyć obciążenie → podaj jawne, różne nazwy konsumentów.

7. Wydajność

7.1. Współbieżność i kolejność

Domyślnie zdarzenia przetwarzane są równolegle (concurrency: 16). partitionKey pozwala zachować kolejność tam, gdzie jest istotna:

await bus.subscribe(OrderCreatedEvent, handler, {
  group: 'rozliczenia',
  concurrency: 16,
  partitionKey: (event) => event.payload.orderId,
});

// zdarzenia różnych zamówień       -> równolegle (do 16 naraz)
// zdarzenia tego samego zamówienia -> ściśle po kolei

Bez partitionKey obowiązuje pełna równoległość — handlery muszą być wtedy wzajemnie niezależne.

7.2. Optymalizacje działające automatycznie

| Mechanizm | Efekt | | ----------------------------- | ------------------------------------------------------------------------ | | Publikacja wsadowa (pipeline) | emitMany wysyła całą partię jednym przejściem do brokera | | Potwierdzanie wsadowe | Zmierzone: 500 zdarzeń → 16 poleceń XACK zamiast 500 | | Odzyskiwanie porzuconych | Dwa polecenia na cykl niezależnie od liczby wiadomości (zamiast N+1) | | Leniwy kontekst logów | Obiekty metadanych powstają dopiero po przejściu progu poziomu logowania |

Potwierdzenia są scalane w oknie 5 ms lub do rozmiaru batchSize — co nastąpi pierwsze. Scalanie nie zamienia potwierdzenia w operację „wyślij i zapomnij": zdarzenie jest uznane za przetworzone dopiero, gdy jego partia dotrze do brokera, a bufor jest opróżniany przed zamknięciem konsumenta.

7.3. Przeciwciśnienie

Pętla odczytu nie pobiera kolejnej partii, dopóki liczba zdarzeń w obróbce nie zejdzie poniżej dwukrotności concurrency. Bez tego progu wolne handlery prowadziłyby do nieograniczonego wzrostu kolejki w pamięci procesu, a wszystkie pobrane wiadomości zajmowałyby miejsce na liście oczekujących brokera.

W praktyce batchSize steruje wielkością pojedynczego odczytu, a concurrency faktyczną przepustowością. Ustawianie batchSize znacznie powyżej concurrency nie zwiększa wydajności — zwiększa tylko liczbę wiadomości czekających w liście oczekujących.


8. Gwarancje dostarczenia

Biblioteka realizuje at-least-once.

8.1. Warstwy zabezpieczeń

| Rodzaj awarii | Mechanizm ochronny | | ----------------------------------------------------- | ----------------------------------------------------------------------------- | | Awaria producenta między zapisem do bazy a publikacją | transakcyjny outbox | | Awaria brokera | snapshoty + replika (konfiguracja) | | Awaria konsumenta w trakcie przetwarzania | lista niepotwierdzonych + przejmowanie porzuconych | | Trwale błędny handler | dead-letter queue po wyczerpaniu prób | | Duplikaty wynikające z ponowień | deduplikacja |

Zdarzenie jest potwierdzane wyłącznie po pomyślnym wykonaniu handlera albo po trwałym zapisaniu go w DLQ.

Retencja jest czasowa, nie ilościowa. Zdarzenia żyją ustaloną liczbę dni (domyślnie 7) niezależnie od natężenia ruchu. Przycinaniem zajmuje się osobny proces uruchamiany przez transport — XADD nie kasuje niczego przy publikacji.

const bus = await createRedisEventBus({
  url: process.env.REDIS_URL,
  retention: {
    maxAgeMs: 14 * 24 * 60 * 60 * 1000, // 14 dni
    intervalMs: 5 * 60 * 1000, // co ile przycinać
  },
});

Alert na zaległość grup jest obowiązkowy. Grupa nieaktywna dłużej niż okno retencji straci zdarzenia. Metryka consumerLag istnieje właśnie po to, by wykryć to zawczasu.

8.2. Dead-letter queue

Po wyczerpaniu prób zdarzenie trafia do osobnego strumienia {namespace}:{NazwaZdarzenia}:dlq wraz z przyczyną, liczbą prób i znacznikiem czasu, a oryginał zostaje potwierdzony.

redis-cli XRANGE 'event-stream:OrderCreatedEvent:dlq' - + COUNT 10

Wyłączenie DLQ (deadLetter: false) sprawia, że nieprzetworzone zdarzenie zostaje niepotwierdzone i wraca przez mechanizm przejmowania porzuconych — w nieskończoność. Używaj świadomie.

8.3. Deduplikacja

const bus = await createRedisEventBus({
  url: process.env.REDIS_URL,
  deduplication: true,
});

await bus.subscribe(OrderCreatedEvent, handler, {
  group: 'rozliczenia',
  deduplication: { ttlMs: 24 * 60 * 60 * 1000 },
});

Oznaczanie jest dwufazowe: przed przetworzeniem zapisywane jest oznaczenie krótkotrwałe, a dopiero po powodzeniu przedłużane do pełnego czasu życia. Dzięki temu awaria procesu w trakcie przetwarzania nie blokuje ponownego przetworzenia zdarzenia — naiwna implementacja zamieniłaby w tym miejscu ochronę przed duplikatami w źródło utraty danych.

8.4. Transakcyjny outbox

Broker nie uczestniczy w transakcji bazy danych. Awaria procesu między zatwierdzeniem transakcji a publikacją oznacza trwałą utratę zdarzenia — outbox zamyka tę lukę.

import { OutboxRelay, type OutboxPort } from 'pp-event-bus';

// 1. Implementacja portu po stronie serwisu (schemat bazy należy do aplikacji)
class PrismaOutbox implements OutboxPort<Prisma.TransactionClient> {
  async enqueue(tx, records) {
    await tx.outbox.createMany({
      data: records.map((r) => ({ topic: r.topic, envelope: r.envelope })),
    });
  }

  async fetchPending(limit) {
    // Zalecane: FOR UPDATE SKIP LOCKED, aby repliki relaya się nie dublowały
    const rows = await this.prisma.$queryRaw`
      SELECT id, topic, envelope, attempts FROM outbox
      WHERE dispatched_at IS NULL
      ORDER BY id LIMIT ${limit}
      FOR UPDATE SKIP LOCKED`;
    return rows.map((row) => ({ recordId: String(row.id), ...row }));
  }

  async markDispatched(ids) {
    /* ... */
  }
  async markFailed(id, reason) {
    /* ... */
  }
}

// 2. Zapis atomowy ze zmianą biznesową
await prisma.$transaction(async (tx) => {
  await tx.order.create({ data: order });
  await outbox.enqueue(tx, [{ topic: 'OrderCreatedEvent', envelope: event.toEnvelope() }]);
});

// 3. Relay przy starcie aplikacji
const relay = new OutboxRelay({ outbox, transport, logger, metrics, clock });
relay.start();

Awaria między publikacją a oznaczeniem rekordu skutkuje ponowną publikacją, czyli duplikatem — usuwa go deduplikacja po stronie konsumenta. To świadomy kompromis: duplikat da się odfiltrować, utraty nie da się odwrócić.


9. Struktura projektu

event-bus/
├── docker-compose.yml                      # DragonflyDB (dev + testy integracyjne), Redis pod profilem `redis`
├── eslint.config.js                        # Flat config; reguły granic DDD per warstwa
├── tsconfig.json                           # Build (ES2022, strict + noUncheckedIndexedAccess)
├── tsconfig.eslint.json                    # Lint/typecheck obejmujący też testy
├── jest.setup.js                           # Konfiguracja kontekstu zdarzeń dla testów
│
└── src/
    ├── index.ts                            # PUBLICZNE API + punkt kompozycji (konfiguruje kontekst zdarzeń)
    ├── create-event-bus.ts                 # Fabryki: createRedisEventBus, createInMemoryEventBus
    │
    ├── domain/                             # ZERO zależności technicznych (także bez modułów Node)
    │   ├── event/
    │   │   ├── event.ts                    # Event<TPayload> - klasa bazowa po stronie producenta
    │   │   ├── event-envelope.ts           # Kontrakt on-the-wire + kodowanie/dekodowanie (czyta też format v1)
    │   │   └── event-context.ts            # Uchwyt na IdGenerator i Clock; configureEventContext()
    │   ├── ports/
    │   │   ├── event-transport.port.ts     # GŁÓWNY PORT - jedyne wyjście na brokera
    │   │   ├── outbox.port.ts              # Transakcyjny outbox (implementacja po stronie serwisu)
    │   │   ├── deduplication-store.port.ts # Dwufazowe oznaczanie: markIfFirst / confirm / release
    │   │   ├── metrics.port.ts             # MetricsPort + NoopMetrics
    │   │   ├── logger.port.ts              # LoggerPort z leniwym kontekstem
    │   │   ├── clock.port.ts               # ClockPort + SystemClock
    │   │   └── id-generator.port.ts        # IdGeneratorPort
    │   ├── policy/retry-policy.ts          # Backoff wykładniczy z jitterem; worstCaseTotalDelayMs()
    │   └── errors/event-bus.error.ts       # Hierarchia błędów + maskowanie danych wrażliwych w kontekście
    │
    ├── application/                        # Zna wyłącznie porty z domain/
    │   ├── event-bus.ts                    # Fasada: emit, emitMany, subscribe, health, close
    │   ├── event-bus-defaults.ts           # Ustawienia domyślne (odczyt zmiennych środowiskowych)
    │   ├── subscription-registry.ts        # Strażnik unikalności (topic, group, consumer)
    │   ├── subscription-policy.ts          # Walidacja rozmiaru, budżetu ponowień, planu deduplikacji
    │   ├── dispatch/
    │   │   ├── message-dispatcher.ts       # Cykl życia JEDNEJ wiadomości: dedup → handler → retry → DLQ
    │   │   ├── consumer-supervisor.ts      # Pętla odczytu, wznawianie, przejmowanie porzuconych, stan zdrowia
    │   │   ├── partitioned-queue.ts        # Współbieżność + kolejność per partitionKey + przeciwciśnienie
    │   │   └── ack-batcher.ts              # Scalanie potwierdzeń (Acknowledger, ImmediateAcknowledger)
    │   └── outbox/outbox-relay.ts          # Publikacja rekordów oczekujących z outboxa
    │
    ├── infrastructure/                     # Jedyne miejsce, które wie, jak działa broker
    │   ├── redis-streams/
    │   │   ├── redis-stream-transport.ts   # Implementacja portu: publikacja, grupy, fabryka subskrypcji
    │   │   ├── redis-stream-subscription.ts# Odczyt, potwierdzanie, przejmowanie, DLQ, kwarantanna
    │   │   ├── redis-connection-factory.ts # TLS, sentinel, cluster; redakcja hasła; idempotentne rozłączanie
    │   │   ├── redis-connection-diagnostics.ts # Awarie połączeń do loggera aplikacji, bez zalewania logów
    │   │   ├── redis-retention-janitor.ts  # XTRIM MINID + metryki głębokości i zaległości
    │   │   ├── redis-deduplication-store.ts# SET NX PX / PEXPIRE / DEL
    │   │   ├── redis-reply-parser.ts       # Defensywne parsowanie odpowiedzi brokera
    │   │   ├── redis-error-predicates.ts   # Rozpoznawanie BUSYGROUP / NOGROUP
    │   │   └── redis-keys.ts               # Przestrzeń nazw kluczy
    │   ├── in-memory/in-memory-transport.ts# Druga implementacja portu - testy bez brokera
    │   ├── identity/                       # CryptoIdGenerator, buildConsumerIdentity
    │   ├── observability/console-logger.ts # Logger z poziomami i leniwym kontekstem
    │   ├── config/env.ts                   # Odczyt zmiennych środowiskowych
    │   └── __tests__/transport-contract.suite.ts  # Zestaw kontraktowy dla KAŻDEJ implementacji portu
    │
    └── nestjs/                             # Subpath export `pp-event-bus/nestjs`
        ├── event-bus.module.ts             # forRoot / forRootAsync
        ├── event-bus.explorer.ts           # Wykrywanie @EventHandler, bootstrap i shutdown
        ├── event-bus.health.ts             # Wskaźnik dla @nestjs/terminus
        ├── event-handler.decorator.ts      # @EventHandler
        └── event-bus.constants.ts          # Token EVENT_BUS

10. Zmienne środowiskowe

Każda opcja podana wprost w kodzie ma pierwszeństwo przed zmienną środowiskową, a ta przed wartością domyślną.

10.1. Połączenie i strumienie

| Zmienna | Domyślna | Zakres | Opis | | ----------------------------- | -------------- | ---------------- | ----------------------------------- | | EVENT_BUS_NAMESPACE | event-stream | tekst niepusty | Prefiks kluczy — izolacja środowisk | | EVENT_BUS_CONSUMER_TIMEOUT | 5000 | 100 .. 300 000 | Czas blokady odczytu (ms) | | EVENT_BUS_BATCH_SIZE | 64 | 1 .. 10 000 | Ile wiadomości na jeden odczyt | | EVENT_BUS_MAX_PAYLOAD_BYTES | 52428800 | 1 KiB .. 512 MiB | Limit rozmiaru koperty (50 MiB) | | EVENT_BUS_CONCURRENCY | 16 | 1 .. 1024 | Zdarzenia przetwarzane równolegle |

10.2. Retencja

| Zmienna | Domyślna | Zakres | Opis | | --------------------------------- | ----------- | ---------------- | --------------------------- | | EVENT_BUS_RETENTION_MAX_AGE_MS | 604800000 | 1 min .. 365 dni | Okno retencji (7 dni) | | EVENT_BUS_RETENTION_INTERVAL_MS | 300000 | 1 s .. 24 h | Co ile przycinać strumienie |

10.3. Odzyskiwanie porzuconych wiadomości

| Zmienna | Domyślna | Zakres | Opis | | ------------------------------------- | -------- | ----------- | ------------------------------------ | | EVENT_BUS_PENDING_CHECK_INTERVAL_MS | 30000 | 1 s .. 1 h | Co ile szukać porzuconych wiadomości | | EVENT_BUS_PENDING_IDLE_TIME_MS | 60000 | 1 s .. 24 h | Po jakim czasie przejąć porzuconą | | EVENT_BUS_RECLAIM_BATCH_SIZE | 100 | 1 .. 10 000 | Limit przejęć w jednym cyklu |

10.4. Ponowienia i deduplikacja

| Zmienna | Domyślna | Zakres | Opis | | ---------------------------------- | ---------- | ------------- | ---------------------------------- | | EVENT_BUS_MAX_ATTEMPTS | 5 | 1 .. 100 | Prób przed skierowaniem do DLQ | | EVENT_BUS_RETRY_INITIAL_DELAY_MS | 1000 | 1 ms .. 1 min | Opóźnienie pierwszego ponowienia | | EVENT_BUS_RETRY_MAX_DELAY_MS | 60000 | 1 ms .. 1 h | Górne ograniczenie opóźnienia | | EVENT_BUS_DEDUP_TTL_MS | 86400000 | 1 s .. 30 dni | Czas życia oznaczenia deduplikacji |

10.5. Logowanie

| Zmienna | Domyślna | Akceptowane wartości | | ----------- | -------- | --------------------------------------------------------- | | LOG_LEVEL | log | debug | log | warn | error albo liczba 0-9 |

LOG_LEVEL jest jedyną zmienną bez prefiksu EVENT_BUS_ i jako jedyna nie podlega twardej walidacji — należy do aplikacji, nie do biblioteki. Nierozpoznana wartość daje ostrzeżenie i poziom log, nigdy błąd startu.

Skala liczbowa jest odwzorowaniem konwencji używanej w naszych serwisach (core/auth-service), więc jeden LOG_LEVEL steruje całym stosem:

| Wartość | Poziom biblioteki | Znaczenie | | ------- | ----------------- | --------------------------------- | | 0-3 | error | tylko krytyczne usterki | | 4-7 | log | typowy prod: log, warn, error | | 8-9 | debug | diagnoza (9 = verbose w NestJS) |

REDIS_URL jest konwencją, a nie zmienną odczytywaną przez bibliotekę — adres połączenia przekazujesz jawnie ({ url: process.env.REDIS_URL }). Dzięki temu aplikacja waliduje własną konfigurację w jednym miejscu.

Wartość spoza zakresu zatrzymuje start aplikacji. Dotyczy wyłącznie zmiennych z prefiksem EVENT_BUS_ - biblioteka nie waliduje nazw, których nie jest właścicielem. Walidacja wykonuje się w createRedisEventBus / EventBusModule i zgłasza wszystkie błędne zmienne naraz. Wcześniej EVENT_BUS_RETENTION_MAX_AGE_MS=7d parsowało się po cichu jako 7 milisekund retencji, a pierwszy przebieg janitora kasował całą historię strumieni.

Suma opóźnień ponowień musi mieścić się w oknie przetwarzania, czyli EVENT_BUS_PENDING_IDLE_TIME_MS pomniejszonym o margines (10% progu, nie mniej niż sekunda). Inaczej inna replika uznałaby wciąż przetwarzaną wiadomość za porzuconą i obsłużyła ją równolegle. Biblioteka sprawdza tę relację przy tworzeniu subskrypcji i odrzuca niespójną konfigurację.


11. Integracja z NestJS

import { EventBusModule, EventHandler, EVENT_BUS } from 'pp-event-bus/nestjs';

@Module({
  imports: [
    EventBusModule.forRootAsync({
      imports: [ConfigModule],
      inject: [ConfigService],
      useFactory: (config: ConfigService) => ({
        url: config.getOrThrow('REDIS_URL'),
        namespace: `pp:${config.get('NODE_ENV')}`,
        deduplication: true,
      }),
    }),
  ],
})
export class AppModule {}

@Injectable()
export class BillingHandlers {
  @EventHandler(OrderCreatedEvent, {
    group: 'rozliczenia',
    concurrency: 16,
    partitionKey: (event) => event.payload.orderId,
  })
  public async onOrderCreated(event: ReceivedEvent<OrderPayload>): Promise<void> {
    await this.billing.charge(event.payload.orderId);
  }
}

// Publikacja
@Injectable()
export class OrdersService {
  constructor(@Inject(EVENT_BUS) private readonly eventBus: EventBus) {}
}

Subskrypcje są zakładane w onApplicationBootstrap, a magistrala zamykana w onApplicationShutdown — włącz app.enableShutdownHooks(), aby zdarzenia w trakcie przetwarzania były dokańczane przed wyłączeniem poda.


12. Obserwowalność i kontrola zdrowia

12.1. Metryki

const bus = await createRedisEventBus({
  url: process.env.REDIS_URL,
  metrics: {
    eventsPublished: (topic, count) => counter.inc({ topic }, count),
    eventConsumed: (topic, group, ms) => histogram.observe({ topic, group }, ms),
    eventFailed: (topic, group, attempt) => counter.inc({ topic, group, attempt }),
    eventDeadLettered: (topic, group) => counter.inc({ topic, group }),
    eventDeduplicated: (topic, group) => counter.inc({ topic, group }),
    eventReclaimed: (topic, group, count) => counter.inc({ topic, group }, count),
    streamDepth: (topic, depth) => gauge.set({ topic }, depth),
    consumerLag: (topic, group, lag) => gauge.set({ topic, group }, lag),
  },
});

Metryki wymagające alertów:

| Metryka | Dlaczego | | ------------------- | ------------------------------------------------------------- | | consumerLag | Grupa zalegająca dłużej niż okno retencji traci zdarzenia | | eventDeadLettered | Zdarzenia wymagające interwencji człowieka | | eventReclaimed | Rosnąca wartość oznacza padające lub zawieszające się repliki | | streamDepth | Wzrost względem maxmemory brokera |

12.2. Kontrola zdrowia

Cicha śmierć konsumenta jest najgroźniejszą awarią tej biblioteki: pod nadal odpowiada, ale przestaje odbierać zdarzenia. health() czyni to widocznym.

bus.health();
// {
//   healthy: false,
//   closed: false,
//   consumers: 1,
//   subscriptions: [{
//     topic: 'OrderCreatedEvent', group: 'rozliczenia',
//     running: false,          // także w trakcie odczekiwania między próbami
//     stopped: false,          // odróżnia awarię od świadomego wyłączenia
//     pending: 12,
//     consecutiveFailures: 4,
//     lastError: 'Błąd odczytu ze strumienia',
//     lastFailureAt: 1700000000000,
//   }],
// }

W NestJS dostępny jest gotowy wskaźnik dla @nestjs/terminus:

import { EventBusHealthIndicator } from 'pp-event-bus/nestjs';

@Controller('health')
export class HealthController {
  constructor(
    private readonly health: HealthCheckService,
    private readonly eventBus: EventBusHealthIndicator,
  ) {}

  @Get()
  @HealthCheck()
  public check() {
    return this.health.check([() => this.eventBus.isHealthy('event-bus')]);
  }
}

Podłącz to do readiness probe z failureThreshold >= 3. Bez wskaźnika Kubernetes kieruje ruch do poda, który nie przetwarza zdarzeń; bez progu pojedynczy przejściowy błąd odczytu niepotrzebnie wywróci poda.

Co dokładnie znaczy „niezdrowa". healthy jest fałszywe, gdy magistrala została zamknięta, gdy którakolwiek pętla odczytu nie pracuje (łącznie z odczekiwaniem między próbami wznowienia) albo gdy pod miał kiedyś konsumentów i wszystkich stracił. Ten ostatni warunek jest istotny: [].every() zwraca true, więc bez niego pod bez ani jednego konsumenta raportowałby pełne zdrowie. Pod czysto publikujący, który nigdy nie subskrybował, pozostaje zdrowy.

Wskaźnik Terminusa zwraca listę degraded obejmującą również subskrypcje pracujące, ale z niezerowym licznikiem niepowodzeń - konsument wznawiany w pętli jest widoczny, zanim przestanie działać zupełnie.

12.3. Awarie połączenia z brokerem

Awarie połączeń trafiają do loggera podanego w konfiguracji:

WARN  Błąd połączenia z brokerem  { purpose: 'consumer', connectionName: 'pp-event-bus:consumer:rozliczenia:worker-1',
                                    code: 'ECONNREFUSED', error: 'connect ECONNREFUSED 10.0.0.1:6379' }
DEBUG Ponowny błąd połączenia z brokerem   ← kolejne próby tej samej awarii
LOG   Połączenie z brokerem odtworzone     { previousError: 'ECONNREFUSED' }
ERROR Broker odmówił dostępu               { code: 'WRONGPASS invalid username-password pair' }

| Sytuacja | Poziom | Dlaczego | | ------------------------------------------------- | ------- | -------------------------------------------------------------- | | Zerwane lub odrzucone gniazdo (ECONNREFUSED, …) | warn | Usterka przemijająca — ioredis wraca sam | | Odmowa dostępu (WRONGPASS, NOAUTH, NOPERM) | error | Ponowienia nic nie zmienią bez poprawienia konfiguracji | | Powtórzenie tej samej przyczyny | debug | Jedna awaria to zdarzenie co kilka sekund na każdym połączeniu | | Powrót połączenia po awarii | log | Domyka wpis o awarii; sam start nie jest logowany |

Bez logger w konfiguracji biblioteka nie zakłada własnego nasłuchu i awarie wypisuje ioredis: [ioredis] Unhandled error event: … przez console.error, z pominięciem loggera aplikacji i bez informacji, którego połączenia dotyczą. Wyciszenie ich bez miejsca docelowego byłoby gorsze niż ten komunikat — dlatego podanie logger jest jedynym sposobem, by przejęły je logi aplikacji.

To samo dotyczy klienta wstrzykniętego przez client — jego zdarzenia obsługuje właściciel, biblioteka nie ingeruje w cudzą instancję.


13. Trwałość i konfiguracja DragonflyDB

Fakt kluczowy: DragonflyDB nie implementuje AOF. Redisowe appendonly yes nie ma tu odpowiednika — trwałość opiera się wyłącznie na snapshotach, więc maksymalne okno utraty danych równa się odstępowi między nimi. Jeśli wymagane jest twardsze zabezpieczenie, źródłem prawdy musi być outbox w bazie z dziennikiem transakcji.

13.1. Konfiguracja

services:
  dragonfly:
    # Magazyn danych zawsze pinujemy do konkretnej wersji
    image: 'docker.dragonflydb.io/dragonflydb/dragonfly:v1.27.0'
    restart: always
    ulimits:
      memlock: -1
    command:
      - '--logtostderr'
      - '--dir=/data'
      - '--dbfilename=dump'
      - '--snapshot_cron=* * * * *' # okno utraty = 60 s
      - '--maxmemory=4gb'
      - '--cache_mode=false' # KRYTYCZNE - nigdy nie eksmituj kluczy
      - '--proactor_threads=4'
    volumes:
      - dragonfly_data:/data
    stop_grace_period: 60s # czas na snapshot przy SIGTERM
    healthcheck:
      test: ['CMD', 'redis-cli', '-p', '6379', 'ping']
      interval: 10s
      timeout: 5s
      retries: 5

13.2. Uzasadnienie ustawień

| Ustawienie | Dlaczego | | ----------------------- | --------------------------------------------------------------------------------------------------------------------- | | --dir=/data + wolumen | Bez tego snapshoty lądują w warstwie kontenera i giną przy odtworzeniu | | --snapshot_cron | Jedyny mechanizm trwałości. Odstęp = maksymalne okno utraty danych | | --cache_mode=false | Przy true Dragonfly eksmituje klucze pod presją pamięci, czyli po cichu kasuje strumienie. Ustawiamy jawnie | | --maxmemory | Po wyczerpaniu limitu XADD zwraca błąd, publikacja zawodzi, outbox ponawia. To właściwy tryb awarii, nie utrata | | stop_grace_period | Snapshot powstaje przy gracefulnym SIGTERM. Za krótki czas = SIGKILL = utrata wszystkiego od ostatniego snapshotu | | Pin wersji | Zestaw flag zmieniał się między wydaniami |

13.3. Dostrajanie wydajności

DragonflyDB dzieli przestrzeń kluczy na shardy obsługiwane przez osobne wątki. Ma to dwie konsekwencje dla tej biblioteki:

  • Wiele typów zdarzeń działa na korzyść. Każdy strumień to osobny klucz, więc trafia na potencjalnie inny shard i jest obsługiwany równolegle. Ustaw --proactor_threads mniej więcej na liczbę dostępnych rdzeni.
  • Potwierdzenia trafiają w jeden shard. XACK z wieloma identyfikatorami dotyczy jednego strumienia, więc scalanie potwierdzeń obciąża dokładnie jeden wątek zamiast rozpraszać pracę — i dlatego jest tak skuteczne.

Połączenie współdzielone obsługuje publikację, potwierdzenia i operacje administracyjne. Przy bardzo dużym ruchu publikacyjnym potwierdzenia mogą czekać w kolejce za XADD; jeśli to zauważysz w metrykach, rozdziel role, tworząc osobną instancję transportu dla producenta i konsumenta.

13.4. Produkcja (Kubernetes)

  • StatefulSet + PVC, nie Deployment z emptyDir.
  • terminationGracePeriodSeconds: 60 spójnie z czasem snapshotu.
  • Co najmniej jedna replika (--replicaof) z monitoringiem opóźnienia replikacji — obniża prawdopodobieństwo awarii mocniej niż zagęszczanie snapshotów.
  • Alerty: zużycie pamięci wobec maxmemory, wiek ostatniego snapshotu, długość strumieni oraz consumerLag.

13.5. Zweryfikowane na DragonflyDB v1.27.0

| Element | Wynik | | ------------------------------------------ | ------------------------------------------------------------------- | | --snapshot_cron, --dir, --dbfilename | przyjmowane, widoczne w CONFIG GET | | XPENDING ... IDLE | działa, zwraca [id, konsument, idle, licznik] | | XCLAIM bez JUSTID | inkrementuje licznik dostarczeń | | XINFO GROUPS | udostępnia pole lag, więc metryka zaległości działa | | redis-cli w obrazie | obecny, healthcheck działa | | --cache_mode | przyjmowane jako flaga startowa, nieodczytywalne przez CONFIG GET |

Weryfikacja trwałości po zmianie konfiguracji:

docker compose exec dragonfly redis-cli INFO persistence
docker compose restart dragonfly
docker compose exec dragonfly redis-cli XLEN 'event-stream:OrderCreatedEvent'

14. Bezpieczeństwo

  • Poświadczenia nie trafiają do logów. URL brokera jest redagowany (redis://jan:***@host), a kontekst każdego błędu biblioteki automatycznie maskuje pola o nazwach wskazujących na dane wrażliwe (password, token, apiKey i podobne).
  • TLS działa. Schemat rediss:// faktycznie włącza szyfrowanie (w wersji 1.x był ignorowany).
  • Payloady nie są logowane na żadnym poziomie.
  • DLQ przechowuje pełne koperty, łącznie z payloadem, przez całe okno retencji — to konieczne do odtworzenia zdarzenia, ale istotne przy ocenie zgodności z RODO. Wpisy są przycinane razem ze strumieniem głównym.
  • Komunikat błędu handlera trafia do DLQ. Jeśli może zawierać dane osobowe, podaj hak sanityzacji:
await bus.subscribe(OrderCreatedEvent, handler, {
  group: 'rozliczenia',
  sanitizeError: (error) => error.message.replace(/\S+@\S+/g, '***'),
});
  • Limit rozmiaru zdarzenia (domyślnie 50 MiB) chroni brokera przed wyczerpaniem pamięci przez pojedynczą publikację. Zdarzenie ponad limit jest odrzucane po stronie producenta, zanim dotrze do brokera:
const bus = await createRedisEventBus({
  url: process.env.REDIS_URL,
  maxPayloadBytes: 70 * 1024 * 1024, // domyślnie 50 MiB
});

Limit mierzy kopertę po serializacji, nie plik

To rozróżnienie decyduje o tym, czy załącznik przejdzie. Kodowanie powiększa dane, więc dokument o wadze 50 MB nie mieści się w limicie 50 MiB:

| Sposób przekazania | Rozmiar koperty | Mnożnik | Czas JSON.stringify | | ------------------ | --------------- | ------- | --------------------- | | base64 (typowe) | 66,7 MB | ×1,33 | 280 ms | | Buffer wprost | 200,0 MB | ×4,00 | 3196 ms |

Wartości zmierzone dla pliku 50 MB. Buffer przechodzi przez JSON.stringify jako tablica liczb — stąd czterokrotny narzut i ponad trzy sekundy zablokowanej pętli zdarzeń. Jeśli naprawdę musisz przesyłać pliki tej wielkości, koduj je base64 i podnieś limit do ~70 MiB.

Iloczyn z batchSize decyduje o zużyciu pamięci

Limit dotyczy pojedynczego zdarzenia, ale jeden XREADGROUP pobiera ich batchSize. Przy domyślnych 64 i limicie 50 MiB jeden odczyt może zażądać 3,2 GiB — więcej niż domyślna sterta Node.js.

Biblioteka wykrywa taką parę i ostrzega przy subskrypcji, podając bezpieczną wartość:

WARN  Jeden odczyt może pochłonąć bardzo dużo pamięci - rozważ mniejszy batchSize
      { topic, batchSize: 64, worstCaseBytes: 3355443200, suggestedBatchSize: 10 }

Ostrzeżenie, nie błąd — limit jest górnym ograniczeniem, a nie rozmiarem typowego zdarzenia, więc przerwanie startu subskrypcji przenoszącej małe komunikaty byłoby szkodliwe. Dla tematów z załącznikami ustaw batchSize jawnie:

await bus.subscribe(DocumentUploadedEvent, handler, {
  group: 'archiwum',
  batchSize: 4, // duże załączniki - mały odczyt
});

Rozważ wzorzec claim-check

Przy plikach rzędu dziesiątek megabajtów tańsza jest zwykle inna konstrukcja: plik trafia do MinIO/S3, a zdarzenie przenosi wyłącznie referencję i sumę kontrolną. Zdarzenie schodzi wtedy do ~1 KB, a przy retencji 7 dni różnica w zużyciu pamięci brokera idzie w dziesiątki gigabajtów — Dragonfly trzyma strumienie w RAM, więc każdy załącznik zajmuje ją przez cały okres retencji.


15. Testowanie

15.1. Testowanie handlerów bez brokera

Transport pamięciowy implementuje ten sam port co Redis Streams:

import { createInMemoryEventBus } from 'pp-event-bus';

const { bus, transport } = await createInMemoryEventBus();

await bus.subscribe(OrderCreatedEvent, handler, { group: 'rozliczenia' });
await bus.emit(new OrderCreatedEvent(payload));
await transport.waitForIdle();

expect(handler).toHaveBeenCalledTimes(1);
expect(transport.getDeadLetters()).toHaveLength(0);
expect(transport.getPendingCount('OrderCreatedEvent', 'rozliczenia')).toBe(0);

15.2. Deterministyczne identyfikatory i czas

Warstwa domenowa pobiera identyfikator zdarzenia i znacznik czasu przez porty, konfigurowane w punkcie kompozycji. Domyślne implementacje są ustawiane przy imporcie biblioteki, więc new OrderCreatedEvent(payload) działa bez żadnych dodatkowych kroków — ale w testach można je podmienić:

import { configureEventContext } from 'pp-event-bus';

let counter = 0;

beforeEach(() => {
  configureEventContext({
    idGenerator: { generate: () => `zdarzenie-${(counter += 1)}` },
    clock: { now: () => 1_700_000_000_000 },
  });
});

// new OrderCreatedEvent(payload).id → 'zdarzenie-1'

Ułatwia to asercje na identyfikatorach i testowanie deduplikacji bez zgadywania UUID-ów.


16. Rozwój

16.1. NPM scripts

| Komenda | Opis | | --------------------------------- | --------------------------------------------------------------- | | npm run lint / lint:fix | ESLint flat config (+ reguły granic DDD per warstwa) | | npm run typecheck | tsc --noEmit obejmujący również testy | | npm run format / format:check | Prettier (src/**, README, MIGRATION) | | npm test / test:unit | Testy jednostkowe (bez brokera) | | npm run test:integration | Testy na realnym silniku — wymaga uruchomionego DragonflyDB | | npm run test:coverage | Oba zestawy + progi 80% (bramka CI) | | npm run build | tsc --build do dist/ | | npm run semantic-release | Wydanie (uruchamiane przez CI) |

16.2. Uruchomienie lokalne

npm install
docker compose up -d dragonfly          # broker dla testów integracyjnych

npm run lint && npm run typecheck && npm run format:check
npm run test:coverage
npm run build

Konsola brokera i diagnostyka:

docker compose exec dragonfly redis-cli
docker compose exec dragonfly redis-cli INFO persistence
docker compose down                     # dane zostają w wolumenie
docker compose down -v                  # skasowanie danych

Weryfikacja przeciwko Redisowi zamiast Dragonfly — sprawdza, czy adapter nie polega na zachowaniach jednego silnika:

docker compose --profile redis up -d redis
EVENT_BUS_TEST_REDIS_URL=redis://localhost:6380 npm run test:integration

16.3. Stack technologiczny

| Warstwa | Pakiet | Wersja | Uwagi | | ----------------- | ------------------ | -------------- | -------------------------------------------- | | Runtime | Node.js | 22 LTS | Wymagane ^22.22.2 \|\| >=24.15 dla wydania | | Język | typescript | ^5.9.3 | strict + noUncheckedIndexedAccess | | Transport | ioredis | ^5.11.1 | Jedyna zależność produkcyjna | | Framework (opcj.) | @nestjs/common | ^10 \|\| ^11 | peerDependency, optional: true | | Testy | jest + ts-jest | ^30 / ^29 | Dwa projekty: unit, integration | | Lint | eslint | ^10.8.0 | Flat config + @typescript-eslint 8 | | Format | prettier | ^3.9.6 | 100 kolumn, single quote | | Wydanie | semantic-release | ^25.0.8 | Conventional Commits → GitLab + npm |

16.4. Konwencje kodu

  • Trzy warstwy DDD, zależności wyłącznie do wewnątrz — egzekwowane przez no-restricted-imports.
  • Zero any i zero eslint-disable w kodzie produkcyjnym.
  • JSDoc po polsku dla wszystkich metod publicznych; komentarze wyjaśniają dlaczego, nie co.
  • Limit rozmiaru pliku: zalecane 300 linii, twardy limit 500.
  • Testy nazywane should [zachowanie], dane testowe polskie (Jan Kowalski, ZAM-2024-001).
  • Commity: Conventional Commits, angielski prefiks + polski opis.
  • Pre-commit hook (Husky) uruchamia lint → typecheck → testy --bail i blokuje commit przy niepowodzeniu (~6,7 s).
  • Commit-msg hook waliduje format przez commitlint. To nie kosmetyka: .releaserc.json mapuje typ commita na poziom wydania, więc nierozpoznany typ oznaczałby brak wydania — po cichu.

Transpilacja. Testy chodzą przez @swc/jest (konfiguracja w .swcrc, wskazywana jawnie w jest.config.ts), co skróciło przebieg z 5,7 s do 2,5 s. Bezpieczeństwo typów zapewnia osobny npm run typecheck, zgodnie z firmowym standardem „SWC do transpilacji, tsc --noEmit do kontroli typów".

Build celowo pozostaje na tsc. Wzorcowy auth-service używa SWC przez nest build, ale to aplikacja, która nie publikuje typów. Tutaj dist/ musi zawierać .d.ts, a ich wygenerowanie i tak wymaga pełnej analizy typów przez tsc - zmierzony wariant „SWC + tsc --emitDeclarationOnly" okazał się wolniejszy (2,57 s wobec 2,39 s) przy dwóch krokach zamiast jednego.


17. Ograniczenia

Świadome granice bieżącej wersji — wymienione wprost, żeby nie trzeba było ich odkrywać na produkcji:

  • Brak API odtwarzania z DLQ. Odrzucone zdarzenia można obejrzeć (XRANGE) i opublikować ponownie własnym skryptem, ale biblioteka nie dostarcza do tego gotowej metody.
  • Jedno połączenie na subskrypcję. Blokujący odczyt zajmuje połączenie na wyłączność. Serwis subskrybujący 20 typów zdarzeń utrzyma 21 połączeń. Multipleksowanie wielu strumieni w jednym odczycie jest możliwe, ale nie zostało jeszcze zrobione.
  • Redis Cluster nieprzetestowany. Fabryka połączeń przyjmuje odpowiednie opcje ioredis, jednak testy integracyjne działają wyłącznie na pojedynczej instancji. Przy klastrze trzeba zadbać, by strumienie danej grupy trafiały do tego samego slotu.
  • Ponowienia mieszczą się w oknie przejęcia. Cykl ponowień musi zamknąć się poniżej EVENT_BUS_PENDING_IDLE_TIME_MS. Biblioteka to waliduje, ale oznacza to, że długie backoffy wymagają jednoczesnego podniesienia progu przejęcia.
  • Poison message blokuje swoją partycję. Przy partitionKey wiadomość wyczerpująca budżet ponowień wstrzymuje kolejne zdarzenia o tym samym kluczu na czas całego cyklu (domyślnie do 15 s). Jest to konieczne dla zachowania kolejności, ale wpływa na dobór maxAttempts w systemach wrażliwych na opóźnienie. Zdarzenia o innych kluczach biegną bez przeszkód.
  • Dwie subskrypcje tego samego zdarzenia w tej samej grupie są odrzucane. Dzieliłyby tożsamość konsumenta, co groziłoby podwójnym przetworzeniem — patrz 6.3.
  • Potwierdzenia zbuforowane giną przy twardym ubiciu procesu. Bufor jest opróżniany przy gracefulnym zamknięciu, ale SIGKILL zostawi do batchSize wiadomości niepotwierdzonych. Wrócą jako duplikaty i usunie je deduplikacja — świadomy kompromis na rzecz kilkudziesięciokrotnie mniejszej liczby poleceń.
  • Odczyt formatu 1.x jest tymczasowy. Dekoder akceptuje koperty z prefiksami __ na czas migracji kroczącej i zostanie usunięty po przejściu wszystkich producentów.

18. Dokumentacja

| Dokument | Zawartość | | ------------------------------------------ | ---------------------------------------------------------------- | | CHANGELOG.md | Historia zmian (generowana przez semantic-release) | | MIGRATION.md | Migracja z wersji 1.x wraz z listą kontrolną | | .env.example | Wszystkie zmienne środowiskowe z opisami | | docker-compose.yml | DragonflyDB z konfiguracją trwałości, Redis pod profilem redis |

Testy integracyjne uruchamiają ten sam zestaw kontraktowy co testy transportu pamięciowego (src/infrastructure/__tests__/transport-contract.suite.ts) — nowa implementacja transportu powinna go przechodzić bez modyfikacji.

18.1. Kontakt

Zespół: [email protected] Repozytorium: gitlab.devbooster.pl/lib/event-bus

Licencja

MIT