pp-event-bus
v2.0.0
Published
Distributed event bus library with pluggable transports (Redis Streams / DragonflyDB)
Maintainers
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
- Główne funkcje
- Architektura
- Instalacja
- Szybki start
- Definiowanie zdarzeń
- Publikacja i subskrypcja
- Wydajność
- Gwarancje dostarczenia
- Struktura projektu
- Zmienne środowiskowe
- Integracja z NestJS
- Obserwowalność i kontrola zdrowia
- Trwałość i konfiguracja DragonflyDB
- Bezpieczeństwo
- Testowanie
- Rozwój
- Ograniczenia
- 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. perorderId). - 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 MINIDzamiastMAXLEN; 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-busZależ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 biznesowe5.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 koleiBez 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
consumerLagistnieje 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 10Wyłą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_BUS10. 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_URLjest 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ę wcreateRedisEventBus/EventBusModulei zgłasza wszystkie błędne zmienne naraz. WcześniejEVENT_BUS_RETENTION_MAX_AGE_MS=7dparsował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_MSpomniejszonym 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
loggerw konfiguracji biblioteka nie zakłada własnego nasłuchu i awarie wypisuje ioredis:[ioredis] Unhandled error event: …przezconsole.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 podanieloggerjest 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 yesnie 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: 513.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_threadsmniej więcej na liczbę dostępnych rdzeni. - Potwierdzenia trafiają w jeden shard.
XACKz 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: 60spó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 orazconsumerLag.
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,apiKeyi 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 buildKonsola 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 danychWeryfikacja 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:integration16.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
anyi zeroeslint-disablew 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 --baili blokuje commit przy niepowodzeniu (~6,7 s). - Commit-msg hook waliduje format przez commitlint. To nie kosmetyka:
.releaserc.jsonmapuje 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
partitionKeywiadomość 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órmaxAttemptsw 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
SIGKILLzostawi dobatchSizewiadomoś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
