NewYour coding agent can read the release notes before it upgrades.Set up the MCP server →
Go modules · #1654 by repository stars
Last release 13 days ago
24 Sep 2026
Ships fairly regularly
a new release about every 2 weeks
Rarely documented
notes for 11 of the last 60 stable releases
Nothing withdrawn
no release was ever pulled
2 years old
128 releases · first in 2024
Nothing published for this version
Nothing published for this version
Nothing published for this version
One column per month.
Nothing published for this version
🧮 Transformation rules work on array columns. Now each element is transformed, with an optional array_options to control the output 🔐 New fpe_ff1 tran
⚡ Highlights
🧮 Transformation rules work on array columns. Now each element is transformed, with an optional array_options to control the output
🔐 New fpe_ff1 transformer: format-preserving encryption for phone numbers, emails and codes, uniqueness preserved
🔎 New lookup_choice transformer: replaces a value with one read from a live table column
🔑 Kafka SASL authentication (plain, scram-sha-256, scram-sha-512) on source and target
⚡ Parallel index restore in pg-to-pg schema snapshots via index_restore_workers, and transient pg_restore failures are retried instead of aborting the snapshot
🧩 Transformers accept more column types: numeric in the greenmask number transformers, user-defined enums in greenmask_choice, any type in template
fpe_ff1 format-preserving transformer for emails, phone numbers by @kvch in #1165lookup_choice transformer by @kvch in #1148A rule on an array column such as varchar(500)[] was rejected at validation time with does not support pg data type: _varchar with OID: 1015. Because the rejection happened for the whole rule set, a single array column blocked transformation of every other column in the table. A rule on a one-dimensional array column now applies the transformer to each element.
transformations:
table_transformers:
- schema: public
table: users
column_transformers:
emails: # varchar(500)[]
name: email
parameters:
replacement_domain: "@example.com.invalid"pgstream decodes the array before it calls the transformer, so each element arrives as the value of its element type: an int4[] element is an integer, not the text "1", and a NULL element is distinguishable from the string "NULL". A transformer accepts an array column if it accepts the element type, which extends to extension types (citext[]) and to arrays of enums. The array is re-encoded in the same shape it arrived in, so Kafka, webhook and search payloads for array columns are unchanged.
The optional array_options block controls how the output array is built:
column_transformers:
emails:
name: email
array_options:
generator: random # map (default) or random
min_count: 0
max_count: 12map applies the transformer to each element in order and keeps the length. random emits between min_count and max_count elements, each the transform of a source element drawn at random with replacement, so it never invents a value that was not in the row. min_count and max_count are required with random, rejected with map, and max_count is capped at 10000. Array columns need a source PostgreSQL connection, since the element type is read from the catalog; the connection-less parser rejects array_options.
The uniqueness check reads the classification of the whole array transform: with map the array keeps the transformer's classification, so fpe_ff1 on text[] stays preserved; with random the array is always lossy, since min_count: 0 lets any row produce an empty array. A NULL array stays NULL, a NULL element stays NULL, and the first failing element fails the whole column under the configured on_error policy. Multi-dimensional arrays are not supported and are rejected at startup when declared as such.
Three things the validator now flags: literal_string and pg_anonymizer write into every element on an array column, which is a change for a configuration written before array support; pg_anonymizer sends one query per element, so a wide array multiplies the source load; and a column declared with more than one dimension is rejected up front instead of failing row by row. The first two are warnings in pgstream validate rules and in the run log.
fpe_ff1: format-preserving encryptionThe output has the same length and the same characters as the input. An encrypted phone number is still a phone number, an encrypted product code still fits its varchar(n) column, and FF1 (NIST SP 800-38G) maps each input to one different output, so the transformer is classified preserved and can be used on a unique-indexed column. It is the first uniqueness-preserving transformer whose output keeps the original format.
transformations:
table_transformers:
- schema: public
table: persons
column_transformers:
phone_number:
name: fpe_ff1
parameters:
key_hex: ${FPE_KEY} # 32, 48 or 64 hex chars; openssl rand -hex 32
associated_data: public.persons.phone_number # FF1 tweak, not secret
alphabet: digits # digits, letters, alphanumeric, or a set of characters
passthrough: keep # copy characters outside the alphabet to the output
keep_prefix: 4 # keep the country and operator code
email:
name: fpe_ff1
parameters:
key_hex: ${FPE_KEY}
associated_data: public.persons.email
alphabet: alphanumeric
passthrough: keep
preserve_from: "." # keep the top-level domainkeep_prefix and keep_suffix count only characters in the alphabet, so one rule serves every phone format: +36301234567, +36 30 123 4567 and +36-30-123-4567 all keep their first four digits and their separators. preserve_from keeps the end of the value from the last occurrence of a delimiter: "." keeps the top-level domain so the address stays well-formed, "@" keeps the full domain but then the target contains every source domain. FF1 refuses domains smaller than 106, so a value needs at least 6 digits, or 4 letters or alphanumerics, after prefix, suffix and passthrough characters are removed; shorter values go to the on_error policy rather than reaching the target in clear.
Like encrypted_aes_siv, this is pseudonymization: anyone holding key_hex can decrypt, and equal inputs give equal outputs. Keep the key in a secret store, use a different key per environment, and set associated_data per column so tokens cannot be correlated across columns.
lookup_choice: choose from a live tablegreenmask_choice and greenmask_integer force users to regenerate their rules file whenever a lookup table changes. lookup_choice replaces a value with one taken from a column of another table, read from the source database once at pipeline start:
transformations:
table_transformers:
- schema: public
table: addresses
column_transformers:
country_id:
name: lookup_choice
parameters:
lookup_table: public.countries # schema-qualified, quoted as written
lookup_column: id
generator: deterministic # or random
ignore_values: [0, -1] # drop placeholder rows from the pool
# max_values: 100000 # load fails above this, rather than truncating
# postgres_url: ... # only needed when the source is not PostgreSQLThe type is taken from the lookup column, and a rule that points a column at an incompatible lookup type is rejected at startup; a narrower integer or float is accepted for a wider column. deterministic picks by hash of the incoming value so rows that shared an original still share a replacement, random picks per row. Deterministic mode is rejected for timestamp and timestamptz lookup columns, because snapshot and replication deliver them in forms that do not hash the same.
Please note that the transformer is lossy: any table with more rows than the lookup column has values produces duplicates, so it cannot be used on a unique-indexed column. The values come from the source lookup table, so if that table's key column is itself transformed, or a lookup row is deleted after start, the chosen value does not exist on the target and the foreign key rejects it; under replication that row is dropped with a DATALOSS line unless strict_mode is set. And the whole column is loaded into memory once per rule, capped at max_values, with a 30 second read timeout.
Both the Kafka source and the Kafka target now authenticate with SASL plain, scram-sha-256 or scram-sha-512. The mechanism is required when SASL is enabled rather than defaulting to plain, so a configuration that omits it fails. GSSAPI and OAUTHBEARER are not supported.
source:
kafka:
servers: ["localhost:9092"]
topic:
name: "mytopic"
consumer_group:
id: "mygroup"
tls:
ca_cert: "/path/to/ca.crt"
sasl: # remove this section to disable SASL authentication
mechanism: "scram-sha-512" # plain, scram-sha-256 or scram-sha-512. Required
user: "myuser"
password: "mypassword"The same block goes under target.kafka. The environment variables are PGSTREAM_KAFKA_SASL_ENABLED, PGSTREAM_KAFKA_SASL_MECHANISM, PGSTREAM_KAFKA_SASL_USER and PGSTREAM_KAFKA_SASL_PASSWORD. plain sends the password in clear text, so enable TLS with it. Thanks to @RemiDesgrange for this contribution.
Index restore during a pg-to-pg schema snapshot ran as a single sequential psql session, so a schema with many indexes built them one at a time even though CREATE INDEX statements do not depend on each other. The new index_restore_workers setting restores standalone CREATE INDEX/CREATE UNIQUE INDEX statements concurrently. Everything that can depend on an index existing (constraints added USING INDEX, comments, partition attachments, REPLICA IDENTITY, CLUSTER) is restored afterwards, once every index has been created, in the original order.
source:
postgres:
snapshot:
schema:
mode: pgdump_pgrestore
pgdump_pgrestore:
index_constraint_session_settings:
- maintenance_work_mem=4GB
index_restore_workers: 8 # defaults to 1 (sequential), maximum 32Each worker holds one target connection and applies index_constraint_session_settings to it, so the target needs that many free connection slots and up to workers × maintenance_work_mem of memory. A failing index does not cancel the restores still in flight; the per-index failures are merged and the retry decision is taken over the whole wave. Measured locally: 300 small indexes went from 4.8s to 1.7s with 8 workers, and 8 indexes on a 3M-row table from 34.4s to 11.0s with 4. The environment variable is PGSTREAM_POSTGRES_SNAPSHOT_INDEX_RESTORE_WORKERS.
Any pg_restore failure during the index and constraint phase that was not recognised as already-exists, does-not-exist, permission or constraint became critical and aborted the whole snapshot. A dropped connection, a server still coming up or a lock that could not be taken in time therefore lost a snapshot that a second attempt would have completed. Restore output is now classified against a list of transient fragments (dropped or refused connections, SSL/EOF/broken pipe, server starting up or in recovery, connection-slot exhaustion, lock timeout, deadlock, serialization failure) and the index and constraint restore is retried with exponential backoff: 1s initial, 1min max, 3 retries.
Three additions:
numeric columns are accepted by greenmask_float, greenmask_integer and greenmask_unix_timestamp, which previously failed with does not support pg data type: numeric with OID: 1700. For greenmask_float on a numeric column min_value and max_value are required, since the default range covers all float32 values and would write the same constant into every row. The range is checked against the column's numeric(p,s) at validation time. The source value is converted to float64 for the seed, so digits beyond float64 precision are rounded; this affects only the seed, since the transformer discards the source value.
User-defined enum columns are accepted by greenmask_choice. choices becomes optional on an enum column: pgstream reads the labels from the catalog at startup and logs them, and an explicit choices list is checked against the labels so a typo stops the run rather than every insert. A domain over an enum, and an array of an enum, resolve to the enum. This needs a source PostgreSQL URL; a Kafka-sourced pipeline still has to set choices. Prefer generator: random on an enum: deterministic maps each label to one fixed label with no secret key, and pgstream warns when it sees that combination.
column_transformers:
mood: # a user-defined enum column
name: greenmask_choice # no choices needed; the enum labels are usedThe template transformer declared only string and byte_array as compatible, so pgstream validate rejected it on an int4 column although the documentation said it works on any type. It now declares all, like literal_string and pg_anonymizer. That exposed a second gap: a snapshot row handed the template pgx Go types where replication handed it Postgres text, so {{ .GetValue }} rendered a date as 2024-02-29 00:00:00 +0000 UTC on a snapshot and the target rejected it. Snapshot values are now encoded with pgx's text codecs before the template runs, so a numeric renders as decimal text, a uuid as a0eebc99-…, a date as 2024-02-29, a bytea as \xdeadbeef, and one template gives the same output from either source. Values read with .GetDynamicValue keep their Go time.Time type so the sprig date functions still work on them.
Two related bulk-ingest fixes: a transformer that turns a range column into its literal ([1,10) from template or literal_string) failed binary COPY with cannot find encode plan; the literal is now scanned into the typed range pgx can encode, for int4range, int8range, numrange, daterange, tsrange and tstzrange. And the array transformer built its own type map that knew only built-in OIDs, so a rule on citext[] or a domain over an array passed validation and then nulled 100% of rows at run time under the default on_error: null, one DATALOSS line each, while the run reported success. It now uses the catalog-aware map from the connection, which also installs the raw JSON codec so jsonb[] elements are no longer round-tripped through float64.
A snapshot with transformation rules could stall forever, logging data exception: value too long for type character varying(50) with a growing retry delay and never completing or erroring. The message named only the type, so it read as if a TEXT column were limited to 50 characters. Two causes: MapError kept PgError.Message and dropped the CONTEXT line (COPY ledger_line, line 1, column code) that names the table and column; and class 22 errors were not in the permanent set, so a failure raised by the row values themselves was replayed identically under a retry policy with no attempt bound. The context is now appended to the error, and ErrDataException is no longer retried.
pgstream validate rules reported a failing transformer without saying which rule produced it. Errors now carry the rule index, table and column:
table_transformers[2]: column 'name' in table "public"."test": template_transformer: error parsing template: template: :5: function "literal_string" not defined
Full Changelog: v1.4.2...v1.5.0
Nothing published for this version
🧭 Replicated DDL is applied under the schema it was issued in, so unqualified CREATE TABLE no longer lands in public on the target 🕰 timetz columns ar
⚡ Highlights
🧭 Replicated DDL is applied under the schema it was issued in, so unqualified CREATE TABLE no longer lands in public on the target
🕰 timetz columns are supported
🛂 pgstream check no longer fails its CREATEROLE check with database "<role>" does not exist
🔌 Dropping the database from a connection string no longer drops the user with it, on the pg_restore --create path too
Dependency updates
timetz by @kvch in #1135Full Changelog: v1.4.1...v1.4.2
Nothing published for this version
🧬 Replicated bytea values are no longer corrupted on the way to a Postgres target 🧩 Partitioned tables restore with valid indexes and primary keys, so
⚡ Highlights
🧬 Replicated bytea values are no longer corrupted on the way to a Postgres target
🧩 Partitioned tables restore with valid indexes and primary keys, so they keep replicating after a snapshot
🔐 Binaries are built with Go 1.26.6, patching stdlib and module CVEs
Dependency updates
ALTER INDEX ... ATTACH PARTITION with indices and constraints by @siriusfrk in #1094Full Changelog: v1.4.0...v1.4.1
📊 Native Prometheus scrape endpoint 🎯 Granular object type filtering for pg-to-pg schema snapshots and DDL replication 🛡️ Transformation rules are now
⚡ Highlights
📊 Native Prometheus scrape endpoint
🎯 Granular object type filtering for pg-to-pg schema snapshots and DDL replication
🛡️ Transformation rules are now validated against unique indexes before a run starts, instead of failing mid-load
🔌 Target connection pool limits are configurable
🪝 Webhook delivery retries transient failures and no longer silently drops still-failing events; adds idempotency headers
🐛 Fixed a DDL parsing bug that silently dropped RENAME COLUMN / DROP COLUMN schema changes
🛂 New preflight check for CREATEROLE when a snapshot restores roles
Dependency updates
--reset snapshot behaviour by @kvch in #1059The setting instrumentation.metrics.prometheus.enabled exposes a /metrics endpoint on the existing health server. Metrics can now be scraped directly, with no collector in between:
instrumentation:
metrics:
prometheus:
enabled: true # Defaults to false
endpoint: "/metrics" # Defaults to /metricsThe endpoint returns 404 when disabled.
Both the schema snapshot and DDL replication for pg-to-pg pipelines can now be restricted to a subset of object categories (tables, sequences, types, indexes, functions, views, triggers, etc.) via include_object_types (allowlist) or exclude_object_types (denylist), independently for snapshot and replication. A DDL statement that touches several object types (e.g. CREATE TABLE with a primary key) is only skipped when every object it touches is excluded.
Anonymization is lossy by design: masking with type: id keeps a 6-character prefix, so every value sharing that prefix collapses to the same masked value. Applied to a column covered by a unique index, this previously surfaced as duplicate key value violates unique constraint hours into a data load. Every transformer is now classified as preserved (distinct inputs stay distinct), not_guaranteed (random or hashed output, collisions possible) or lossy (collides by construction).
validate rules and pipeline startup now read pg_index for the tables in scope and check the configured rules against unique indexes, primary keys and unique constraints. With a Postgres target, a lossy transformer on a covered column fails at parse time with transformation rules break a unique index; with other targets (Kafka, Elasticsearch/OpenSearch, webhooks) the same finding is a warning, since there is no unique index to violate. A not_guaranteed transformer always warns. Set allow_uniqueness_loss: true on a column rule to keep a lossy transformer anyway. A custom mask configured to mask zero characters (mask_begin equal to mask_end, or unmask_begin: 0) is now rejected at construction rather than silently passing values through unmasked. Use noop to pass a column through on purpose.
target.postgres.max_connections (PGSTREAM_POSTGRES_WRITER_MAX_CONNECTIONS) now controls the writer pool size, defaulting to 50 and honoring pool_max_conns from the target URL when set explicitly. Previously this setting also sized the schema observer's pool, so a target capped at 90 connections to stay under a 100-connection server limit could still open up to 180. The schema observer is now sized separately, at min(writer pool, 16). So the process now opens at most the configured writer limit plus 16.
A failed webhook POST is now retried with configurable backoff (network errors, 429, 5xx); other non-2xx responses are treated as permanent and not retried. A delivery that is still failing after retries is logged and dropped so one broken subscriber can't block delivery to everyone else, and at-least-once semantics apply on restart. New X-Pgstream-LSN / Idempotency-Key headers let subscribers dedupe redelivered events (omitted for snapshot events, which share a zero LSN). DISABLE_RETRIES now actually disables retries.
Full Changelog: v1.3.1...v1.4.0
📌 Data snapshots read the columns the schema snapshot captured, so a column added to the source mid-snapshot no longer breaks the load ⚠️ Schema drift
⚡ Highlights
📌 Data snapshots read the columns the schema snapshot captured, so a column added to the source mid-snapshot no longer breaks the load
⚠️ Schema drift during a snapshot is now reported: a warning when it happens during the dump, an error when a table's captured columns are gone
🛂 New preflight check for CREATEDB on the target when the snapshot creates the target database
🔐 Documented how DDL replication can carry privileges from a less trusted source to the target
Dependency updates
Full Changelog: v1.3.0...v1.3.1
474ff1c Added tests and docs for snapshot+stream from a read replica
Full Changelog: v1.2.5...v1.3.0
🎲 Template and JSON transformers no longer corrupt their random generator under multi-worker snapshots 📉 New metrics and totals for what ignore_send_e
⚡ Highlights
🎲 Template and JSON transformers no longer corrupt their random generator under multi-worker snapshots
📉 New metrics and totals for what ignore_send_errors silently discards, plus a startup warning when suppression is on
🧾 Snapshot failures are logged as errors instead of rendering as {}
🔁 Snapshot request status updates are retried, and report data loss when they still cannot be persisted
🔍 A failed restore now carries the output that explains it, with credentials redacted
No configuration changes are required.
Contributors. make lint now installs golangci-lint v2.10.1 to match CI, which reports findings the previously pinned v2.5.0 did not.
Full Changelog: v1.2.4...v1.2.5
🔑 Replicated DELETE s on tables with a user-defined enum primary key no longer fail 🧬 Enum arrays and domains over enums are now handled in delete ide
⚡ Highlights
🔑 Replicated DELETEs on tables with a user-defined enum primary key no longer fail
🧬 Enum arrays and domains over enums are now handled in delete identities, not just plain enums
🛟 A failed query can no longer take its whole batch down with it, including writes to unrelated tables
No configuration changes are required, and the SQL emitted for tables without enum identity columns is unchanged.
Full Changelog: v1.2.3...v1.2.4
🎨 Bulk ingest no longer fails on tables with user-defined enum columns 🔑 Kafka primary_key partitioning no longer collides distinct rows onto one mess
⚡ Highlights
🎨 Bulk ingest no longer fails on tables with user-defined enum columns
🔑 Kafka primary_key partitioning no longer collides distinct rows onto one message key
🩺 Two preflight checks that reported false failures on correctly-configured sources (wal2json, replica identity) are fixed
🔍 Preflight can now run source-only checks (--source) and verifies major PostgreSQL version compatibility
Dependency updates
Full Changelog: v1.2.2...v1.2.3
🗜 Configurable session settings for the snapshot index/constraint restore phase (via PGOPTIONS ) 🛡 Index constraint settings validated at startup 🩺 Pr
⚡ Highlights
🗜 Configurable session settings for the snapshot index/constraint restore phase (via PGOPTIONS)
🛡 Index constraint settings validated at startup
🩺 Preflight now warns when a snapshot source looks load-balanced (Aurora/RDS reader, instance-spanning pooler)
💥 Fixed a stack-depth overflow on bulk deletes of tables with composite primary keys
🛑 pg_restore now honors context cancellation, so shutdowns and timeouts no longer hang
Full Changelog: v1.2.1...v1.2.2
📊 Pipeline phase visibility via /status and a new OTel gauge 🧬 Preflight checks no longer reports custom enums as unsupported types 🗜 Snapshot request
⚡ Highlights
📊 Pipeline phase visibility via /status and a new OTel gauge
🧬 Preflight checks no longer reports custom enums as unsupported types
🗜 Snapshot request unique index hashes table_names: fixes btree overflow with many tables
🐛 Bug fixes (deferred parsing of validated schema objects, comment restore downgraded to warning)
Dependency updates
/status and OTel gauge by @Shaurya2k06 in #1000table_names in snapshot request unique index to avoid btree overflow by @kvch in #1013Full Changelog: v1.2.0...v1.2.1
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Your coding agent can read these notes before it upgrades. Set up the MCP server →