posthog/posthog-foss

ingestion-pipeline-doctor-nodejs

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.

First seen Jun 24, 2026

Installation

$ npx skills add posthog/posthog-foss --skill ingestion-pipeline-doctor-nodejs

Similar popular skills

Related neighbors and high-traction skills in the same topics — useful to compare before installing.

Also in this package

Other skills from posthog/posthog-foss · top by installs.

npx skills add posthog/posthog-foss

Browse all from posthog/posthog-foss

More details

Agent compatibility

Declared targets from SKILL.md / docs. Unmarked agents are not listed — the skill may still install via the CLI.

Claude Code Declared
Cursor Not declared
Codex Not declared
GitHub Copilot Not declared
Windsurf Not declared
Gemini CLI Not declared
Cline Not declared
OpenCode Not declared

Repository health

Stars 516
License LICENSE
Default branch master
Status Active

Skill metadata

Parsed from SKILL.md frontmatter.

Declared agents claude-code

Package contents

Files included with this skill beyond the listing page.

  • skill md SKILL.md 4,902 B
  • docs SUMMARY.md 236 B

History

  1. First seen on skills.sh
  2. First recorded snapshot · 5 installs

SKILL.md

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:

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.