@cqrs-ddd/pipeline-resilience
Πολιτικές ανοχής σφαλμάτων και αξιοπιστίας βασισμένες στο cockatiel, εφαρμοζόμενες σε δύο διακριτά αρχιτεκτονικά επίπεδα:
- Ανθεκτικότητα σε επίπεδο λειτουργίας (
ResilienceBehavior): Περιβάλλει ολόκληρους handlers ή μεμονωμένες λειτουργίες με απομόνωση retry, timeout και bulkhead. - Πολιτικές εξωτερικών εξαρτήσεων (
ResiliencePolicies): Ένα κεντρικό μητρώο επώνυμων πολιτικών ανθεκτικότητας (circuit breakers, retry, bulkhead, timeout, fallback) για κλήσεις RPC, payment APIs, database queries και εξωτερικά webhooks.
Επιβάλλει αυστηρούς ελέγχους κατά το startup προς αποφυγή επικίνδυνων επανεκτελέσεων εντολών με παρενέργειες.
Εγκατάσταση
Ενότητα με τίτλο «Εγκατάσταση»pnpm add @cqrs-ddd/pipeline-resilience @cqrs-ddd/pipelineΑπαιτεί Node.js 22.12 ή νεότερο. Περιλαμβάνει το cockatiel 4 ως εξάρτηση.
Δύο αρχιτεκτονικά επίπεδα
Ενότητα με τίτλο «Δύο αρχιτεκτονικά επίπεδα»Εισερχόμενο Command / Query │ ▼[ResilienceBehavior] (Επίπεδο Λειτουργίας) ├─ Bulkhead: Περιορίζει τις ταυτόχρονες εκτελέσεις στον συγκεκριμένο handler ├─ Timeout: Θέτει ανώτατο όριο στον συνολικό χρόνο εκτέλεσης του handler └─ Retry: Επανεκτελεί τη λειτουργία (σε commands και events μόνο με replaySafe) │ ▼[Επιχειρησιακή Λογική Handler] │ ▼ Εξερχόμενη Απομακρυσμένη Κλήση (π.χ. Stripe, AWS S3, SendGrid)[ResiliencePolicies] (Επίπεδο Εξάρτησης) ├─ Circuit Breaker: Ανοίγει σε διαδοχικές αποτυχίες του εξωτερικού συστήματος ├─ Outbound Timeout: Περιορίζει τη διάρκεια του HTTP socket ├─ Outbound Retry: Επαναλαμβάνει παροδικές πτώσεις σύνδεσης └─ Fallback: Επιστρέφει cached ή υποβαθμισμένη απάντηση σε διακοπήΧρήση σε επίπεδο λειτουργίας (ResilienceBehavior)
Ενότητα με τίτλο «Χρήση σε επίπεδο λειτουργίας (ResilienceBehavior)»import { createPipeline } from '@cqrs-ddd/pipeline';import { resilience, getResilienceAbortSignal } from '@cqrs-ddd/pipeline-resilience';
const pipeline = createPipeline();
export const fetchRemotePrice = pipeline.wrap( { name: 'fetchRemotePrice', kind: 'query' }, resilience({ retry: { maxAttempts: 3, backoff: { type: 'exponential', initialDelay: 200, maxDelay: 2_000 }, }, timeout: { duration: 3_000, strategy: 'cooperative' }, bulkhead: { limit: 20, queue: 10 }, handle: (error) => error instanceof UpstreamNetworkError, }),)(async (sku: string) => { const signal = getResilienceAbortSignal(); return pricingApi.getPrice(sku, { signal });});Invariants ασφαλείας & Έλεγχοι
Ενότητα με τίτλο «Invariants ασφαλείας & Έλεγχοι»Επειδή τα retries επανεκτελούν ολόκληρο το περιεχόμενο του behavior, το ResilienceBehavior επικυρώνει αυστηρούς κανόνες:
- Επιλεκτικό φίλτρο σφαλμάτων (
handle): Απαιτείται ρητή συνάρτησηhandle: (error) => booleanπου επιλέγει παροδικά σφάλματα (πτώσεις δικτύου, deadlocks). Για επανάληψη σε όλα τα σφάλματα, απαιτείταιhandleAllErrors: true. - Ασφάλεια επανάληψης σε Commands/Events (
retry.replaySafe): Commands και events μεταβάλλουν κατάσταση. Η ρύθμιση retry σε command χωρίς ρητόretry: { replaySafe: true }αποτυγχάνει κατά την εκκίνηση μεPipelineConfigurationError. - Στρατηγική Timeout (
timeout.strategy):'aggressive'(προεπιλογή): Απαντά άμεσα στον καλούντα με σφάλμα timeout ενώ ο handler συνεχίζει να εκτελείται στο background. Σε commands ή events απαιτείtimeout: { replaySafe: true }.'cooperative': Ακυρώνει τοAbortSignalαπό τοgetResilienceAbortSignal()και περιμένει τον handler να ολοκληρώσει πριν απαντήσει.
- Όχι Circuit Breakers σε Use Cases: Η δήλωση circuit breaker απευθείας σε έναν handler απορρίπτεται. Τα circuit breakers παρακολουθούν την υγεία εξωτερικών εξαρτήσεων και όχι τα εσωτερικά use cases. Δηλώνονται στο
ResiliencePolicies.
Πολιτικές εξωτερικών εξαρτήσεων (ResiliencePolicies)
Ενότητα με τίτλο «Πολιτικές εξωτερικών εξαρτήσεων (ResiliencePolicies)»Κεντρικό μητρώο επώνυμων πολιτικών για εξωτερικά συστήματα:
import { ResiliencePolicies } from '@cqrs-ddd/pipeline-resilience';
export const externalPolicies = new ResiliencePolicies({ paymentGateway: { handle: (error) => isTransientHttpError(error), retry: { maxAttempts: 2, backoff: { type: 'exponential', initialDelay: 500 } }, circuitBreaker: { halfOpenAfter: 30_000, breaker: { type: 'consecutive', threshold: 5 }, }, timeout: { duration: 5_000 }, }, smsProvider: { timeout: { duration: 2_000 }, bulkhead: { limit: 10, queue: 5 }, },});
const charge = await externalPolicies.execute('paymentGateway', async ({ signal }) => { return stripe.charges.create(chargeParams, { signal });});Όλοι οι handlers μοιράζονται την ίδια κατάσταση circuit breaker για το 'paymentGateway', προστατεύοντας τα εξωτερικά APIs από αλυσιδωτές αποτυχίες (cascading failures).
Ενσωμάτωση NestJS (@cqrs-ddd/nestjs)
Ενότητα με τίτλο «Ενσωμάτωση NestJS (@cqrs-ddd/nestjs)»Καταχωρίστε το ResiliencePolicies ως injectable provider:
import { Module, Global } from '@nestjs/common';import { ResiliencePolicies } from '@cqrs-ddd/pipeline-resilience';
@Global()@Module({ providers: [ { provide: ResiliencePolicies, useFactory: () => new ResiliencePolicies({ stripeApi: { handle: (error) => isTransientHttpError(error), retry: { maxAttempts: 3 }, circuitBreaker: { halfOpenAfter: 15_000, breaker: { type: 'consecutive', threshold: 3 } }, timeout: { duration: 4_000 }, }, }), }, ], exports: [ResiliencePolicies],})export class ResilienceModule {}Διακοσμήστε handlers με resilience():
import { CommandHandler, ICommandHandler } from '@nestjs/cqrs';import { UsePipeline } from '@cqrs-ddd/pipeline';import { resilience } from '@cqrs-ddd/pipeline-resilience';
@CommandHandler(SyncInventoryCommand)@UsePipeline( resilience({ retry: { maxAttempts: 3, replaySafe: true }, timeout: { duration: 10_000, strategy: 'cooperative' }, handle: (err) => err instanceof DatabaseDeadlockError, }),)export class SyncInventoryHandler implements ICommandHandler<SyncInventoryCommand> { async execute(command: SyncInventoryCommand) {}}Σειρά των επιπέδων (Layer Ordering)
Ενότητα με τίτλο «Σειρά των επιπέδων (Layer Ordering)»Το ResilienceBehavior συνθέτει τις πολιτικές με την ακόλουθη προεπιλεγμένη σειρά (από το εξωτερικό προς το εσωτερικό):
['retry', 'bulkhead', 'timeout']Μπορείτε να προσαρμόσετε τη σειρά μέσω της επιλογής order:
resilience({ order: ['bulkhead', 'retry', 'timeout'], retry: { maxAttempts: 3 }, bulkhead: { limit: 10 }, timeout: { duration: 2_000 },})Τύποι σφαλμάτων & Guards
Ενότητα με τίτλο «Τύποι σφαλμάτων & Guards»Επανεξάγει τις κλάσεις σφαλμάτων του Cockatiel και αντίστοιχα guards:
BrokenCircuitError/isBrokenCircuitError(err): Ρίχνεται όταν μια κλήση μπλοκάρεται από ανοιχτό circuit breaker.BulkheadRejectedError/isBulkheadRejectedError(err): Ρίχνεται όταν τα όρια ταυτόχρονης εκτέλεσης και η ουρά εξαντληθούν.TaskCancelledError/isTaskCancelledError(err): Ρίχνεται όταν η εκτέλεση υπερβεί το timeout.