Ingestion pipeline architecture overview and convention reference. Use when you need a quick orientation to the pipeline framework or want to know which doctor agent to use for a specific concern.
Scanned 9/1/2026
Install to Claude Code
npx -y skills add PostHog/posthog --skill ingestion-pipeline-doctor-nodejs --agent claude-codeInstalls into .claude/skills of the current project.
Are you the author of Ingestion Pipeline Doctor Nodejs?
Add the live security badge to your README — it updates automatically with every re-scan.
[](https://www.skillsdirectory.com/skills/posthog-ingestion-pipeline-doctor-nodejs-posthog)More formats (shields.io, HTML) on the badges page.
---
name: ingestion-pipeline-doctor-nodejs
description: >
Ingestion pipeline architecture overview and convention reference.
Use when you need a quick orientation to the pipeline framework
or want to know which doctor agent to use for a specific concern.
---
# Pipeline Doctor
Quick reference for PostHog's ingestion pipeline framework and its convention-checking agents.
## Architecture overview
The ingestion pipeline processes events through a typed, composable step chain:
```text
Kafka message
→ messageAware()
→ parse headers/body
→ sequentially() for preprocessing
→ filterMap() to enrich context (e.g., team lookup)
→ teamAware()
→ concurrentlyPerGroup(token:distinctId) for per-entity processing
→ gather()
→ pipeChunk() for chunk operations
→ handleIngestionWarnings()
→ handleResults()
→ handleSideEffects()
→ build()
```
See `nodejs/src/ingestion/pipelines/analytics/joined-ingestion-pipeline.ts` for the real implementation.
## Key file locations
| What | Where |
| ----------------- | ----------------------------------------------------------------------- |
| Step type | `nodejs/src/ingestion/framework/steps.ts` |
| Result types | `nodejs/src/ingestion/framework/results.ts` |
| Doc-test chapters | `nodejs/src/ingestion/framework/docs/*.test.ts` |
| Joined pipeline | `nodejs/src/ingestion/pipelines/analytics/joined-ingestion-pipeline.ts` |
| Doctor agents | `.claude/agents/ingestion/` |
| Test helpers | `nodejs/src/ingestion/framework/docs/helpers.ts` |
## Which agent to use
| Concern | Agent | When to use |
| --------------- | ----------------------------- | --------------------------------------------------------- |
| Step structure | `pipeline-step-doctor` | Factory pattern, type extension, config injection, naming |
| Result handling | `pipeline-result-doctor` | ok/dlq/drop/redirect, side effects, ingestion warnings |
| Composition | `pipeline-composition-doctor` | Builder chain, concurrency, grouping, branching, retries |
| Testing | `pipeline-testing-doctor` | Test helpers, assertions, fake timers, doc-test style |
## Quick convention reference
**Steps**: Factory function returning a named inner function. Generic `<T extends Input>` for type extension. No `any`. Config via closure.
**Results**: Use `ok()`, `dlq()`, `drop()`, `redirect()` constructors. Side effects as promises in `ok(value, [effects])`. Warnings as third parameter.
**Composition**: `messageAware` wraps the pipeline. `handleResults` inside `messageAware`. `handleSideEffects` after. `concurrentlyPerGroup` for per-entity work. `gather` before chunk steps.
**Batching lifecycle hooks** (`BatchingPipeline` beforeBatch/afterBatch): enrich-only. Hooks may enrich elements and batch context but must return exactly the elements they received — a count change is a broken invariant and `feed()` throws. Filtering belongs in sub-pipeline steps that return `drop()`. An empty `feed()` is a no-op (no hooks, no capacity). Details: `nodejs/src/ingestion/framework/docs/14-batching.test.ts`.
**Fan-out/fan-in** (`fanOut(fn).via((sub) => …).fanIn(fn)`): per-element sub-work with cardinality restored — one element fans out to N sub-elements (e.g. per-blob uploads), a regular sub-pipeline processes them (`maxConcurrency` on the sub `concurrently` block, `retry` on the per-sub step), and fan-in folds the OK results back into the parent. Reach for it over `concurrently`/`concurrentlyPerGroup` when the unit of concurrency is smaller than the element; hand-rolled `p-limit`/`Promise.all` inside a step is the tell. Sequencing is compile-time enforced (an unclosed stage cannot build). Sub-result contract: OK collected; DROP excludes the sub silently; DLQ fails the parent with aggregated reasons; REDIRECT is excluded with a warning — sub redirects never escape the stage. Sub-pipelines are context-agnostic: team/message data goes in the sub-element value, and context-gated surface (`teamAware`, `handleIngestionWarnings`, …) is uncallable. Fan-out/fan-in functions are cheap, synchronous, and named. Parents emit unordered as they complete. Details: `nodejs/src/ingestion/framework/docs/17-fan-out-fan-in.test.ts`.
**Testing**: Step tests call factory directly. Use `consumeAll()`/`collectChunks()` helpers. Fake timers for async. Type guards for result assertions. No `any`.
## Running all doctors
Ask Claude to "run all pipeline doctors on my recent changes" to get a comprehensive review across all 4 concern areas.
Is this your skill, or is something wrong with this listing? Request removal or report an issue. Author removals are honored within 72 hours.
No comments yet. Be the first to comment!