Skip to content

Stateful Pipeline Transforms

Stateful Pipeline Transforms

Stateful transforms extend PipelineService with middleware that can enrich, aggregate, deduplicate, and filter events before pipeline emission.

flowchart LR
  E[Pipeline Event] --> C[StatefulPipelineService]
  C --> T1[Deduplication]
  T1 --> T2[Enrichment]
  T2 --> T3[Threat Intel Tagging]
  T3 --> T4[Aggregation]
  T4 --> T5[Rate Limit]
  T5 --> P[PipelineService.send]

Contracts

  • IPipelineTransform<TIn, TOut>
  • TransformContext (kv, logger, now)
  • PipelineTransformChain fluent composition

Built-in Transforms

  • DeduplicationTransform
  • EnrichmentTransform
  • ThreatIntelTaggingTransform
  • AggregationTransform
  • RateLimitTransform

Fail-safe Behavior

  • Transform failures are logged.
  • Default behavior is fail-open (event passes through unchanged).
  • failOpen: false enables fail-closed dropping on transform errors.

Custom Transform Authoring

  1. Implement IPipelineTransform.
  2. Register via .use(...) on StatefulPipelineService.
  3. Keep transform side effects idempotent.
  4. Write focused tests for drop/pass-through behavior.