npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

pg-cdc

v0.1.0

Published

Postgres change data capture (CDC) as Effect streams, built on logical replication and pgoutput

Readme

pg-cdc

Postgres change data capture (CDC) as Effect streams.

pg-cdc connects to Postgres over the logical replication protocol, decodes pgoutput messages, and exposes committed transactions as a typed Stream. Each transaction carries an acknowledge effect, so you control exactly when the replication slot advances.

Status: early release. The API may change before 1.0. effect is currently a 4.0 release candidate.

Install

pnpm add pg-cdc effect pg

effect and pg are peer dependencies.

Postgres setup

Logical replication has to be enabled on the server, and you need a publication for the tables you want to capture.

-- postgresql.conf (requires a restart)
-- wal_level = logical

CREATE PUBLICATION my_pub FOR TABLE todos;

-- Optional: include full old-row values on UPDATE / DELETE.
-- Without this, `before` only contains the primary key.
ALTER TABLE todos REPLICA IDENTITY FULL;

The connecting role needs the REPLICATION attribute (or be a superuser). The replication slot is created on first run if it does not exist.

examples/init.sql contains a complete schema (table, publication, and seed rows) that works with examples/demo.ts.

Usage

With make

make returns a scoped Effect. The replication connection is closed when the scope ends.

import { Effect, Stream } from "effect"
import { NodeRuntime } from "@effect/platform-node"
import * as PostgresCDC from "pg-cdc"

const program = Effect.gen(function* () {
  const cdc = yield* PostgresCDC.make({
    connectionString: "postgres://postgres:postgres@localhost:5432/postgres",
    publication: "my_pub",
    slot: "app_slot",
  })

  yield* cdc.transaction.pipe(
    Stream.runForEach((tx) =>
      Effect.gen(function* () {
        yield* publish(tx.changes) // your side effect
        yield* tx.acknowledge      // then advance the slot
      })
    )
  )
})

program.pipe(Effect.scoped, NodeRuntime.runMain)

With a Layer

layer(config) provides the PostgresCDC service. The connection lives as long as the layer.

import { Effect, Stream } from "effect"
import { PostgresCDC, layer } from "pg-cdc"

const program = Effect.gen(function* () {
  const cdc = yield* PostgresCDC
  yield* cdc.transaction.pipe(Stream.runForEach(handleTransaction))
})

program.pipe(
  Effect.provide(layer({ connectionString: "...", publication: "my_pub", slot: "app_slot" }))
)

Because PostgresCDC is a service, tests can swap in a fake:

Layer.succeed(PostgresCDC, {
  transaction: Stream.make(fakeTransaction),
  changes: Stream.make(fakeChange),
})

API

make(config): Effect<PostgresCDCService, PostgresConnectionError | PgReplError, Scope>

| option | description | | ------------------ | -------------------------------------------------------- | | connectionString | Postgres connection string. The role needs REPLICATION. | | publication | Name of an existing publication. | | slot | Replication slot name. Created if it does not exist. |

PostgresCDCService

interface PostgresCDCService {
  transaction: Stream<CDCTransaction, PgReplError | RelationNotFound>
  changes: Stream<CDCChange, PgReplError | RelationNotFound>
}

transaction — one element per committed transaction. Nothing is acknowledged until you run tx.acknowledge. Use this when you need at-least-once delivery into another system (Kafka, another database, ...).

changes — the same data flattened to individual row changes. Acknowledgement happens automatically once all changes of a transaction have been pulled through your consumer. Keep your side effect directly in Stream.runForEach (or an equivalent sequential sink); inserting Stream.buffer, Stream.mapEffect with concurrency, or a fork between the stream and your side effect lets the acknowledgement run before your work is done.

Only one of transaction / changes may be consumed per make (or per layer). Both share a single replication connection and slot, and Postgres allows one active consumer per slot.

CDCTransaction

interface CDCTransaction {
  xid: number                          // Postgres transaction id
  commitLSN: string                    // e.g. "0/19B58D8"
  changes: CDCChange[]
  acknowledge: Effect<void, PgReplError>
}

CDCChange

A Data.TaggedEnum with three cases. before is only present when Postgres sent old-row data (see REPLICA IDENTITY above).

type CDCChange =
  | { _tag: "Insert"; schema: string; table: string; after: Record<string, unknown> }
  | { _tag: "Update"; schema: string; table: string; before?: Before; after: Record<string, unknown> }
  | { _tag: "Delete"; schema: string; table: string; before?: Before }

type Before = { kind: "full" | "key"; rows: Record<string, unknown> }

Match on it with CDCChange.$match or Match.tag.

Errors

| error | when | | ------------------------- | --------------------------------------------------------------------- | | PostgresConnectionError | the replication connection could not be opened (cause has details) | | RelationNotFound | a row change referenced a table no Relation message was seen for | | PgReplError | protocol / decoding error from pg-replicator |

Delivery semantics

  • At-least-once. If the process dies before acknowledge, Postgres redelivers the transaction on the next connection. Make your consumers idempotent.
  • Unacknowledged transactions retain WAL. Postgres keeps WAL segments until the slot is advanced. If you never acknowledge, disk usage grows without bound.
  • Ordering. Transactions are delivered in commit order.

Development

pnpm install
pnpm typecheck
pnpm build
pnpm example   # runs examples/demo.ts against a local Postgres

License

MIT