Skip to content
Back to skills

Ingestion Pipeline Doctor Nodejs

ASecurity

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.

  • 40,048 stars
  • 0 votes
  • 0 copies
  • 2 views
  • Added September 1, 2026
developmentgonodenodejstesting

Security analysis

A100/100

Scanned September 1, 2026

npx -y skills add PostHog/posthog --skill ingestion-pipeline-doctor-nodejs --agent claude-code

Installs 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 with every re-scan.

Security grade badge for Ingestion Pipeline Doctor Nodejs
[![Security: A — Skills Directory](https://www.skillsdirectory.com/api/skills/posthog-ingestion-pipeline-doctor-nodejs-posthog/badge)](https://www.skillsdirectory.com/skills/posthog-ingestion-pipeline-doctor-nodejs-posthog)

More formats (shields.io, HTML) on the badges page. Keep it an A: scan every change in CI with Pro.

Download with Pro
SKILL.md
---
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.

Attribution

Is this your skill, or is something wrong with this listing? Request removal or report an issue. Author removals are honored within 72 hours.

Comments

Loading comments…