From cf0363a4e13e090a53ad0e4d964802117acfcc28 Mon Sep 17 00:00:00 2001 From: avison9 Date: Fri, 25 Sep 2026 22:25:43 +0100 Subject: [PATCH 1/3] A table.include.list entry that matches no table is raised, with the corrected entry when the mistake is a missing schema or a shell glob Debezium logs a warning when a table.include.list entry matches nothing and keeps running, so the topic is simply never produced. People ask it to fail instead (debezium/dbz#872); a maintainer answered that a table may be created later. cdclint raises it as a warning for the same reason: when a sink actually reads the table, sink-table-not-captured already fails the pull request. The two shapes behind the most-viewed questions get a pointed fix, because both follow from Debezium's documented matching (each entry is a regular expression matched against the whole schema.table name): an entry without its schema (Stack Overflow 74103659: ipaddrs for myschema.ipaddrs) gets "write myschema\.ipaddrs", and a shell glob (Stack Overflow 51345636: public.bg_* names no table as a regex) gets "write public\.bg_.* to capture ..." with the tables it would capture. An entry that already contains .* is read as the regex it is. Exclude-list entries are not checked: one that matches nothing excludes nothing, the same reasoning as captured-column-missing. Three corpus entries, written before the code: include-table-no-schema, include-table-glob, include-table-typo. --- README.md | 4 +- corpus/README.md | 3 + corpus/include-table-glob/connector.json | 9 ++ corpus/include-table-glob/expected.txt | 9 ++ .../migrations/0001_tables.sql | 21 ++++ .../sink/0001_bg_orders.sql | 12 +++ corpus/include-table-no-schema/connector.json | 9 ++ corpus/include-table-no-schema/expected.txt | 9 ++ .../migrations/0001_ipaddrs.sql | 11 +++ .../sink/0001_ipaddrs.sql | 14 +++ corpus/include-table-typo/connector.json | 9 ++ corpus/include-table-typo/expected.txt | 8 ++ .../migrations/0001_tables.sql | 14 +++ .../include-table-typo/sink/0001_orders.sql | 12 +++ internal/capture/debezium/debezium.go | 11 +++ internal/engine/engine.go | 96 ++++++++++++++++++- 16 files changed, 247 insertions(+), 4 deletions(-) create mode 100644 corpus/include-table-glob/connector.json create mode 100644 corpus/include-table-glob/expected.txt create mode 100644 corpus/include-table-glob/migrations/0001_tables.sql create mode 100644 corpus/include-table-glob/sink/0001_bg_orders.sql create mode 100644 corpus/include-table-no-schema/connector.json create mode 100644 corpus/include-table-no-schema/expected.txt create mode 100644 corpus/include-table-no-schema/migrations/0001_ipaddrs.sql create mode 100644 corpus/include-table-no-schema/sink/0001_ipaddrs.sql create mode 100644 corpus/include-table-typo/connector.json create mode 100644 corpus/include-table-typo/expected.txt create mode 100644 corpus/include-table-typo/migrations/0001_tables.sql create mode 100644 corpus/include-table-typo/sink/0001_orders.sql diff --git a/README.md b/README.md index 6ed9833..543e302 100644 --- a/README.md +++ b/README.md @@ -53,9 +53,10 @@ your repository, and cdclint reads it: | A column is `0`, `''`, `NULL` or `1970-01-01` on every row in ClickHouse, BigQuery, Snowflake or Iceberg, with nothing in the logs | the column is not on the connector's `column.include.list` (or is on its exclude list), so Debezium drops it before Kafka | `sink-column-not-captured` | | A column added in Postgres never shows up downstream | the migration added it and nobody added it to the include list | `schema-before-connector` (with `--base`), `sink-column-not-captured` | | Listing the columns of one table made another table's columns disappear | `column.include.list` is one list for every captured table | `sink-column-not-captured` | -| "`table.include.list` not working", or the connector is `RUNNING` and no topic appears | the entry has no schema (`orders` for `public.orders`), or is a glob (`public.bg_*`) where Debezium expects a regex (`public\.bg_.*`); Debezium matches each entry against the whole `schema.table` name, never a substring | `sink-table-not-captured`, when a sink reads the table | +| "`table.include.list` not working", or the connector is `RUNNING` and no topic appears | the entry has no schema (`orders` for `public.orders`), or is a glob (`public.bg_*`) where Debezium expects a regex (`public\.bg_.*`); Debezium matches each entry against the whole `schema.table` name, never a substring | `captured-table-missing` names the entry and the fix; `sink-table-not-captured` when a sink reads the table | | A Kafka-engine table or sink receives nothing, or the wrong table fills | the topic it reads is not the one the connector produces (prefix, schema, `RegexRouter`) | `topic-table-mapping` | | A typo in the include list, and a column quietly missing | the pattern matches no column in the source | `captured-column-missing` | +| A typo in `table.include.list`, and a topic that never appears | the entry matches no table in the source | `captured-table-missing` | | A ClickHouse materialized view writes defaults | a refreshable view's `SELECT` order differs from the target's (it matches by position), or a streaming view's names differ (it matches by name) | `mv-column-match` | It reads files only: Postgres sources today, no type checks, no connection to @@ -70,6 +71,7 @@ anything running. | `sink-column-unknown` | the sink expects a field the source table does not have (warning; renames and computed fields are legitimate) | v0.1 | | `source-column-not-captured` | a source column nothing captures and nothing reads yet, so the day something asks for it is the day it is found missing (info) | v0.1 | | `captured-column-missing` | the include list names a column the source does not have (warning) | v0.1 | +| `captured-table-missing` | a `table.include.list` entry matches no table: a typo, a missing schema (`orders` for `public.orders`), or a shell glob (`public.bg_*`) where Debezium reads a regular expression; the last two get the corrected entry as the fix (warning) | next | | `topic-table-mapping` | a Kafka-engine table reads a topic the connector will not produce | v0.1 | | `mv-column-match` | ClickHouse streaming materialized views match by name, refreshable ones by position. ClickHouse 25.4 and later reject a streaming view that writes a column the target lacks when it is created; a refreshable view's order mismatch was loud on 24.8 and is silent on 26.8 | v0.1 | | `schema-before-connector` | this change adds a column to a captured table and leaves it out of the stream without deciding to, the trap itself, judged on the diff (warning: leaving PII off is right, so it asks for the decision) | v0.2 | diff --git a/corpus/README.md b/corpus/README.md index c2dddf5..41fdbbe 100644 --- a/corpus/README.md +++ b/corpus/README.md @@ -28,6 +28,9 @@ The rest exercise one path each: | `include-list-typo` | an include-list pattern that matches no column, and the read it silently breaks | | `golang-migrate-down-files` | Mattermost v10.11.0's `000092_add_createat_to_teammembers.down.sql` drops a column from the wrong table; read in name order it ran just before its up file and deleted `reactions.createat`. Down files are skipped | | `goose-down-section` | a goose file's down section follows its up section; applied whole, every table was created and dropped again. Down sections are skipped | +| `include-table-no-schema` | Stack Overflow 74103659: `table.include.list` names `ipaddrs` without its schema, the connector runs and no topic appears | +| `include-table-glob` | Stack Overflow 51345636: `public.bg_*` written as a shell glob; as a regular expression it names no table | +| `include-table-typo` | a misspelled `table.include.list` entry, which Debezium only logs as a warning (debezium/dbz#872) | | `diff-adds-column-connector-untouched` | the diff rule: `base/` holds the migrations and connector before the change; the change adds a column and leaves the connector alone | | `diff-connector-captures-one-of-two` | #963 with a hurried fix: two of three new columns go on the include list in the same change; the third is raised, since adding two says nothing about it | | `diff-connector-touched-other-table` | RefuseRadar #962 and #963 in one range: the connector gains report_validations columns, the migration adds reports columns; the reports ones are raised | diff --git a/corpus/include-table-glob/connector.json b/corpus/include-table-glob/connector.json new file mode 100644 index 0000000..96dc386 --- /dev/null +++ b/corpus/include-table-glob/connector.json @@ -0,0 +1,9 @@ +{ + "name": "shop-source", + "config": { + "connector.class": "io.debezium.connector.postgresql.PostgresConnector", + "plugin.name": "pgoutput", + "topic.prefix": "shop", + "table.include.list": "public.bg_*" + } +} diff --git a/corpus/include-table-glob/expected.txt b/corpus/include-table-glob/expected.txt new file mode 100644 index 0000000..121f076 --- /dev/null +++ b/corpus/include-table-glob/expected.txt @@ -0,0 +1,9 @@ +warning captured-table-missing include-table-glob/connector.json + public.bg_* in table.include.list matches no table in the source schema + Debezium reads each entry as a regular expression, where * repeats the character before it + fix: write public\.bg_.* to capture public.bg_items and public.bg_orders, or remove it +error sink-table-not-captured include-table-glob/sink/0001_bg_orders.sql:1 + kafka_bg_orders reads topic shop.public.bg_orders, but public.bg_orders is not matched by table.include.list in include-table-glob/connector.json + nothing will ever arrive on that topic + fix: add public.bg_orders to table.include.list and deploy the connector before the sink schema +1 error(s), 1 warning(s), 0 info diff --git a/corpus/include-table-glob/migrations/0001_tables.sql b/corpus/include-table-glob/migrations/0001_tables.sql new file mode 100644 index 0000000..67aa027 --- /dev/null +++ b/corpus/include-table-glob/migrations/0001_tables.sql @@ -0,0 +1,21 @@ +-- Reduced from Stack Overflow question 51345636 (5.9k views): the include +-- list was written as a shell glob, public.bg_*, meaning every table whose +-- name starts with bg_. Debezium reads it as a regular expression, where * +-- repeats the character before it, so it names public.bg, public.bg_, +-- public.bg__ and no real table. +CREATE TABLE bg_orders ( + id BIGSERIAL PRIMARY KEY, + amount NUMERIC(12, 2) NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE bg_items ( + id BIGSERIAL PRIMARY KEY, + order_id BIGINT NOT NULL REFERENCES bg_orders (id), + sku TEXT NOT NULL +); + +CREATE TABLE cp_users ( + id BIGSERIAL PRIMARY KEY, + email TEXT NOT NULL +); diff --git a/corpus/include-table-glob/sink/0001_bg_orders.sql b/corpus/include-table-glob/sink/0001_bg_orders.sql new file mode 100644 index 0000000..3343339 --- /dev/null +++ b/corpus/include-table-glob/sink/0001_bg_orders.sql @@ -0,0 +1,12 @@ +CREATE TABLE kafka_bg_orders +( + `id` Int64, + `amount` String, + `created_at` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shop.public.bg_orders', + kafka_group_name = 'clickhouse-bg-orders', + kafka_format = 'JSONEachRow'; diff --git a/corpus/include-table-no-schema/connector.json b/corpus/include-table-no-schema/connector.json new file mode 100644 index 0000000..2714711 --- /dev/null +++ b/corpus/include-table-no-schema/connector.json @@ -0,0 +1,9 @@ +{ + "name": "testdb", + "config": { + "connector.class": "io.debezium.connector.postgresql.PostgresConnector", + "plugin.name": "pgoutput", + "topic.prefix": "test", + "table.include.list": "ipaddrs" + } +} diff --git a/corpus/include-table-no-schema/expected.txt b/corpus/include-table-no-schema/expected.txt new file mode 100644 index 0000000..79c00e0 --- /dev/null +++ b/corpus/include-table-no-schema/expected.txt @@ -0,0 +1,9 @@ +warning captured-table-missing include-table-no-schema/connector.json + ipaddrs in table.include.list matches no table in the source schema + Debezium matches each entry against the whole schema.table name + fix: write myschema\.ipaddrs, or remove it +error sink-table-not-captured include-table-no-schema/sink/0001_ipaddrs.sql:3 + kafka_ipaddrs reads topic test.myschema.ipaddrs, but myschema.ipaddrs is not matched by table.include.list in include-table-no-schema/connector.json + nothing will ever arrive on that topic + fix: add myschema.ipaddrs to table.include.list and deploy the connector before the sink schema +1 error(s), 1 warning(s), 0 info diff --git a/corpus/include-table-no-schema/migrations/0001_ipaddrs.sql b/corpus/include-table-no-schema/migrations/0001_ipaddrs.sql new file mode 100644 index 0000000..de36077 --- /dev/null +++ b/corpus/include-table-no-schema/migrations/0001_ipaddrs.sql @@ -0,0 +1,11 @@ +-- Reduced from Stack Overflow question 74103659 ("Debezium PostgreSQL +-- connector not creating topic", 6.9k views) and the Debezium Google Group +-- thread Kp_9jbTTEkY: the table lives in a schema, the include list names it +-- without one, the connector reports RUNNING and nothing is ever captured. +CREATE SCHEMA myschema; + +CREATE TABLE myschema.ipaddrs ( + id BIGSERIAL PRIMARY KEY, + address INET NOT NULL, + seen_at TIMESTAMPTZ NOT NULL DEFAULT now() +); diff --git a/corpus/include-table-no-schema/sink/0001_ipaddrs.sql b/corpus/include-table-no-schema/sink/0001_ipaddrs.sql new file mode 100644 index 0000000..1c86740 --- /dev/null +++ b/corpus/include-table-no-schema/sink/0001_ipaddrs.sql @@ -0,0 +1,14 @@ +-- The warehouse side, waiting on the topic Debezium would produce for +-- myschema.ipaddrs once the table is captured. +CREATE TABLE kafka_ipaddrs +( + `id` Int64, + `address` String, + `seen_at` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'test.myschema.ipaddrs', + kafka_group_name = 'clickhouse-ipaddrs', + kafka_format = 'JSONEachRow'; diff --git a/corpus/include-table-typo/connector.json b/corpus/include-table-typo/connector.json new file mode 100644 index 0000000..ef131ba --- /dev/null +++ b/corpus/include-table-typo/connector.json @@ -0,0 +1,9 @@ +{ + "name": "shop-source", + "config": { + "connector.class": "io.debezium.connector.postgresql.PostgresConnector", + "plugin.name": "pgoutput", + "topic.prefix": "shop", + "table.include.list": "public.customers,public.oders" + } +} diff --git a/corpus/include-table-typo/expected.txt b/corpus/include-table-typo/expected.txt new file mode 100644 index 0000000..f19d386 --- /dev/null +++ b/corpus/include-table-typo/expected.txt @@ -0,0 +1,8 @@ +warning captured-table-missing include-table-typo/connector.json + public.oders in table.include.list matches no table in the source schema + fix: remove it, or check the spelling against the migrations +error sink-table-not-captured include-table-typo/sink/0001_orders.sql:1 + kafka_orders reads topic shop.public.orders, but public.orders is not matched by table.include.list in include-table-typo/connector.json + nothing will ever arrive on that topic + fix: add public.orders to table.include.list and deploy the connector before the sink schema +1 error(s), 1 warning(s), 0 info diff --git a/corpus/include-table-typo/migrations/0001_tables.sql b/corpus/include-table-typo/migrations/0001_tables.sql new file mode 100644 index 0000000..e7c5914 --- /dev/null +++ b/corpus/include-table-typo/migrations/0001_tables.sql @@ -0,0 +1,14 @@ +-- The plain shape: one include-list entry misspelled. Debezium logs a +-- warning when an entry matches nothing and keeps running (Debezium issue +-- dbz#872 asks it to fail instead), so the orders topic is simply never +-- produced. +CREATE TABLE customers ( + id BIGSERIAL PRIMARY KEY, + email TEXT NOT NULL +); + +CREATE TABLE orders ( + id BIGSERIAL PRIMARY KEY, + customer_id BIGINT NOT NULL REFERENCES customers (id), + total NUMERIC(12, 2) NOT NULL +); diff --git a/corpus/include-table-typo/sink/0001_orders.sql b/corpus/include-table-typo/sink/0001_orders.sql new file mode 100644 index 0000000..e755bb2 --- /dev/null +++ b/corpus/include-table-typo/sink/0001_orders.sql @@ -0,0 +1,12 @@ +CREATE TABLE kafka_orders +( + `id` Int64, + `customer_id` Int64, + `total` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shop.public.orders', + kafka_group_name = 'clickhouse-orders', + kafka_format = 'JSONEachRow'; diff --git a/internal/capture/debezium/debezium.go b/internal/capture/debezium/debezium.go index 2ef7f03..5b6dc31 100644 --- a/internal/capture/debezium/debezium.go +++ b/internal/capture/debezium/debezium.go @@ -33,6 +33,7 @@ type Contract struct { columnListKey string tableListKey string columnPatterns []string + tablePatterns []string routers []router Config map[string]string Class string @@ -95,6 +96,11 @@ func Parse(b []byte) (*Contract, error) { if c.tableInclude, c.tableListKey, err = list(cfg, "table.include.list", "table.whitelist"); err != nil { return nil, err } + for _, p := range strings.Split(first(cfg, "table.include.list", "table.whitelist"), ",") { + if p = strings.TrimSpace(p); p != "" { + c.tablePatterns = append(c.tablePatterns, p) + } + } if c.tableExclude, _, err = list(cfg, "table.exclude.list", "table.blacklist"); err != nil { return nil, err } @@ -232,4 +238,9 @@ func (c *Contract) ColumnListSetting() string { return c.columnListKey } // captured-column-missing rule. Exclude lists are not returned: a pattern // there that matches nothing excludes nothing, which is harmless. func (c *Contract) ColumnPatterns() []string { return c.columnPatterns } + +// TablePatterns returns the table include-list patterns as written, for the +// captured-table-missing rule. Exclude lists are left out for the same +// reason as columns. +func (c *Contract) TablePatterns() []string { return c.tablePatterns } func (c *Contract) TableListSetting() string { return c.tableListKey } diff --git a/internal/engine/engine.go b/internal/engine/engine.go index 8370e92..199eb36 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -6,6 +6,8 @@ package engine import ( "fmt" + "regexp" + "sort" "strings" "github.com/avison9/cdclint/internal/model" @@ -18,7 +20,7 @@ type Input struct { Source *model.Source Contract model.Contract // Patterns, when the contract can list them, lets captured-column-missing - // check each include pattern against the schema. + // and captured-table-missing check each include pattern against the schema. Patterns PatternLister Sinks []*model.Sink Views []clickhouse.View @@ -27,10 +29,11 @@ type Input struct { Base *Base } -// PatternLister is implemented by contracts whose column list is a set of -// patterns that can be checked one by one. +// PatternLister is implemented by contracts whose column and table lists are +// sets of patterns that can be checked one by one. type PatternLister interface { ColumnPatterns() []string + TablePatterns() []string } // Read is one sink column resolved to the source column it expects. @@ -134,6 +137,7 @@ func Run(in *Input) []model.Finding { fs = append(fs, diff...) fs = append(fs, sourceColumnNotCaptured(in, reads, raised)...) fs = append(fs, capturedColumnMissing(in)...) + fs = append(fs, capturedTableMissing(in)...) fs = append(fs, mvColumnMatch(in)...) model.Sort(fs) return fs @@ -291,6 +295,92 @@ func capturedColumnMissing(in *Input) []model.Finding { return fs } +// capturedTableMissing checks each table include-list entry against the +// schema. Debezium logs a warning when an entry matches no table and keeps +// running; debezium/dbz#872 asks it to fail instead, and a maintainer +// answered that a table may be created later. That is why this is a warning: +// when a sink reads the table, sink-table-not-captured raises the error. +// +// Two shapes get a pointed fix because they are how people get it wrong in +// practice (Stack Overflow 74103659 and 51345636): an entry without its schema, +// and a shell glob. Both follow from Debezium's documented matching: each +// entry is a regular expression matched against the whole schema.table name, +// never a substring. +func capturedTableMissing(in *Input) []model.Finding { + if in.Patterns == nil { + return nil + } + var fs []model.Finding + for _, p := range in.Patterns.TablePatterns() { + for _, name := range expand(p) { + m, err := compileAnchored(name) + if err != nil || tablesMatching(in, func(t *model.Table) bool { return m(t.Qualified()) }) != nil { + continue + } + message := fmt.Sprintf("%s in %s matches no table in the source schema", name, in.Contract.TableListSetting()) + fix := "remove it, or check the spelling against the migrations" + if bare := tablesMatching(in, func(t *model.Table) bool { return m(t.Name) }); bare != nil { + message += "\nDebezium matches each entry against the whole schema.table name" + var qualified []string + for _, t := range bare { + qualified = append(qualified, regexp.QuoteMeta(t.Qualified())) + } + fix = "write " + strings.Join(qualified, " or ") + ", or remove it" + } else if re, caught := globReading(in, name); caught != nil { + message += "\nDebezium reads each entry as a regular expression, where * repeats the character before it" + var names []string + for _, t := range caught { + names = append(names, t.Qualified()) + } + fix = "write " + re + " to capture " + joinAnd(names) + ", or remove it" + } + fs = append(fs, model.Finding{ + Rule: "captured-table-missing", Severity: model.Warning, Pos: in.Contract.Pos(), + Message: message, + Fix: fix, + }) + } + } + return fs +} + +// tablesMatching returns the source tables for which match is true, in +// qualified-name order, or nil when there are none. +func tablesMatching(in *Input, match func(*model.Table) bool) []*model.Table { + var out []*model.Table + for _, t := range in.Source.Tables { + if match(t) { + out = append(out, t) + } + } + sort.Slice(out, func(i, j int) bool { return out[i].Qualified() < out[j].Qualified() }) + return out +} + +// globReading reads an entry the way its author most likely meant it, as a +// shell glob where * is any run of characters, and returns the regular +// expression that says so and the tables it would capture. An entry that +// already contains .* was written as a regular expression and is left alone. +func globReading(in *Input, entry string) (string, []*model.Table) { + if !strings.Contains(entry, "*") || strings.Contains(entry, ".*") { + return "", nil + } + re := strings.ReplaceAll(regexp.QuoteMeta(entry), `\*`, ".*") + m, err := compileAnchored(re) + if err != nil { + return "", nil + } + return re, tablesMatching(in, func(t *model.Table) bool { return m(t.Qualified()) }) +} + +// joinAnd lists names as "a", "a and b" or "a, b and c". +func joinAnd(names []string) string { + if len(names) < 2 { + return strings.Join(names, "") + } + return strings.Join(names[:len(names)-1], ", ") + " and " + names[len(names)-1] +} + // expand turns the common `schema.table\.(a|b|c)` shape into one pattern // per alternative so each column is checked on its own. Any other shape is // returned as is. From 6133eb3aadd64ab9d491deb9c96bc173579ffa27 Mon Sep 17 00:00:00 2001 From: avison9 Date: Fri, 25 Sep 2026 22:36:05 +0100 Subject: [PATCH 2/3] Sink columns are matched to the field that actually arrives after the connectors' transforms, and a table that misses a flatten is raised A Flatten transform on the Debezium envelope names every field after., and sinks match fields to columns by name. A ClickHouse table that names its columns as Postgres does therefore receives a row of zeros and empty strings for every change, with no error (ClickHouse/clickhouse-kafka-connect discussion 182). cdclint said "ok" on that exact configuration, because it took every sink column name as a row column name. The same assumption made it raise twelve false sink-column-unknown warnings on ClickHouse's own CDC example, whose table names its columns before. and after., and miss that one of those columns was never captured. A new internal/smt package works out the value's shape from each connector's transforms, in order: the envelope for a Debezium connector (or after.state.only false), the row after ExtractNewRecordState or the Iceberg DebeziumTransform, and the flattened envelope after Flatten, with its delimiter. The sink connector's transforms apply after the source's. When the result is a flattened envelope, after. and before. resolve to source column c, so the capture rules see through them; op, ts_ms and source/transaction fields are envelope metadata; and a column named as the row names it is raised by the new sink-column-flattened rule, once per table. Only documented transforms are modelled. Routers, key transforms and Filter change no field name; any other value transform, or a modelled one under a predicate, makes the shape unknown, and an unknown shape keeps the old behaviour and produces no new finding. Rejected: raising a bare envelope with no transform at all. Whether a sink handles Debezium's envelope itself, or a converter reshapes it, is not in the files, so that would be a guess. Corpus first: flatten-bare-columns (discussion 182's configs), flatten-envelope-columns (ClickHouse/examples cdc/postgresql at ae417ae, with an include list that leaves postcode2 out), flatten-after-unwrap (a guard). RefuseRadar, which unwraps, is unchanged at 0/0/159. --- README.md | 13 +++ corpus/README.md | 3 + corpus/flatten-after-unwrap/connector.json | 13 +++ corpus/flatten-after-unwrap/expected.txt | 1 + .../migrations/0001_orders.sql | 6 ++ .../flatten-after-unwrap/sink-connector.json | 6 ++ .../flatten-after-unwrap/sink/0001_orders.sql | 7 ++ corpus/flatten-bare-columns/connector.json | 21 ++++ corpus/flatten-bare-columns/expected.txt | 5 + .../migrations/0001_testradiomeas.sql | 11 ++ .../flatten-bare-columns/sink-connector.json | 9 ++ .../sink/0001_testradiomeas.sql | 13 +++ .../flatten-envelope-columns/connector.json | 23 ++++ corpus/flatten-envelope-columns/expected.txt | 9 ++ .../migrations/0001_uk_price_paid.sql | 12 +++ .../sink-connector.json | 10 ++ .../sink/0001_uk_price_paid_changes.sql | 17 +++ internal/capture/debezium/debezium.go | 13 +++ internal/engine/diff.go | 4 +- internal/engine/engine.go | 102 ++++++++++++++++-- internal/model/model.go | 27 +++++ internal/sink/connect/connect.go | 10 +- internal/smt/smt.go | 94 ++++++++++++++++ internal/smt/smt_test.go | 85 +++++++++++++++ 24 files changed, 504 insertions(+), 10 deletions(-) create mode 100644 corpus/flatten-after-unwrap/connector.json create mode 100644 corpus/flatten-after-unwrap/expected.txt create mode 100644 corpus/flatten-after-unwrap/migrations/0001_orders.sql create mode 100644 corpus/flatten-after-unwrap/sink-connector.json create mode 100644 corpus/flatten-after-unwrap/sink/0001_orders.sql create mode 100644 corpus/flatten-bare-columns/connector.json create mode 100644 corpus/flatten-bare-columns/expected.txt create mode 100644 corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql create mode 100644 corpus/flatten-bare-columns/sink-connector.json create mode 100644 corpus/flatten-bare-columns/sink/0001_testradiomeas.sql create mode 100644 corpus/flatten-envelope-columns/connector.json create mode 100644 corpus/flatten-envelope-columns/expected.txt create mode 100644 corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql create mode 100644 corpus/flatten-envelope-columns/sink-connector.json create mode 100644 corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql create mode 100644 internal/smt/smt.go create mode 100644 internal/smt/smt_test.go diff --git a/README.md b/README.md index 543e302..1376a63 100644 --- a/README.md +++ b/README.md @@ -57,6 +57,7 @@ your repository, and cdclint reads it: | A Kafka-engine table or sink receives nothing, or the wrong table fills | the topic it reads is not the one the connector produces (prefix, schema, `RegexRouter`) | `topic-table-mapping` | | A typo in the include list, and a column quietly missing | the pattern matches no column in the source | `captured-column-missing` | | A typo in `table.include.list`, and a topic that never appears | the entry matches no table in the source | `captured-table-missing` | +| The Kafka Connect sink writes a row of `0` and empty strings for every change, and no error ("sink replicates zero values or nulls") | a `Flatten` transform keeps Debezium's envelope and names every field `after.`, while the sink table names its columns as Postgres does; sinks match fields to columns by name | `sink-column-flattened` | | A ClickHouse materialized view writes defaults | a refreshable view's `SELECT` order differs from the target's (it matches by position), or a streaming view's names differ (it matches by name) | `mv-column-match` | It reads files only: Postgres sources today, no type checks, no connection to @@ -73,6 +74,7 @@ anything running. | `captured-column-missing` | the include list names a column the source does not have (warning) | v0.1 | | `captured-table-missing` | a `table.include.list` entry matches no table: a typo, a missing schema (`orders` for `public.orders`), or a shell glob (`public.bg_*`) where Debezium reads a regular expression; the last two get the corrected entry as the fix (warning) | next | | `topic-table-mapping` | a Kafka-engine table reads a topic the connector will not produce | v0.1 | +| `sink-column-flattened` | a sink table names its columns as the source does while a `Flatten` transform on the way delivers the envelope as `after.`, so every column stays at its default | next | | `mv-column-match` | ClickHouse streaming materialized views match by name, refreshable ones by position. ClickHouse 25.4 and later reject a streaming view that writes a column the target lacks when it is created; a refreshable view's order mismatch was loud on 24.8 and is silent on 26.8 | v0.1 | | `schema-before-connector` | this change adds a column to a captured table and leaves it out of the stream without deciding to, the trap itself, judged on the diff (warning: leaving PII off is right, so it asks for the decision) | v0.2 | | `replica-identity` | a captured table's replica identity cannot supply what the sink reads | next | @@ -215,6 +217,17 @@ same name, and the finding says so. Debezium's `RegexRouter` transform is applied when computing topic names. Other sources (MySQL, SQL Server) and sinks are packages behind the same two interfaces. +Transforms on either connector decide what the sink receives, and cdclint +follows the documented ones in order: Debezium's `ExtractNewRecordState` (or the +Iceberg `DebeziumTransform`) unwraps the envelope into the row, and Kafka +Connect's `Flatten` turns the envelope into `after.`, `before.`, +`source.`, `op` and `ts_ms` (its `delimiter` included). A sink table +named that way, as in ClickHouse's own CDC example, is matched through +`after.` and `before.` to the source columns. Routers, key transforms and +`Filter` change no field name. Any other transform that touches the value, or a +modelled one applied under a predicate, leaves the shape unknown, and cdclint +then reads the sink's column names as the row's, as it always has. + ## Design - **One job.** Lint the contract. Not a CDC platform, not a migration runner, diff --git a/corpus/README.md b/corpus/README.md index 41fdbbe..e171a23 100644 --- a/corpus/README.md +++ b/corpus/README.md @@ -31,6 +31,9 @@ The rest exercise one path each: | `include-table-no-schema` | Stack Overflow 74103659: `table.include.list` names `ipaddrs` without its schema, the connector runs and no topic appears | | `include-table-glob` | Stack Overflow 51345636: `public.bg_*` written as a shell glob; as a regular expression it names no table | | `include-table-typo` | a misspelled `table.include.list` entry, which Debezium only logs as a warning (debezium/dbz#872) | +| `flatten-bare-columns` | ClickHouse/clickhouse-kafka-connect discussion 182: a `Flatten` transform on the Debezium connector sends `after.`, the ClickHouse table names its columns as Postgres does, and every change lands as a row of zeros | +| `flatten-envelope-columns` | ClickHouse's own Postgres CDC example (ClickHouse/examples `cdc/postgresql` at ae417ae): `before.` and `after.` columns resolve to the source columns, so the capture rules see `postcode2`, which the include list added for this entry leaves out. The sink connector config there has no `connector.class`; it is added so the file names its sink | +| `flatten-after-unwrap` | a guard: unwrapping first leaves a flat row, so a `Flatten` after it renames nothing | | `diff-adds-column-connector-untouched` | the diff rule: `base/` holds the migrations and connector before the change; the change adds a column and leaves the connector alone | | `diff-connector-captures-one-of-two` | #963 with a hurried fix: two of three new columns go on the include list in the same change; the third is raised, since adding two says nothing about it | | `diff-connector-touched-other-table` | RefuseRadar #962 and #963 in one range: the connector gains report_validations columns, the migration adds reports columns; the reports ones are raised | diff --git a/corpus/flatten-after-unwrap/connector.json b/corpus/flatten-after-unwrap/connector.json new file mode 100644 index 0000000..9b2c930 --- /dev/null +++ b/corpus/flatten-after-unwrap/connector.json @@ -0,0 +1,13 @@ +{ + "name": "shop-source", + "config": { + "connector.class": "io.debezium.connector.postgresql.PostgresConnector", + "plugin.name": "pgoutput", + "topic.prefix": "shop", + "table.include.list": "public.orders", + "transforms": "unwrap,flatten", + "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", + "transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value", + "transforms.flatten.delimiter": "_" + } +} diff --git a/corpus/flatten-after-unwrap/expected.txt b/corpus/flatten-after-unwrap/expected.txt new file mode 100644 index 0000000..cd1f70b --- /dev/null +++ b/corpus/flatten-after-unwrap/expected.txt @@ -0,0 +1 @@ +ok: source, connector and sink agree diff --git a/corpus/flatten-after-unwrap/migrations/0001_orders.sql b/corpus/flatten-after-unwrap/migrations/0001_orders.sql new file mode 100644 index 0000000..f5d419b --- /dev/null +++ b/corpus/flatten-after-unwrap/migrations/0001_orders.sql @@ -0,0 +1,6 @@ +-- A guard: unwrapping first leaves a flat row, so a flatten after it changes +-- no field name and the plain column names are right. +CREATE TABLE orders ( + id BIGSERIAL PRIMARY KEY, + total NUMERIC(12, 2) NOT NULL +); diff --git a/corpus/flatten-after-unwrap/sink-connector.json b/corpus/flatten-after-unwrap/sink-connector.json new file mode 100644 index 0000000..c950d30 --- /dev/null +++ b/corpus/flatten-after-unwrap/sink-connector.json @@ -0,0 +1,6 @@ +{ + "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector", + "database": "shop", + "topics": "shop.public.orders", + "topic2TableMap": "shop.public.orders=orders" +} diff --git a/corpus/flatten-after-unwrap/sink/0001_orders.sql b/corpus/flatten-after-unwrap/sink/0001_orders.sql new file mode 100644 index 0000000..6ca16c4 --- /dev/null +++ b/corpus/flatten-after-unwrap/sink/0001_orders.sql @@ -0,0 +1,7 @@ +CREATE TABLE shop.orders +( + `id` Int64, + `total` String +) +ENGINE = MergeTree +ORDER BY id; diff --git a/corpus/flatten-bare-columns/connector.json b/corpus/flatten-bare-columns/connector.json new file mode 100644 index 0000000..8cb27eb --- /dev/null +++ b/corpus/flatten-bare-columns/connector.json @@ -0,0 +1,21 @@ +{ + "config": { + "connector.class": "io.debezium.connector.postgresql.PostgresConnector", + "plugin.name": "pgoutput", + "database.server.name": "postgres", + "table.include.list": "public.testradiomeas", + "topic.prefix": "testprefix", + "publication.name": "myint_logical_replication", + "value.converter": "org.apache.kafka.connect.json.JsonConverter", + "value.converter.schemas.enable": "false", + "snapshot.mode": "never", + "tombstones.on.delete": "false", + "decimal.handling.mode": "double", + "transforms": "set_topic,flatten", + "transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value", + "transforms.flatten.delimiter": ".", + "transforms.set_topic.type": "org.apache.kafka.connect.transforms.RegexRouter", + "transforms.set_topic.regex": ".*", + "transforms.set_topic.replacement": "testradiomeas" + } +} diff --git a/corpus/flatten-bare-columns/expected.txt b/corpus/flatten-bare-columns/expected.txt new file mode 100644 index 0000000..959c48a --- /dev/null +++ b/corpus/flatten-bare-columns/expected.txt @@ -0,0 +1,5 @@ +error sink-column-flattened flatten-bare-columns/sink/0001_testradiomeas.sql:4 + default.testradiomeas (is fed topic testradiomeas by flatten-bare-columns/sink-connector.json) reads measDate, measLongitude, measLatitude, measNetNoCoverage, measNetType by their names in public.testradiomeas, but transforms.flatten in flatten-bare-columns/connector.json flattens each change event, so every one arrives as after. + every row will carry the columns' defaults, with no error anywhere + fix: rename them after.measDate and so on, or unwrap the event with io.debezium.transforms.ExtractNewRecordState in place of transforms.flatten +1 error(s), 0 warning(s), 0 info diff --git a/corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql b/corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql new file mode 100644 index 0000000..1d44320 --- /dev/null +++ b/corpus/flatten-bare-columns/migrations/0001_testradiomeas.sql @@ -0,0 +1,11 @@ +-- Reduced from ClickHouse/clickhouse-kafka-connect discussion 182 and issue +-- 183 ("ClickHouseSinkConnector & Debezium PostgresConnector creates rows +-- zero values or nulls", 2023-09): every insert, update and delete produced +-- a ClickHouse row of zeros and empty strings, with no error anywhere. +CREATE TABLE testradiomeas ( + "measDate" INTEGER NOT NULL, + "measLongitude" DOUBLE PRECISION, + "measLatitude" DOUBLE PRECISION, + "measNetNoCoverage" INTEGER, + "measNetType" TEXT +); diff --git a/corpus/flatten-bare-columns/sink-connector.json b/corpus/flatten-bare-columns/sink-connector.json new file mode 100644 index 0000000..71501d8 --- /dev/null +++ b/corpus/flatten-bare-columns/sink-connector.json @@ -0,0 +1,9 @@ +{ + "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector", + "tasks.max": "1", + "database": "default", + "exactlyOnce": "false", + "topics": "testradiomeas", + "value.converter": "org.apache.kafka.connect.json.JsonConverter", + "value.converter.schemas.enable": "false" +} diff --git a/corpus/flatten-bare-columns/sink/0001_testradiomeas.sql b/corpus/flatten-bare-columns/sink/0001_testradiomeas.sql new file mode 100644 index 0000000..97e4e9c --- /dev/null +++ b/corpus/flatten-bare-columns/sink/0001_testradiomeas.sql @@ -0,0 +1,13 @@ +-- The table as first written: the columns named as they are in Postgres. +-- The flatten transform names every field after., and the sink +-- matches fields to columns by name, so nothing matched. +CREATE TABLE default.testradiomeas +( + `measDate` Int32, + `measLongitude` Float64, + `measLatitude` Float64, + `measNetNoCoverage` Int32, + `measNetType` String +) +ENGINE = MergeTree +ORDER BY `measDate`; diff --git a/corpus/flatten-envelope-columns/connector.json b/corpus/flatten-envelope-columns/connector.json new file mode 100644 index 0000000..66a906d --- /dev/null +++ b/corpus/flatten-envelope-columns/connector.json @@ -0,0 +1,23 @@ +{ + "connector.class": "PostgresCdcSource", + "name": "uk_price_paid_changes", + "database.server.name": "postgres", + "database.dbname": "postgres", + "publication.name": "dbz_publication", + "publication.autocreate.mode": "filtered", + "table.include.list": "public.uk_price_paid", + "column.include.list": "public\\.uk_price_paid\\.(id|price|postcode1|town)", + "snapshot.mode": "never", + "tombstones.on.delete": "false", + "plugin.name": "pgoutput", + "output.data.format": "JSON", + "after.state.only": "false", + "output.key.format": "STRING", + "tasks.max": "1", + "transforms": "flatten,set_topic", + "transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value", + "transforms.flatten.delimiter": ".", + "transforms.set_topic.type": "io.confluent.connect.cloud.transforms.TopicRegexRouter", + "transforms.set_topic.regex": ".*", + "transforms.set_topic.replacement": "uk_price_paid_changes" +} diff --git a/corpus/flatten-envelope-columns/expected.txt b/corpus/flatten-envelope-columns/expected.txt new file mode 100644 index 0000000..26ecad7 --- /dev/null +++ b/corpus/flatten-envelope-columns/expected.txt @@ -0,0 +1,9 @@ +error sink-column-not-captured flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql:6 + public.uk_price_paid.postcode2 is read by default.uk_price_paid_changes (is fed topic uk_price_paid_changes by flatten-envelope-columns/sink-connector.json) but is not matched by column.include.list in flatten-envelope-columns/connector.json + every row will carry the column's default, with no error anywhere + fix: add public.uk_price_paid.postcode2 to column.include.list, deploy the connector, then apply the sink schema; rows already written need a snapshot +error sink-column-not-captured flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql:11 + public.uk_price_paid.postcode2 is read by default.uk_price_paid_changes (is fed topic uk_price_paid_changes by flatten-envelope-columns/sink-connector.json) but is not matched by column.include.list in flatten-envelope-columns/connector.json + every row will carry the column's default, with no error anywhere + fix: add public.uk_price_paid.postcode2 to column.include.list, deploy the connector, then apply the sink schema; rows already written need a snapshot +2 error(s), 0 warning(s), 0 info diff --git a/corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql b/corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql new file mode 100644 index 0000000..96ae38e --- /dev/null +++ b/corpus/flatten-envelope-columns/migrations/0001_uk_price_paid.sql @@ -0,0 +1,12 @@ +-- Reduced from ClickHouse's own Postgres CDC example (ClickHouse/examples, +-- cdc/postgresql at ae417ae, and the blog post "ClickHouse PostgreSQL change +-- data capture, part 2"): the change events keep Debezium's envelope, a +-- flatten transform names the fields before. and after., and +-- the ClickHouse table names its columns the same way. +CREATE TABLE uk_price_paid ( + id BIGSERIAL PRIMARY KEY, + price INTEGER NOT NULL, + postcode1 TEXT, + postcode2 TEXT, + town TEXT +); diff --git a/corpus/flatten-envelope-columns/sink-connector.json b/corpus/flatten-envelope-columns/sink-connector.json new file mode 100644 index 0000000..5ab233d --- /dev/null +++ b/corpus/flatten-envelope-columns/sink-connector.json @@ -0,0 +1,10 @@ +{ + "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector", + "database": "default", + "exactlyOnce": "false", + "schemas.enable": "false", + "topics": "uk_price_paid_changes", + "value.converter": "org.apache.kafka.connect.json.JsonConverter", + "value.converter.schemas.enable": "false", + "key.converter": "org.apache.kafka.connect.storage.StringConverter" +} diff --git a/corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql b/corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql new file mode 100644 index 0000000..a3c8ab7 --- /dev/null +++ b/corpus/flatten-envelope-columns/sink/0001_uk_price_paid_changes.sql @@ -0,0 +1,17 @@ +CREATE TABLE default.uk_price_paid_changes +( + `before.id` Nullable(UInt64), + `before.price` Nullable(UInt32), + `before.postcode1` Nullable(String), + `before.postcode2` Nullable(String), + `before.town` Nullable(String), + `after.id` Nullable(UInt64), + `after.price` Nullable(UInt32), + `after.postcode1` Nullable(String), + `after.postcode2` Nullable(String), + `after.town` Nullable(String), + `op` LowCardinality(String), + `ts_ms` UInt64 +) +ENGINE = MergeTree +ORDER BY tuple(); diff --git a/internal/capture/debezium/debezium.go b/internal/capture/debezium/debezium.go index 5b6dc31..e18b91e 100644 --- a/internal/capture/debezium/debezium.go +++ b/internal/capture/debezium/debezium.go @@ -19,6 +19,7 @@ import ( "strings" "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/smt" ) // Contract implements model.Contract for a Postgres, MySQL or SQL Server @@ -35,6 +36,7 @@ type Contract struct { columnPatterns []string tablePatterns []string routers []router + shape model.Shape Config map[string]string Class string Name string @@ -126,6 +128,7 @@ func Parse(b []byte) (*Contract, error) { c.tableListKey = "table.include.list" } c.routers = routers(cfg) + c.shape = smt.Apply(cfg, "", smt.Start(cfg)) return c, nil } @@ -195,6 +198,16 @@ func anyMatch(res []*regexp.Regexp, s string) bool { func (c *Contract) Pos() model.Pos { return c.pos } +// ValueShape is what the change events look like after this connector's own +// transforms, for the engine to match sink columns against. +func (c *Contract) ValueShape() model.Shape { + s := c.shape + if s.Kind == model.ShapeFlattened { + s.File = c.pos.File + } + return s +} + func (c *Contract) CapturesTable(schema, table string) bool { q := schema + "." + table if len(c.tableInclude) > 0 { diff --git a/internal/engine/diff.go b/internal/engine/diff.go index bcc90d6..154bb3a 100644 --- a/internal/engine/diff.go +++ b/internal/engine/diff.go @@ -79,8 +79,8 @@ func schemaBeforeConnector(in *Input, reads []Read) ([]model.Finding, map[string } read := map[string]bool{} for _, r := range reads { - if r.Table != nil { - read[strings.ToLower(r.Table.Qualified()+"."+r.Column.Name)] = true + if r.Table != nil && !r.Unreached { + read[strings.ToLower(r.Table.Qualified()+"."+r.Field)] = true } } var fs []model.Finding diff --git a/internal/engine/engine.go b/internal/engine/engine.go index 199eb36..f312090 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -44,6 +44,49 @@ type Read struct { Mapping model.Mapping Topic string Mapper *connect.Mapper + // Field is the source column the sink column receives: its own name, or + // the column under after or before in a flattened envelope. It is + // empty for a column no field reaches (see Unreached). + Field string + // Unreached marks a column named after a source column that receives + // nothing, because the value arrives under after instead. + Unreached bool + Shape model.Shape +} + +// ShapeLister is implemented by contracts that know what their events look +// like after their own transforms. +type ShapeLister interface { + ValueShape() model.Shape +} + +// envelopeFields are the flattened envelope's own fields besides before and +// after (Debezium's change event structure; ts_us and ts_ns since 2.6). +var envelopeFields = map[string]bool{"op": true, "ts_ms": true, "ts_us": true, "ts_ns": true} + +// field works out which source column a sink column receives under shape. +// skip is true for fields the connector or sink adds around the row. +func field(shape model.Shape, table *model.Table, name string) (col string, unreached, skip bool) { + if isMetadata(name) { + return "", false, true + } + if shape.Kind != model.ShapeFlattened { + return name, false, false + } + d := shape.Delimiter + lower := strings.ToLower(name) + for _, image := range []string{"after", "before"} { + if strings.HasPrefix(lower, image+d) { + return name[len(image)+len(d):], false, false + } + } + if envelopeFields[lower] || strings.HasPrefix(lower, "source"+d) || strings.HasPrefix(lower, "transaction"+d) { + return "", false, true + } + if table.Column(name) != nil { + return "", true, false + } + return name, false, false } // isMetadata reports whether a sink column is one the connector or the @@ -64,6 +107,10 @@ func isMetadata(name string) bool { // a topic nothing produces. func resolve(in *Input) []Read { var reads []Read + var source model.Shape + if sl, ok := in.Contract.(ShapeLister); ok { + source = sl.ValueShape() + } byTopic := map[string]*model.Table{} for _, t := range in.Source.Tables { byTopic[in.Contract.Topic(t.Schema, t.Name)] = t @@ -112,11 +159,17 @@ func resolve(in *Input) []Read { } continue } + shape := source + if mapper != nil { + shape = mapper.Reshape(source) + } for _, c := range st.Columns { - if isMetadata(c.Name) { + f, unreached, skip := field(shape, table, c.Name) + if skip { continue } - reads = append(reads, Read{Sink: st, Column: c, Table: table, Mapping: mapping, Topic: topic, Mapper: mapper}) + reads = append(reads, Read{Sink: st, Column: c, Table: table, Mapping: mapping, Topic: topic, Mapper: mapper, + Field: f, Unreached: unreached, Shape: shape}) } } } @@ -131,6 +184,7 @@ func Run(in *Input) []model.Finding { fs = append(fs, sinkTableNotCaptured(in, reads)...) fs = append(fs, sinkColumnNotCaptured(in, reads)...) fs = append(fs, sinkColumnUnknown(in, reads)...) + fs = append(fs, sinkColumnFlattened(reads)...) // The diff rule runs before the inventory so a column it raised as a // warning is not listed again as information. diff, raised := schemaBeforeConnector(in, reads) @@ -216,8 +270,8 @@ func sinkColumnNotCaptured(in *Input, reads []Read) []model.Finding { if r.Table == nil || !in.Contract.CapturesTable(r.Table.Schema, r.Table.Name) { continue } - src := r.Table.Column(r.Column.Name) - if src == nil || in.Contract.CapturesColumn(r.Table.Schema, r.Table.Name, src.Name) { + src := r.Table.Column(r.Field) + if r.Unreached || src == nil || in.Contract.CapturesColumn(r.Table.Schema, r.Table.Name, src.Name) { continue } q := fmt.Sprintf("%s.%s.%s", r.Table.Schema, r.Table.Name, src.Name) @@ -230,10 +284,44 @@ func sinkColumnNotCaptured(in *Input, reads []Read) []model.Finding { return fs } +// sinkColumnFlattened raises a sink table whose columns carry the row's +// names while a flatten transform delivers the envelope: every field arrives +// as aftercolumn, the sink matches by name, and the columns are left at +// their defaults (ClickHouse/clickhouse-kafka-connect discussion 182). One +// finding per table, since one transform causes all of them. +func sinkColumnFlattened(reads []Read) []model.Finding { + var fs []model.Finding + var order []*model.SinkTable + names := map[*model.SinkTable][]string{} + first := map[*model.SinkTable]Read{} + for _, r := range reads { + if !r.Unreached { + continue + } + if _, ok := names[r.Sink]; !ok { + order = append(order, r.Sink) + first[r.Sink] = r + } + names[r.Sink] = append(names[r.Sink], r.Column.Name) + } + for _, st := range order { + r := first[st] + cols := names[st] + fs = append(fs, model.Finding{ + Rule: "sink-column-flattened", Severity: model.Error, Pos: st.Pos, + Message: fmt.Sprintf("%s (%s) reads %s by their names in %s.%s, but %s in %s flattens each change event, so every one arrives as after%s\nevery row will carry the columns' defaults, with no error anywhere", + st.Name, how(r), strings.Join(cols, ", "), r.Table.Schema, r.Table.Name, r.Shape.Transform, r.Shape.File, r.Shape.Delimiter), + Fix: fmt.Sprintf("rename them after%s%s and so on, or unwrap the event with io.debezium.transforms.ExtractNewRecordState in place of %s", + r.Shape.Delimiter, cols[0], r.Shape.Transform), + }) + } + return fs +} + func sinkColumnUnknown(in *Input, reads []Read) []model.Finding { var fs []model.Finding for _, r := range reads { - if r.Table == nil || r.Table.Column(r.Column.Name) != nil { + if r.Table == nil || r.Unreached || r.Table.Column(r.Field) != nil { continue } fs = append(fs, model.Finding{ @@ -248,8 +336,8 @@ func sourceColumnNotCaptured(in *Input, reads []Read, raised map[string]bool) [] var fs []model.Finding declared := map[string]bool{} for _, r := range reads { - if r.Table != nil { - declared[strings.ToLower(r.Table.Qualified()+"."+r.Column.Name)] = true + if r.Table != nil && !r.Unreached { + declared[strings.ToLower(r.Table.Qualified()+"."+r.Field)] = true } } for _, t := range in.Source.Tables { diff --git a/internal/model/model.go b/internal/model/model.go index 3eb70c4..df8cbb3 100644 --- a/internal/model/model.go +++ b/internal/model/model.go @@ -197,3 +197,30 @@ func (f Finding) Format() string { } return b.String() } + +// ShapeKind says what a change event's value looks like when a sink reads it. +type ShapeKind int + +const ( + // ShapeUnknown: a transform cdclint does not model changed the value, or + // the connector is not one whose output it knows. Sink columns are taken + // as the row's column names, as they always were. + ShapeUnknown ShapeKind = iota + // ShapeEnvelope: Debezium's change event, before/after/source/op/ts_ms. + ShapeEnvelope + // ShapeRow: the row itself, unwrapped from the envelope. + ShapeRow + // ShapeFlattened: the envelope flattened into aftercolumn, + // beforecolumn, sourcefield, op, ts_ms. + ShapeFlattened +) + +// Shape is the value a sink receives after every transform on the way. +type Shape struct { + Kind ShapeKind + // Delimiter joins envelope and column names when Kind is ShapeFlattened. + Delimiter string + // Transform and File name the flatten transform, for findings. + Transform string + File string +} diff --git a/internal/sink/connect/connect.go b/internal/sink/connect/connect.go index 0ac38c7..2d2907f 100644 --- a/internal/sink/connect/connect.go +++ b/internal/sink/connect/connect.go @@ -13,6 +13,7 @@ import ( "strings" "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/smt" ) // Kind is which sink connector family a config belongs to. @@ -37,6 +38,7 @@ type Mapper struct { setting string // which setting held the map routes []route // Iceberg: table -> route regex single string // Iceberg: the one table when iceberg.tables lists one + cfg map[string]string // the whole config, for its transforms } type route struct { @@ -44,6 +46,12 @@ type route struct { re *regexp.Regexp } +// Reshape applies this sink connector's transforms to the value the source +// connector produced; a sink connector can unwrap or flatten on its side. +func (m *Mapper) Reshape(in model.Shape) model.Shape { + return smt.Apply(m.cfg, m.Pos.File, in) +} + // ReadFile parses one sink connector JSON file. func ReadFile(path string) (*Mapper, error) { b, err := os.ReadFile(path) @@ -72,7 +80,7 @@ func Parse(b []byte) (*Mapper, error) { for k, v := range src { cfg[k] = fmt.Sprint(v) } - m := &Mapper{Class: cfg["connector.class"], explicit: map[string]string{}} + m := &Mapper{Class: cfg["connector.class"], explicit: map[string]string{}, cfg: cfg} lc := strings.ToLower(m.Class) switch { case strings.Contains(lc, "snowflake"): diff --git a/internal/smt/smt.go b/internal/smt/smt.go new file mode 100644 index 0000000..99b6f8a --- /dev/null +++ b/internal/smt/smt.go @@ -0,0 +1,94 @@ +// Package smt works out what a change event's value looks like after a +// connector's single message transforms, so a sink column can be matched +// to the field that actually arrives. +// +// Only transforms whose effect on field names is documented are modelled: +// Debezium's ExtractNewRecordState (and its old name UnwrapFromEnvelope) and +// the Iceberg DebeziumTransform, which unwrap the envelope into the row, and +// Kafka Connect's Flatten, which joins nested names with a delimiter ("." by +// default). Transforms that only route (RegexRouter, TimestampRouter, +// ByLogicalTableRouter) or drop whole records (Filter) change no field name. +// Anything else that touches the value, or a modelled transform applied +// under a predicate to only some records, makes the shape unknown, and an +// unknown shape produces no finding. +package smt + +import ( + "strings" + + "github.com/avison9/cdclint/internal/model" +) + +// Start is the value a source connector emits before its transforms. It is +// Debezium's envelope for a Debezium connector, and for a connector that +// states it with after.state.only (Confluent's managed Postgres CDC source +// keeps the envelope when it is false and sends only the row when it is +// true). Anything else is unknown. +func Start(cfg map[string]string) model.Shape { + switch strings.ToLower(strings.TrimSpace(cfg["after.state.only"])) { + case "false": + return model.Shape{Kind: model.ShapeEnvelope} + case "true": + return model.Shape{Kind: model.ShapeRow} + } + if strings.HasPrefix(cfg["connector.class"], "io.debezium.") { + return model.Shape{Kind: model.ShapeEnvelope} + } + return model.Shape{} +} + +// Apply walks cfg's transforms in the order transforms= lists them. +func Apply(cfg map[string]string, file string, in model.Shape) model.Shape { + s := in + for _, name := range strings.Split(cfg["transforms"], ",") { + name = strings.TrimSpace(name) + if name == "" || s.Kind == model.ShapeUnknown { + continue + } + key := "transforms." + name + typ := strings.TrimSpace(cfg[key+".type"]) + conditional := cfg[key+".predicate"] != "" + switch { + case isRouter(typ), strings.HasSuffix(typ, "$Key"): + // Routes the record or changes the key; the value is untouched. + case typ == "org.apache.kafka.connect.transforms.Filter": + // Drops whole records; the ones that pass are unchanged. + case isUnwrap(typ) && !conditional: + if s.Kind == model.ShapeEnvelope { + s = model.Shape{Kind: model.ShapeRow} + } else { + s = model.Shape{} + } + case typ == "org.apache.kafka.connect.transforms.Flatten$Value" && !conditional: + if s.Kind == model.ShapeEnvelope { + d := cfg[key+".delimiter"] + if d == "" { + d = "." + } + s = model.Shape{Kind: model.ShapeFlattened, Delimiter: d, Transform: key, File: file} + } + // A row is already flat, and a flattened envelope has nothing + // left to flatten: either way the names do not change. + default: + s = model.Shape{} + } + } + return s +} + +func isRouter(typ string) bool { + return strings.HasSuffix(typ, "RegexRouter") || + typ == "org.apache.kafka.connect.transforms.TimestampRouter" || + typ == "io.debezium.transforms.ByLogicalTableRouter" +} + +func isUnwrap(typ string) bool { + switch typ { + case "io.debezium.transforms.ExtractNewRecordState", + "io.debezium.transforms.UnwrapFromEnvelope", + "io.tabular.iceberg.connect.transforms.DebeziumTransform", + "org.apache.iceberg.connect.transforms.DebeziumTransform": + return true + } + return false +} diff --git a/internal/smt/smt_test.go b/internal/smt/smt_test.go new file mode 100644 index 0000000..792e63c --- /dev/null +++ b/internal/smt/smt_test.go @@ -0,0 +1,85 @@ +package smt + +import ( + "testing" + + "github.com/avison9/cdclint/internal/model" +) + +const ( + debezium = "io.debezium.connector.postgresql.PostgresConnector" + unwrap = "io.debezium.transforms.ExtractNewRecordState" + flatten = "org.apache.kafka.connect.transforms.Flatten$Value" + router = "org.apache.kafka.connect.transforms.RegexRouter" +) + +func TestShapeAfterTransforms(t *testing.T) { + for _, tc := range []struct { + name string + cfg map[string]string + want model.Shape + }{ + {"a Debezium connector with no transforms emits the envelope", + map[string]string{"connector.class": debezium}, + model.Shape{Kind: model.ShapeEnvelope}}, + {"after.state.only false keeps the envelope whatever the class", + map[string]string{"connector.class": "PostgresCdcSource", "after.state.only": "false"}, + model.Shape{Kind: model.ShapeEnvelope}}, + {"after.state.only true sends the row", + map[string]string{"connector.class": "PostgresCdcSource", "after.state.only": "true"}, + model.Shape{Kind: model.ShapeRow}}, + {"a connector cdclint does not know is unknown", + map[string]string{"connector.class": "PostgresCdcSource"}, + model.Shape{}}, + {"unwrap turns the envelope into the row", + map[string]string{"connector.class": debezium, "transforms": "unwrap", "transforms.unwrap.type": unwrap}, + model.Shape{Kind: model.ShapeRow}}, + {"flatten on the envelope uses . by default", + map[string]string{"connector.class": debezium, "transforms": "flat", "transforms.flat.type": flatten}, + model.Shape{Kind: model.ShapeFlattened, Delimiter: ".", Transform: "transforms.flat", File: "c.json"}}, + {"flatten keeps its configured delimiter, and routers around it change nothing", + map[string]string{"connector.class": debezium, "transforms": "route, flat ,route2", + "transforms.route.type": router, "transforms.flat.type": flatten, "transforms.flat.delimiter": "_", + "transforms.route2.type": "io.confluent.connect.cloud.transforms.TopicRegexRouter"}, + model.Shape{Kind: model.ShapeFlattened, Delimiter: "_", Transform: "transforms.flat", File: "c.json"}}, + {"flatten after unwrap leaves the row as it is", + map[string]string{"connector.class": debezium, "transforms": "unwrap,flat", + "transforms.unwrap.type": unwrap, "transforms.flat.type": flatten}, + model.Shape{Kind: model.ShapeRow}}, + {"unwrap after flatten is not modelled", + map[string]string{"connector.class": debezium, "transforms": "flat,unwrap", + "transforms.unwrap.type": unwrap, "transforms.flat.type": flatten}, + model.Shape{}}, + {"a flatten under a predicate applies to some records only, so the shape is unknown", + map[string]string{"connector.class": debezium, "transforms": "flat", + "transforms.flat.type": flatten, "transforms.flat.predicate": "isOrders"}, + model.Shape{}}, + {"a transform that renames value fields makes the shape unknown", + map[string]string{"connector.class": debezium, "transforms": "flat,rename", + "transforms.flat.type": flatten, "transforms.rename.type": "org.apache.kafka.connect.transforms.ReplaceField$Value"}, + model.Shape{}}, + {"key and filter transforms leave the value alone", + map[string]string{"connector.class": debezium, "transforms": "key,drop,flat", + "transforms.key.type": "org.apache.kafka.connect.transforms.ExtractField$Key", + "transforms.drop.type": "org.apache.kafka.connect.transforms.Filter", + "transforms.flat.type": flatten}, + model.Shape{Kind: model.ShapeFlattened, Delimiter: ".", Transform: "transforms.flat", File: "c.json"}}, + } { + t.Run(tc.name, func(t *testing.T) { + if got := Apply(tc.cfg, "c.json", Start(tc.cfg)); got != tc.want { + t.Errorf("got %+v, want %+v", got, tc.want) + } + }) + } +} + +func TestSinkTransformsApplyAfterTheSource(t *testing.T) { + sink := map[string]string{"transforms": "unwrap", "transforms.unwrap.type": unwrap} + if got := Apply(sink, "sink.json", model.Shape{Kind: model.ShapeEnvelope}); got.Kind != model.ShapeRow { + t.Errorf("unwrap on the sink side: got %+v, want the row", got) + } + flat := map[string]string{"transforms": "f", "transforms.f.type": flatten} + if got := Apply(flat, "sink.json", model.Shape{Kind: model.ShapeRow}); got.Kind != model.ShapeRow { + t.Errorf("flatten on the sink side of an unwrapped row: got %+v, want the row", got) + } +} From 291bdb7361dcaa3764667145cbbb23a2f91cfb8b Mon Sep 17 00:00:00 2001 From: avison9 Date: Fri, 25 Sep 2026 22:52:23 +0100 Subject: [PATCH 3/3] MySQL and MariaDB migrations are read, so every rule works on a MySQL source Debezium's MySQL connector is where the most include-list questions come from, and cdclint could only read Postgres. MySQL has no schemas: Debezium names a table database.table and a column database.table.column, so a new reader puts the database where the Postgres reader puts the schema and the rules, include lists and topic names work unchanged. The connector's class picks the reader (MySqlConnector, MariaDbConnector); --migrations mysql:DIR or postgres:DIR says it outright. The database unqualified migrations run in comes from a USE statement, or from a connector whose database.include.list names one database (or whose table.include.list entries all start with the same one); when neither says, cdclint stops and says how to tell it, since guessing would name the wrong table in every finding. database.include.list and database.exclude.list are applied before the table lists, as Debezium does, and sink-table-not-captured names the list that actually keeps a table out, with the entry that belongs in it: "add billing to database.include.list", not billing.invoices. The reader covers CREATE TABLE (backticks, ENGINE options, KEY, INDEX, FULLTEXT and SPATIAL lines, LIKE, keywords glued to their column list), ALTER TABLE (ADD with FIRST or AFTER and several columns in parentheses, CHANGE, MODIFY, DROP, RENAME COLUMN, RENAME TO another database), RENAME TABLE, DROP TABLE and USE. The splitter gains a MySQL mode: # comments, backslash escapes, no dollar quoting, and DELIMITER. MySQL has no ADD COLUMN IF NOT EXISTS, so re-runnable migrations put the DDL in a string run through PREPARE; the reader applies the table statements in those strings, which is the end state the migration guarantees. Proven against a real repository: Mattermost v10.11.0 keeps the same schema as 140 MySQL and 140 Postgres migrations. Both readers find the same 71 tables; of 610 columns, every difference left is on the Postgres reader's side (a glued UNIQUE( read as a column, and columns added inside DO blocks). Before reading the prepared-statement strings the MySQL reader missed about 90 columns there. Down migrations are left out here as in the Postgres reader (#17): golang-migrate *.down.sql files are skipped and goose, sql-migrate and dbmate down sections blanked, through the same internal/migrate calls. Corpus: mysql-column-never-captured, mysql-change-column, mysql-database-not-included, mysql-diff-adds-column. RefuseRadar is unchanged at 0/0/159. --- README.md | 27 +- cmd/cdclint/corpus_test.go | 6 +- cmd/cdclint/main.go | 52 +- cmd/cdclint/source.go | 49 +- corpus/README.md | 4 + corpus/mysql-change-column/connector.json | 12 + corpus/mysql-change-column/expected.txt | 8 + .../migrations/V1__orders.sql | 4 + .../migrations/V2__rename_amount.sql | 4 + .../mysql-change-column/sink/0001_orders.sql | 11 + .../connector.json | 13 + .../mysql-column-never-captured/expected.txt | 5 + .../migrations/V1__orders.sql | 10 + .../migrations/V2__orders_currency.sql | 2 + .../sink/0001_orders.sql | 14 + .../connector.json | 10 + .../mysql-database-not-included/expected.txt | 5 + .../migrations/V1__shop.sql | 4 + .../migrations/V2__billing.sql | 9 + .../sink/0001_invoices.sql | 12 + .../base/connector.json | 11 + .../base/migrations/V1__orders.sql | 4 + corpus/mysql-diff-adds-column/connector.json | 11 + corpus/mysql-diff-adds-column/expected.txt | 5 + .../migrations/V1__orders.sql | 4 + .../migrations/V2__orders_coupon.sql | 3 + .../sink/0001_orders.sql | 11 + internal/capture/debezium/debezium.go | 76 +++ internal/engine/engine.go | 15 +- internal/source/mysql/mysql.go | 554 ++++++++++++++++++ internal/source/mysql/mysql_test.go | 128 ++++ internal/source/postgres/postgres.go | 26 +- internal/source/source.go | 41 ++ internal/sqlsplit/split.go | 48 +- internal/sqlsplit/split_test.go | 27 + 35 files changed, 1171 insertions(+), 54 deletions(-) create mode 100644 corpus/mysql-change-column/connector.json create mode 100644 corpus/mysql-change-column/expected.txt create mode 100644 corpus/mysql-change-column/migrations/V1__orders.sql create mode 100644 corpus/mysql-change-column/migrations/V2__rename_amount.sql create mode 100644 corpus/mysql-change-column/sink/0001_orders.sql create mode 100644 corpus/mysql-column-never-captured/connector.json create mode 100644 corpus/mysql-column-never-captured/expected.txt create mode 100644 corpus/mysql-column-never-captured/migrations/V1__orders.sql create mode 100644 corpus/mysql-column-never-captured/migrations/V2__orders_currency.sql create mode 100644 corpus/mysql-column-never-captured/sink/0001_orders.sql create mode 100644 corpus/mysql-database-not-included/connector.json create mode 100644 corpus/mysql-database-not-included/expected.txt create mode 100644 corpus/mysql-database-not-included/migrations/V1__shop.sql create mode 100644 corpus/mysql-database-not-included/migrations/V2__billing.sql create mode 100644 corpus/mysql-database-not-included/sink/0001_invoices.sql create mode 100644 corpus/mysql-diff-adds-column/base/connector.json create mode 100644 corpus/mysql-diff-adds-column/base/migrations/V1__orders.sql create mode 100644 corpus/mysql-diff-adds-column/connector.json create mode 100644 corpus/mysql-diff-adds-column/expected.txt create mode 100644 corpus/mysql-diff-adds-column/migrations/V1__orders.sql create mode 100644 corpus/mysql-diff-adds-column/migrations/V2__orders_coupon.sql create mode 100644 corpus/mysql-diff-adds-column/sink/0001_orders.sql create mode 100644 internal/source/mysql/mysql.go create mode 100644 internal/source/mysql/mysql_test.go create mode 100644 internal/source/source.go diff --git a/README.md b/README.md index 1376a63..c6df9f0 100644 --- a/README.md +++ b/README.md @@ -12,7 +12,7 @@ GitHub Marketplace. A change-data-capture pipeline has three schemas that must agree and nothing that makes them agree: -1. **The source schema.** Postgres tables, evolved by migrations, changed +1. **The source schema.** Postgres or MySQL tables, evolved by migrations, changed weekly by the application team. 2. **The capture contract.** The Debezium connector JSON: `table.include.list`, `column.include.list`, replica identity, topic naming. Written once by @@ -51,7 +51,7 @@ your repository, and cdclint reads it: | what you see | what is usually wrong | rule | |---|---|---| | A column is `0`, `''`, `NULL` or `1970-01-01` on every row in ClickHouse, BigQuery, Snowflake or Iceberg, with nothing in the logs | the column is not on the connector's `column.include.list` (or is on its exclude list), so Debezium drops it before Kafka | `sink-column-not-captured` | -| A column added in Postgres never shows up downstream | the migration added it and nobody added it to the include list | `schema-before-connector` (with `--base`), `sink-column-not-captured` | +| A column added in Postgres or MySQL never shows up downstream | the migration added it and nobody added it to the include list | `schema-before-connector` (with `--base`), `sink-column-not-captured` | | Listing the columns of one table made another table's columns disappear | `column.include.list` is one list for every captured table | `sink-column-not-captured` | | "`table.include.list` not working", or the connector is `RUNNING` and no topic appears | the entry has no schema (`orders` for `public.orders`), or is a glob (`public.bg_*`) where Debezium expects a regex (`public\.bg_.*`); Debezium matches each entry against the whole `schema.table` name, never a substring | `captured-table-missing` names the entry and the fix; `sink-table-not-captured` when a sink reads the table | | A Kafka-engine table or sink receives nothing, or the wrong table fills | the topic it reads is not the one the connector produces (prefix, schema, `RegexRouter`) | `topic-table-mapping` | @@ -60,7 +60,7 @@ your repository, and cdclint reads it: | The Kafka Connect sink writes a row of `0` and empty strings for every change, and no error ("sink replicates zero values or nulls") | a `Flatten` transform keeps Debezium's envelope and names every field `after.`, while the sink table names its columns as Postgres does; sinks match fields to columns by name | `sink-column-flattened` | | A ClickHouse materialized view writes defaults | a refreshable view's `SELECT` order differs from the target's (it matches by position), or a streaming view's names differ (it matches by name) | `mv-column-match` | -It reads files only: Postgres sources today, no type checks, no connection to +It reads files only: Postgres and MySQL (or MariaDB) sources, no type checks, no connection to anything running. ## What it checks @@ -202,6 +202,23 @@ line; `fail-on`, `version` and `working-directory` are the other inputs. The action downloads the release binary for the runner and verifies its checksum before running it. +## Sources + +The migrations are read in filename order, and the connector's class picks +the dialect: `MySqlConnector` and `MariaDbConnector` read MySQL, every other +Debezium source reads Postgres. `--migrations mysql:DIR` or `postgres:DIR` +says it outright. + +MySQL has no schemas: Debezium names a table `database.table`, so cdclint +needs the database the migrations run in. It takes it from a `USE db;` in the +migrations, or from the connector when `database.include.list` names one +database (or every `table.include.list` entry starts with the same one), and +stops with that advice when neither says. `database.include.list` and +`database.exclude.list` are applied before the table lists, as Debezium does, +and a finding names whichever list keeps a table out. MySQL's re-runnable +migrations wrap `ALTER TABLE` in a string run through `PREPARE`, because MySQL +has no `ADD COLUMN IF NOT EXISTS`; cdclint reads the DDL in those strings too. + ## Sinks, out of the box in v1 | sink | how tables are found | how a table maps to a topic | @@ -214,7 +231,7 @@ checksum before running it. Without a sink connector config, a sink table maps to the source table of the same name, and the finding says so. Debezium's `RegexRouter` transform is -applied when computing topic names. Other sources (MySQL, SQL Server) and +applied when computing topic names. Other sources (SQL Server, Oracle) and sinks are packages behind the same two interfaces. Transforms on either connector decide what the sink receives, and cdclint @@ -232,7 +249,7 @@ then reads the sink's column names as the row's, as it always has. - **One job.** Lint the contract. Not a CDC platform, not a migration runner, not monitoring. -- **Pluggable ends.** Postgres and Debezium as the source and capture +- **Pluggable ends.** Postgres or MySQL, and Debezium, as the source and capture readers; ClickHouse, BigQuery, Snowflake and Iceberg as sink readers, each behind a small interface so the next one is a package, not a rewrite. - **The corpus is the spec.** [`corpus/`](corpus/) holds one directory per diff --git a/cmd/cdclint/corpus_test.go b/cmd/cdclint/corpus_test.go index 47f1448..c2e76d9 100644 --- a/cmd/cdclint/corpus_test.go +++ b/cmd/cdclint/corpus_test.go @@ -9,7 +9,7 @@ import ( "testing" "github.com/avison9/cdclint/internal/engine" - "github.com/avison9/cdclint/internal/source/postgres" + "github.com/avison9/cdclint/internal/source" ) // update rewrites every expected.txt from the current output. Run it after @@ -110,13 +110,13 @@ func baseFromDir(t *testing.T, dir, headConnector string) *engine.Base { if err != nil { t.Fatal(err) } - var files []postgres.NamedFile + var files []source.NamedFile for _, e := range entries { b, err := os.ReadFile(filepath.Join(dir, "migrations", e.Name())) if err != nil { t.Fatal(err) } - files = append(files, postgres.NamedFile{Path: filepath.Join(dir, "migrations", e.Name()), Text: string(b)}) + files = append(files, source.NamedFile{Path: filepath.Join(dir, "migrations", e.Name()), Text: string(b)}) } sort.Slice(files, func(i, j int) bool { return files[i].Path < files[j].Path }) // No connector.json under base/ means the connector is new in the diff --git a/cmd/cdclint/main.go b/cmd/cdclint/main.go index 0dbabd5..b2ee952 100644 --- a/cmd/cdclint/main.go +++ b/cmd/cdclint/main.go @@ -24,7 +24,7 @@ import ( "github.com/avison9/cdclint/internal/sink/clickhouse" "github.com/avison9/cdclint/internal/sink/connect" "github.com/avison9/cdclint/internal/sink/sqlddl" - "github.com/avison9/cdclint/internal/source/postgres" + "github.com/avison9/cdclint/internal/source" ) // version is set by the release build with -ldflags "-X main.version=...". @@ -42,7 +42,7 @@ func main() { func run(args []string) int { fs := flag.NewFlagSet("cdclint", flag.ContinueOnError) var ( - migrations = fs.String("migrations", "", "directory of source migrations, applied in name order (Postgres)") + migrations = fs.String("migrations", "", "directory of source migrations, applied in name order; Postgres or MySQL, from the connector class, or say it with postgres:DIR or mysql:DIR") connector = fs.String("connector", "", "Debezium source connector JSON") sinks multi sinkConns multi @@ -144,11 +144,15 @@ func load(migrations, connector string, sinks, sinkConns []string) (*engine.Inpu // Load is exported for the corpus test, which runs the tool exactly as the // command line does. func Load(migrations, connector string, sinks, sinkConns []string) (*engine.Input, error) { - src, _, err := readSource(migrations) + c, err := debezium.ReadFile(connector) if err != nil { return nil, err } - c, err := debezium.ReadFile(connector) + reader, dir, err := pickSource(migrations, c) + if err != nil { + return nil, err + } + src, _, err := reader.readDir(dir) if err != nil { return nil, err } @@ -198,55 +202,73 @@ func LoadBase(ref, migrations, connector string) (*engine.Base, error) { if err != nil { return nil, fmt.Errorf("--base: %w", err) } + // The dialect and the database come from the connector as it is now: + // the base's may not exist, and the source did not change dialect. + head, err := debezium.ReadFile(connector) + if err != nil { + return nil, err + } + reader, migrations, err := pickSource(migrations, head) + if err != nil { + return nil, err + } files, err := gitread.Dir(ref, migrations) if err != nil { return nil, fmt.Errorf("--base: %w", err) } - var named []postgres.NamedFile + var named []source.NamedFile for _, f := range files { - named = append(named, postgres.NamedFile{Path: f.Path, Text: f.Text}) + named = append(named, source.NamedFile{Path: f.Path, Text: f.Text}) } exists, err := gitread.Exists(ref, connector) if err != nil { return nil, fmt.Errorf("--base: %w", err) } if !exists { - return baseFrom(id, named, nil, true) + return baseFrom(id, reader, named, nil, true) } changed, err := gitread.Changed(ref, connector) if err != nil { return nil, fmt.Errorf("--base: %w", err) } if !changed { - return baseFrom(id, named, nil, false) + return baseFrom(id, reader, named, nil, false) } text, err := gitread.Show(ref, connector) if err != nil { return nil, fmt.Errorf("--base: %w", err) } - return baseFrom(id, named, []byte(text), true) + return baseFrom(id, reader, named, []byte(text), true) } // BaseFromFiles builds the base from migrations already in memory and the // connector's text at the base and now; the corpus test feeds it from a // base/ directory, so the diff rule is tested without a repository. A nil // baseConnector is a connector that did not exist at the base. -func BaseFromFiles(ref string, migrations []postgres.NamedFile, baseConnector, headConnector []byte) (*engine.Base, error) { +func BaseFromFiles(ref string, migrations []source.NamedFile, baseConnector, headConnector []byte) (*engine.Base, error) { + head, err := debezium.Parse(headConnector) + if err != nil { + return nil, err + } + reader, _, err := pickSource("", head) + if err != nil { + return nil, err + } if baseConnector == nil { - return baseFrom(ref, migrations, nil, true) + return baseFrom(ref, reader, migrations, nil, true) } if string(baseConnector) == string(headConnector) { - return baseFrom(ref, migrations, nil, false) + return baseFrom(ref, reader, migrations, nil, false) } - return baseFrom(ref, migrations, baseConnector, true) + return baseFrom(ref, reader, migrations, baseConnector, true) } // baseFrom parses the base's migrations and, when given, its connector. A // parse failure there is an error, because the base once ran and its // config once parsed, so the failure is in the reader and hiding it would // hide the rule. -func baseFrom(ref string, migrations []postgres.NamedFile, connector []byte, changed bool) (*engine.Base, error) { - src, err := postgres.ReadFiles(migrations) +func baseFrom(ref string, reader sourceReader, migrations []source.NamedFile, connector []byte, changed bool) (*engine.Base, error) { + src, err := reader.readFiles(migrations) if err != nil { return nil, fmt.Errorf("--base %s: %w", ref, err) } diff --git a/cmd/cdclint/source.go b/cmd/cdclint/source.go index 4b63229..d879555 100644 --- a/cmd/cdclint/source.go +++ b/cmd/cdclint/source.go @@ -1,12 +1,55 @@ package main import ( + "fmt" + "strings" + + "github.com/avison9/cdclint/internal/capture/debezium" "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/source" + "github.com/avison9/cdclint/internal/source/mysql" "github.com/avison9/cdclint/internal/source/postgres" ) -// readSource is the one place the source dialect is chosen. Postgres is the -// only one in v1; a MySQL reader slots in here. -func readSource(dir string) (*model.Source, []string, error) { +// sourceReader is the one place the source dialect is chosen. A +// "postgres:" or "mysql:" prefix on --migrations says it outright; +// otherwise the connector's class decides, since a MySqlConnector (or +// MariaDbConnector) reads a MySQL database and every other Debezium source +// cdclint knows reads Postgres. +type sourceReader struct { + mysql bool + database string // MySQL: the database unqualified migrations run in +} + +func pickSource(migrations string, c *debezium.Contract) (sourceReader, string, error) { + r := sourceReader{mysql: c.MySQL()} + dir := migrations + if i := strings.Index(migrations, ":"); i > 0 && !strings.Contains(migrations[:i], "/") { + switch strings.ToLower(migrations[:i]) { + case "mysql", "mariadb": + r.mysql, dir = true, migrations[i+1:] + case "postgres", "postgresql": + r.mysql, dir = false, migrations[i+1:] + default: + return r, "", fmt.Errorf("unknown source dialect %q in --migrations %q", migrations[:i], migrations) + } + } + if r.mysql { + r.database = c.DefaultDatabase() + } + return r, dir, nil +} + +func (r sourceReader) readDir(dir string) (*model.Source, []string, error) { + if r.mysql { + return mysql.ReadDir(dir, r.database) + } return postgres.ReadDir(dir) } + +func (r sourceReader) readFiles(files []source.NamedFile) (*model.Source, error) { + if r.mysql { + return mysql.ReadFiles(files, r.database) + } + return postgres.ReadFiles(files) +} diff --git a/corpus/README.md b/corpus/README.md index e171a23..06a890a 100644 --- a/corpus/README.md +++ b/corpus/README.md @@ -34,6 +34,10 @@ The rest exercise one path each: | `flatten-bare-columns` | ClickHouse/clickhouse-kafka-connect discussion 182: a `Flatten` transform on the Debezium connector sends `after.`, the ClickHouse table names its columns as Postgres does, and every change lands as a row of zeros | | `flatten-envelope-columns` | ClickHouse's own Postgres CDC example (ClickHouse/examples `cdc/postgresql` at ae417ae): `before.` and `after.` columns resolve to the source columns, so the capture rules see `postcode2`, which the include list added for this entry leaves out. The sink connector config there has no `connector.class`; it is added so the file names its sink | | `flatten-after-unwrap` | a guard: unwrapping first leaves a flat row, so a `Flatten` after it renames nothing | +| `mysql-column-never-captured` | the column-never-captured incident in MySQL syntax: backticks, `KEY` lines, `ADD COLUMN ... AFTER`, a `database.include.list` naming the database the migrations run in | +| `mysql-change-column` | MySQL renames a column with `CHANGE`; the include list still names the old column and the new one stops arriving | +| `mysql-database-not-included` | a table in a second database (`USE billing`) that `database.include.list` leaves out; the fix names the database list and the database | +| `mysql-diff-adds-column` | the diff rule through the MySQL reader: `base/` without the new column, the change adds it and leaves the connector alone | | `diff-adds-column-connector-untouched` | the diff rule: `base/` holds the migrations and connector before the change; the change adds a column and leaves the connector alone | | `diff-connector-captures-one-of-two` | #963 with a hurried fix: two of three new columns go on the include list in the same change; the third is raised, since adding two says nothing about it | | `diff-connector-touched-other-table` | RefuseRadar #962 and #963 in one range: the connector gains report_validations columns, the migration adds reports columns; the reports ones are raised | diff --git a/corpus/mysql-change-column/connector.json b/corpus/mysql-change-column/connector.json new file mode 100644 index 0000000..6e9140f --- /dev/null +++ b/corpus/mysql-change-column/connector.json @@ -0,0 +1,12 @@ +{ + "name": "shop-mysql-source", + "config": { + "connector.class": "io.debezium.connector.mysql.MySqlConnector", + "topic.prefix": "shopdb", + "database.include.list": "shop", + "table.include.list": "shop.orders", + "column.include.list": "shop\\.orders\\.(id|amount)", + "transforms": "unwrap", + "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState" + } +} diff --git a/corpus/mysql-change-column/expected.txt b/corpus/mysql-change-column/expected.txt new file mode 100644 index 0000000..b100789 --- /dev/null +++ b/corpus/mysql-change-column/expected.txt @@ -0,0 +1,8 @@ +warning captured-column-missing mysql-change-column/connector.json + shop\.orders\.amount in column.include.list matches no column in the source schema + fix: remove it, or check the spelling against the migrations +error sink-column-not-captured mysql-change-column/sink/0001_orders.sql:4 + shop.orders.total_amount is read by kafka_orders (reads topic shopdb.shop.orders) but is not matched by column.include.list in mysql-change-column/connector.json + every row will carry the column's default, with no error anywhere + fix: add shop.orders.total_amount to column.include.list, deploy the connector, then apply the sink schema; rows already written need a snapshot +1 error(s), 1 warning(s), 0 info diff --git a/corpus/mysql-change-column/migrations/V1__orders.sql b/corpus/mysql-change-column/migrations/V1__orders.sql new file mode 100644 index 0000000..49013b9 --- /dev/null +++ b/corpus/mysql-change-column/migrations/V1__orders.sql @@ -0,0 +1,4 @@ +CREATE TABLE orders ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + amount DECIMAL(12,2) NOT NULL +) ENGINE=InnoDB; diff --git a/corpus/mysql-change-column/migrations/V2__rename_amount.sql b/corpus/mysql-change-column/migrations/V2__rename_amount.sql new file mode 100644 index 0000000..be6e1d5 --- /dev/null +++ b/corpus/mysql-change-column/migrations/V2__rename_amount.sql @@ -0,0 +1,4 @@ +-- MySQL renames a column with CHANGE, which also restates its type. The +-- include list still names amount, which no longer exists, and the new +-- name is on no list, so the renamed column stops arriving. +ALTER TABLE orders CHANGE COLUMN amount total_amount DECIMAL(14,2) NOT NULL; diff --git a/corpus/mysql-change-column/sink/0001_orders.sql b/corpus/mysql-change-column/sink/0001_orders.sql new file mode 100644 index 0000000..a6d2b1e --- /dev/null +++ b/corpus/mysql-change-column/sink/0001_orders.sql @@ -0,0 +1,11 @@ +CREATE TABLE kafka_orders +( + `id` UInt64, + `total_amount` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shopdb.shop.orders', + kafka_group_name = 'clickhouse-orders', + kafka_format = 'JSONEachRow'; diff --git a/corpus/mysql-column-never-captured/connector.json b/corpus/mysql-column-never-captured/connector.json new file mode 100644 index 0000000..4a8a125 --- /dev/null +++ b/corpus/mysql-column-never-captured/connector.json @@ -0,0 +1,13 @@ +{ + "name": "shop-mysql-source", + "config": { + "connector.class": "io.debezium.connector.mysql.MySqlConnector", + "database.server.id": "184054", + "topic.prefix": "shopdb", + "database.include.list": "shop", + "table.include.list": "shop.orders", + "column.include.list": "shop\\.orders\\.(id|customer_id|total|created_at)", + "transforms": "unwrap", + "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState" + } +} diff --git a/corpus/mysql-column-never-captured/expected.txt b/corpus/mysql-column-never-captured/expected.txt new file mode 100644 index 0000000..01fae05 --- /dev/null +++ b/corpus/mysql-column-never-captured/expected.txt @@ -0,0 +1,5 @@ +error sink-column-not-captured mysql-column-never-captured/sink/0001_orders.sql:6 + shop.orders.currency is read by kafka_orders (reads topic shopdb.shop.orders) but is not matched by column.include.list in mysql-column-never-captured/connector.json + every row will carry the column's default, with no error anywhere + fix: add shop.orders.currency to column.include.list, deploy the connector, then apply the sink schema; rows already written need a snapshot +1 error(s), 0 warning(s), 0 info diff --git a/corpus/mysql-column-never-captured/migrations/V1__orders.sql b/corpus/mysql-column-never-captured/migrations/V1__orders.sql new file mode 100644 index 0000000..48a72a9 --- /dev/null +++ b/corpus/mysql-column-never-captured/migrations/V1__orders.sql @@ -0,0 +1,10 @@ +# The Postgres column-never-captured incident, in MySQL: a column added by a +# later migration and never put on the connector's column.include.list. +CREATE TABLE `orders` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `customer_id` BIGINT UNSIGNED NOT NULL, + `total` DECIMAL(12,2) NOT NULL, + `created_at` DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6), + PRIMARY KEY (`id`), + KEY `idx_orders_customer` (`customer_id`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; diff --git a/corpus/mysql-column-never-captured/migrations/V2__orders_currency.sql b/corpus/mysql-column-never-captured/migrations/V2__orders_currency.sql new file mode 100644 index 0000000..5189b43 --- /dev/null +++ b/corpus/mysql-column-never-captured/migrations/V2__orders_currency.sql @@ -0,0 +1,2 @@ +ALTER TABLE `orders` + ADD COLUMN `currency` CHAR(3) NOT NULL DEFAULT 'EUR' AFTER `total`; diff --git a/corpus/mysql-column-never-captured/sink/0001_orders.sql b/corpus/mysql-column-never-captured/sink/0001_orders.sql new file mode 100644 index 0000000..1da26e2 --- /dev/null +++ b/corpus/mysql-column-never-captured/sink/0001_orders.sql @@ -0,0 +1,14 @@ +CREATE TABLE kafka_orders +( + `id` UInt64, + `customer_id` UInt64, + `total` String, + `currency` String, + `created_at` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shopdb.shop.orders', + kafka_group_name = 'clickhouse-orders', + kafka_format = 'JSONEachRow'; diff --git a/corpus/mysql-database-not-included/connector.json b/corpus/mysql-database-not-included/connector.json new file mode 100644 index 0000000..3176604 --- /dev/null +++ b/corpus/mysql-database-not-included/connector.json @@ -0,0 +1,10 @@ +{ + "name": "shop-mysql-source", + "config": { + "connector.class": "io.debezium.connector.mysql.MySqlConnector", + "topic.prefix": "shopdb", + "database.include.list": "shop", + "transforms": "unwrap", + "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState" + } +} diff --git a/corpus/mysql-database-not-included/expected.txt b/corpus/mysql-database-not-included/expected.txt new file mode 100644 index 0000000..cd0e99b --- /dev/null +++ b/corpus/mysql-database-not-included/expected.txt @@ -0,0 +1,5 @@ +error sink-table-not-captured mysql-database-not-included/sink/0001_invoices.sql:1 + kafka_invoices reads topic shopdb.billing.invoices, but billing is not matched by database.include.list in mysql-database-not-included/connector.json + nothing will ever arrive on that topic + fix: add billing to database.include.list and deploy the connector before the sink schema +1 error(s), 0 warning(s), 0 info diff --git a/corpus/mysql-database-not-included/migrations/V1__shop.sql b/corpus/mysql-database-not-included/migrations/V1__shop.sql new file mode 100644 index 0000000..eedb0bb --- /dev/null +++ b/corpus/mysql-database-not-included/migrations/V1__shop.sql @@ -0,0 +1,4 @@ +CREATE TABLE orders ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + total DECIMAL(12,2) NOT NULL +); diff --git a/corpus/mysql-database-not-included/migrations/V2__billing.sql b/corpus/mysql-database-not-included/migrations/V2__billing.sql new file mode 100644 index 0000000..760df29 --- /dev/null +++ b/corpus/mysql-database-not-included/migrations/V2__billing.sql @@ -0,0 +1,9 @@ +-- A second database on the same server, created by the same migrations. +-- The connector's database.include.list names only shop, so nothing from +-- billing is captured, whatever the table lists say. +USE billing; +CREATE TABLE invoices ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + order_id BIGINT UNSIGNED NOT NULL, + issued_at DATETIME NOT NULL +); diff --git a/corpus/mysql-database-not-included/sink/0001_invoices.sql b/corpus/mysql-database-not-included/sink/0001_invoices.sql new file mode 100644 index 0000000..45966e5 --- /dev/null +++ b/corpus/mysql-database-not-included/sink/0001_invoices.sql @@ -0,0 +1,12 @@ +CREATE TABLE kafka_invoices +( + `id` UInt64, + `order_id` UInt64, + `issued_at` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shopdb.billing.invoices', + kafka_group_name = 'clickhouse-invoices', + kafka_format = 'JSONEachRow'; diff --git a/corpus/mysql-diff-adds-column/base/connector.json b/corpus/mysql-diff-adds-column/base/connector.json new file mode 100644 index 0000000..aa0018d --- /dev/null +++ b/corpus/mysql-diff-adds-column/base/connector.json @@ -0,0 +1,11 @@ +{ + "name": "shop-mysql-source", + "config": { + "connector.class": "io.debezium.connector.mysql.MySqlConnector", + "topic.prefix": "shopdb", + "database.include.list": "shop", + "column.include.list": "shop\\.orders\\.(id|total)", + "transforms": "unwrap", + "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState" + } +} diff --git a/corpus/mysql-diff-adds-column/base/migrations/V1__orders.sql b/corpus/mysql-diff-adds-column/base/migrations/V1__orders.sql new file mode 100644 index 0000000..eedb0bb --- /dev/null +++ b/corpus/mysql-diff-adds-column/base/migrations/V1__orders.sql @@ -0,0 +1,4 @@ +CREATE TABLE orders ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + total DECIMAL(12,2) NOT NULL +); diff --git a/corpus/mysql-diff-adds-column/connector.json b/corpus/mysql-diff-adds-column/connector.json new file mode 100644 index 0000000..aa0018d --- /dev/null +++ b/corpus/mysql-diff-adds-column/connector.json @@ -0,0 +1,11 @@ +{ + "name": "shop-mysql-source", + "config": { + "connector.class": "io.debezium.connector.mysql.MySqlConnector", + "topic.prefix": "shopdb", + "database.include.list": "shop", + "column.include.list": "shop\\.orders\\.(id|total)", + "transforms": "unwrap", + "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState" + } +} diff --git a/corpus/mysql-diff-adds-column/expected.txt b/corpus/mysql-diff-adds-column/expected.txt new file mode 100644 index 0000000..2be666f --- /dev/null +++ b/corpus/mysql-diff-adds-column/expected.txt @@ -0,0 +1,5 @@ +warning schema-before-connector mysql-diff-adds-column/migrations/V2__orders_coupon.sql:3 + this change adds shop.orders.coupon_code to a captured table without adding it to column.include.list in mysql-diff-adds-column/connector.json (compared with base) + the column will not be in the stream; if a sink is later given it, every row will be the default until a snapshot + fix: add shop.orders.coupon_code to column.include.list in the same change, or leave it off on purpose and let this warning stand as the record of that (it blocks only under --fail-on warning) +0 error(s), 1 warning(s), 0 info diff --git a/corpus/mysql-diff-adds-column/migrations/V1__orders.sql b/corpus/mysql-diff-adds-column/migrations/V1__orders.sql new file mode 100644 index 0000000..eedb0bb --- /dev/null +++ b/corpus/mysql-diff-adds-column/migrations/V1__orders.sql @@ -0,0 +1,4 @@ +CREATE TABLE orders ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + total DECIMAL(12,2) NOT NULL +); diff --git a/corpus/mysql-diff-adds-column/migrations/V2__orders_coupon.sql b/corpus/mysql-diff-adds-column/migrations/V2__orders_coupon.sql new file mode 100644 index 0000000..8565e39 --- /dev/null +++ b/corpus/mysql-diff-adds-column/migrations/V2__orders_coupon.sql @@ -0,0 +1,3 @@ +# The change under review adds a column to a captured table and leaves the +# connector alone. +ALTER TABLE orders ADD COLUMN coupon_code VARCHAR(32) NULL; diff --git a/corpus/mysql-diff-adds-column/sink/0001_orders.sql b/corpus/mysql-diff-adds-column/sink/0001_orders.sql new file mode 100644 index 0000000..722b557 --- /dev/null +++ b/corpus/mysql-diff-adds-column/sink/0001_orders.sql @@ -0,0 +1,11 @@ +CREATE TABLE kafka_orders +( + `id` UInt64, + `total` String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = '${KAFKA_BROKERS}', + kafka_topic_list = 'shopdb.shop.orders', + kafka_group_name = 'clickhouse-orders', + kafka_format = 'JSONEachRow'; diff --git a/internal/capture/debezium/debezium.go b/internal/capture/debezium/debezium.go index e18b91e..ab9e0b8 100644 --- a/internal/capture/debezium/debezium.go +++ b/internal/capture/debezium/debezium.go @@ -29,6 +29,10 @@ type Contract struct { prefix string tableInclude []*regexp.Regexp tableExclude []*regexp.Regexp + dbInclude []*regexp.Regexp // MySQL and MariaDB: database.include.list + dbExclude []*regexp.Regexp + dbIncludeKey string + dbExcludeKey string columnInclude []*regexp.Regexp columnExclude []*regexp.Regexp columnListKey string @@ -127,6 +131,16 @@ func Parse(b []byte) (*Contract, error) { if c.tableListKey == "" { c.tableListKey = "table.include.list" } + if c.MySQL() { + // MySQL has no schemas: Debezium filters databases first, and a + // table must pass both lists. + if c.dbInclude, c.dbIncludeKey, err = list(cfg, "database.include.list", "database.whitelist"); err != nil { + return nil, err + } + if c.dbExclude, c.dbExcludeKey, err = list(cfg, "database.exclude.list", "database.blacklist"); err != nil { + return nil, err + } + } c.routers = routers(cfg) c.shape = smt.Apply(cfg, "", smt.Start(cfg)) return c, nil @@ -208,7 +222,69 @@ func (c *Contract) ValueShape() model.Shape { return s } +// MySQL reports whether this is Debezium's MySQL or MariaDB connector, whose +// tables are databaseName.tableName and whose migrations are MySQL's. +func (c *Contract) MySQL() bool { + lc := strings.ToLower(c.Class) + return strings.Contains(lc, "mysql") || strings.Contains(lc, "mariadb") +} + +// DefaultDatabase is the database unqualified MySQL migrations run in, when +// the connector names exactly one: a database.include.list holding a single +// plain name, or failing that a table.include.list whose entries all start +// with the same plain database name. It is "" when the connector does not +// say, and the reader then needs a USE statement. +func (c *Contract) DefaultDatabase() string { + if dbs := strings.Split(first(c.Config, "database.include.list", "database.whitelist"), ","); len(dbs) == 1 { + if name, ok := literal(dbs[0]); ok { + return name + } + } + common := "" + for _, p := range strings.Split(first(c.Config, "table.include.list", "table.whitelist"), ",") { + p = strings.TrimSpace(p) + dot := strings.Index(p, `\.`) + if dot < 0 { + dot = strings.Index(p, ".") + } + if dot <= 0 { + return "" + } + name, ok := literal(p[:dot]) + if !ok || common != "" && !strings.EqualFold(common, name) { + return "" + } + common = name + } + return common +} + +// literal reports whether a list entry is a plain name rather than a pattern. +func literal(p string) (string, bool) { + p = strings.TrimSpace(p) + if p == "" || strings.ContainsAny(p, `.*+?()[]{}|^$\`) { + return "", false + } + return p, true +} + +// TableBlock names the list that keeps schema.table out of the stream and +// the entry that would let it in: the database for a MySQL database list, +// the qualified table otherwise. +func (c *Contract) TableBlock(schema, table string) (setting, entry string) { + if len(c.dbInclude) > 0 && !anyMatch(c.dbInclude, schema) { + return c.dbIncludeKey, schema + } + if len(c.dbExclude) > 0 && anyMatch(c.dbExclude, schema) { + return c.dbExcludeKey, schema + } + return c.tableListKey, schema + "." + table +} + func (c *Contract) CapturesTable(schema, table string) bool { + if len(c.dbInclude) > 0 && !anyMatch(c.dbInclude, schema) || len(c.dbExclude) > 0 && anyMatch(c.dbExclude, schema) { + return false + } q := schema + "." + table if len(c.tableInclude) > 0 { return anyMatch(c.tableInclude, q) diff --git a/internal/engine/engine.go b/internal/engine/engine.go index f312090..c89acd5 100644 --- a/internal/engine/engine.go +++ b/internal/engine/engine.go @@ -29,6 +29,13 @@ type Input struct { Base *Base } +// TableBlocker is implemented by contracts with more than one list that can +// keep a table out (MySQL's database lists besides the table list), so a +// finding names the list to edit and what to put in it. +type TableBlocker interface { + TableBlock(schema, table string) (setting, entry string) +} + // PatternLister is implemented by contracts whose column and table lists are // sets of patterns that can be checked one by one. type PatternLister interface { @@ -255,10 +262,14 @@ func sinkTableNotCaptured(in *Input, reads []Read) []model.Finding { continue } seen[r.Sink] = true + setting, entry := in.Contract.TableListSetting(), r.Table.Schema+"."+r.Table.Name + if tb, ok := in.Contract.(TableBlocker); ok { + setting, entry = tb.TableBlock(r.Table.Schema, r.Table.Name) + } fs = append(fs, model.Finding{ Rule: "sink-table-not-captured", Severity: model.Error, Pos: r.Sink.Pos, - Message: fmt.Sprintf("%s %s, but %s.%s %s in %s\nnothing will ever arrive on that topic", r.Sink.Name, how(r), r.Table.Schema, r.Table.Name, leftOut(in.Contract.TableListSetting()), in.Contract.Pos().File), - Fix: edit(in.Contract.TableListSetting(), r.Table.Schema+"."+r.Table.Name) + " and deploy the connector before the sink schema", + Message: fmt.Sprintf("%s %s, but %s %s in %s\nnothing will ever arrive on that topic", r.Sink.Name, how(r), entry, leftOut(setting), in.Contract.Pos().File), + Fix: edit(setting, entry) + " and deploy the connector before the sink schema", }) } return fs diff --git a/internal/source/mysql/mysql.go b/internal/source/mysql/mysql.go new file mode 100644 index 0000000..afdcd2a --- /dev/null +++ b/internal/source/mysql/mysql.go @@ -0,0 +1,554 @@ +// Package mysql reads a directory of MySQL or MariaDB migrations, applies +// them in filename order, and returns the schema they leave behind. Like the +// Postgres reader it reads only the statements that shape a table's columns +// (CREATE, ALTER, RENAME and DROP TABLE, and USE) and ignores indexes, +// routines, triggers, grants and data. +// +// MySQL has no schemas. Debezium names a table databaseName.tableName and a +// column databaseName.tableName.columnName, so the database goes where the +// Postgres reader puts the schema, and the rules, include lists and topic +// names work unchanged. Replica identity is a Postgres setting and is left +// empty here. +package mysql + +import ( + "fmt" + "strings" + + "github.com/avison9/cdclint/internal/ddl" + "github.com/avison9/cdclint/internal/migrate" + "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/source" + "github.com/avison9/cdclint/internal/sqlsplit" +) + +// ReadDir applies every *.sql file in dir, sorted by name. database is the +// database an unqualified table belongs to until a USE statement names +// another; see ReadFiles. +func ReadDir(dir, database string) (*model.Source, []string, error) { + named, files, err := source.ReadDir(dir) + if err != nil { + return nil, nil, err + } + src, err := ReadFiles(named, database) + return src, files, err +} + +// ReadFiles applies migrations already in memory, in the order given, and +// leaves out down migrations as the Postgres reader does (package migrate). A +// migration runner connects to one database and runs unqualified statements +// in it, and that name is in neither the files nor the statements; database +// supplies it (the caller takes it from the connector). A table the reader +// cannot place in a database is an error, because every finding about it +// would name the wrong table. +func ReadFiles(files []source.NamedFile, database string) (*model.Source, error) { + r := &reader{src: &model.Source{}, db: database} + for _, f := range files { + if migrate.Down(f.Path) { + continue + } + if err := r.apply(f.Path, f.Text); err != nil { + return nil, fmt.Errorf("%s: %w", f.Path, err) + } + } + return r.src, nil +} + +type reader struct { + src *model.Source + db string +} + +func (r *reader) apply(file, text string) error { + for _, st := range sqlsplit.SplitMySQL(migrate.Up(text)) { + if err := r.statement(file, st); err != nil { + return err + } + } + return nil +} + +func (r *reader) statement(file string, st sqlsplit.Statement) error { + body := doubleEscapedQuotes(st.Text) + w := ddl.Words(body) + // MySQL has no ADD COLUMN IF NOT EXISTS, so re-runnable migrations put + // the DDL in a string and run it through a prepared statement: + // SET @s = (SELECT IF(, 'SELECT 1', 'ALTER TABLE t ADD c INT')); + // PREPARE stmt FROM @s; EXECUTE stmt; + // The IF only skips the DDL when its effect is already there, and adding + // a column that exists or dropping one that does not is a no-op here, so + // applying every table statement found in such a string gives the end + // state the migration guarantees. 100 of Mattermost's 140 MySQL + // migrations are written this way. + if (ddl.HasPrefixFold(w, "SET") && len(w) > 1 && strings.HasPrefix(w[1], "@")) || ddl.HasPrefixFold(w, "PREPARE") { + for _, lit := range stringLiterals(st) { + if err := r.statement(file, lit); err != nil { + return err + } + } + return nil + } + pos := model.Pos{File: file, Line: st.Line} + var err error + switch { + case ddl.HasPrefixFold(w, "USE") && len(w) >= 2: + r.db = ddl.Unquote(w[1]) + case ddl.HasPrefixFold(w, "CREATE", "TEMPORARY"): + // Session-scoped and never written to the binlog as rows. + case ddl.HasPrefixFold(w, "CREATE", "TABLE"): + err = r.createTable(body, w, pos) + case ddl.HasPrefixFold(w, "ALTER") && tableKeyword(w) > 0: + err = r.alterTable(w, tableKeyword(w), pos) + case ddl.HasPrefixFold(w, "RENAME", "TABLE"): + err = r.renameTables(w[2:]) + case ddl.HasPrefixFold(w, "DROP", "TABLE"), ddl.HasPrefixFold(w, "DROP", "TEMPORARY", "TABLE"): + err = r.dropTables(w) + } + if err != nil { + return fmt.Errorf("line %d: %w", st.Line, err) + } + return nil +} + +// stringLiterals returns the single-quoted literals in st that hold a table +// statement, as statements of their own, each on the line it starts on. +func stringLiterals(st sqlsplit.Statement) []sqlsplit.Statement { + var out []sqlsplit.Statement + s := st.Text + line := st.Line + for i := 0; i < len(s); i++ { + switch s[i] { + case '\n': + line++ + case '`', '"': + if end := strings.IndexByte(s[i+1:], s[i]); end >= 0 { + line += strings.Count(s[i:i+1+end], "\n") + i += end + 1 + } + case '\'': + var lit strings.Builder + start := line + j := i + 1 + for ; j < len(s); j++ { + c := s[j] + if c == '\\' && j+1 < len(s) { + lit.WriteByte(s[j+1]) + j++ + continue + } + if c == '\'' { + if j+1 < len(s) && s[j+1] == '\'' { + lit.WriteByte('\'') + j++ + continue + } + break + } + if c == '\n' { + line++ + } + lit.WriteByte(c) + } + i = j + text := strings.TrimSuffix(strings.TrimSpace(lit.String()), ";") + lw := ddl.Words(text) + if ddl.HasPrefixFold(lw, "ALTER") || ddl.HasPrefixFold(lw, "CREATE", "TABLE") || + ddl.HasPrefixFold(lw, "DROP", "TABLE") || ddl.HasPrefixFold(lw, "RENAME", "TABLE") { + out = append(out, sqlsplit.Statement{Text: text, Line: start}) + } + } + } + return out +} + +// tableKeyword finds TABLE in ALTER [ONLINE] [IGNORE] TABLE and returns its +// index, or 0 when the statement alters something else. +func tableKeyword(w []string) int { + for i := 1; i < len(w) && i <= 3; i++ { + switch strings.ToUpper(w[i]) { + case "TABLE": + return i + case "ONLINE", "OFFLINE", "IGNORE": + default: + return 0 + } + } + return 0 +} + +// qualify splits db.table, placing an unqualified table in the current +// database. +func (r *reader) qualify(name string) (db, table string, err error) { + parts := ddl.SplitName(name) + if len(parts) >= 2 { + return parts[len(parts)-2], parts[len(parts)-1], nil + } + if r.db == "" { + return "", "", fmt.Errorf("table %s names no database and the migrations do not say which one they run in: "+ + "give the connector a database.include.list naming that one database, or start the first migration with USE ;", parts[0]) + } + return r.db, parts[0], nil +} + +func (r *reader) createTable(text string, w []string, pos model.Pos) error { + i := 2 + if ddl.HasPrefixFold(w[i:], "IF", "NOT", "EXISTS") { + i += 3 + } + if i >= len(w) { + return nil + } + // The name may be glued to the paren: "reports(" is legal. + name := w[i] + if p := strings.IndexByte(name, '('); p > 0 { + name = name[:p] + } + db, table, err := r.qualify(name) + if err != nil { + return err + } + if r.src.Table(db, table) != nil { + // CREATE TABLE IF NOT EXISTS for a table that exists is a no-op; a + // plain CREATE would have failed. Either way keep the first. + return nil + } + t := &model.Table{Schema: db, Name: table, Pos: pos} + rest := w[i+1:] + if strings.HasPrefix(w[i], name+"(") { + rest = append([]string{w[i][len(name):]}, rest...) + } + // CREATE TABLE t LIKE s, or (LIKE s): the same columns as s. + if ddl.HasPrefixFold(rest, "LIKE") && len(rest) >= 2 { + return r.createLike(t, rest[1]) + } + if len(rest) > 0 && strings.HasPrefix(rest[0], "(") { + inner, _, _ := ddl.Body(rest[0]) + if iw := ddl.Words(inner); ddl.HasPrefixFold(iw, "LIKE") && len(iw) >= 2 { + return r.createLike(t, iw[1]) + } + } + // The column list is the first parenthesis after TABLE. Searching for + // the name's word instead fails when the name is glued to a paren that + // opens a multi-line list ("orders(\n id ..."): the word splitter has + // folded that list's whitespace, so the word is not in the text. + body, _, ok := ddl.Body(text[strings.Index(strings.ToUpper(text), "TABLE")+len("TABLE"):]) + if !ok { + // CREATE TABLE ... AS SELECT without a column list: its columns are + // the query's, which a file reader cannot know. Leaving the table + // out is better than inventing its columns. + return nil + } + cursor := 0 + for _, item := range ddl.SplitTop(body) { + line := pos.Line + ddl.ItemLine(text, item, &cursor) - 1 + iw := ddl.Words(item) + if len(iw) == 0 { + continue + } + if isKey(iw) { + if pk := primaryKey(iw, item); pk != nil { + t.PrimaryKey = pk + } + continue + } + col := model.Column{Name: ddl.Unquote(iw[0]), Type: strings.Join(typeWords(iw[1:]), " "), Pos: model.Pos{File: pos.File, Line: line}} + if containsFold(iw, "PRIMARY") { + t.PrimaryKey = []string{col.Name} + } + t.Columns = append(t.Columns, col) + } + r.src.Tables = append(r.src.Tables, t) + return nil +} + +func (r *reader) createLike(t *model.Table, like string) error { + db, table, err := r.qualify(like) + if err != nil { + return err + } + from := r.src.Table(db, table) + if from == nil { + return nil + } + t.Columns = append([]model.Column(nil), from.Columns...) + t.PrimaryKey = append([]string(nil), from.PrimaryKey...) + r.src.Tables = append(r.src.Tables, t) + return nil +} + +// isKey reports whether a table element defines an index or a constraint +// rather than a column. SPATIAL is not reserved, so a column may be named +// spatial; it only starts an index when INDEX, KEY or a column list follows. +func isKey(iw []string) bool { + // A keyword may be glued to its column list: UNIQUE(a, b). + first := iw[0] + if p := strings.IndexByte(first, '('); p > 0 { + first = first[:p] + } + switch strings.ToUpper(first) { + case "PRIMARY", "KEY", "INDEX", "UNIQUE", "FULLTEXT", "CONSTRAINT", "FOREIGN", "CHECK": + return true + case "SPATIAL": + return len(iw) > 1 && (strings.EqualFold(iw[1], "INDEX") || strings.EqualFold(iw[1], "KEY") || strings.HasPrefix(iw[1], "(")) + } + return false +} + +// primaryKey returns the key's columns when the element is a primary key. +// Key parts may carry a prefix length or an order: name(10) DESC. +func primaryKey(iw []string, item string) []string { + if !containsFold(iw, "PRIMARY") { + return nil + } + body, _, ok := ddl.Body(item) + if !ok { + return nil + } + var cols []string + for _, c := range ddl.SplitTop(body) { + name := ddl.Words(c)[0] + if p := strings.IndexByte(name, '('); p > 0 { + name = name[:p] + } + cols = append(cols, ddl.Unquote(name)) + } + return cols +} + +// typeWords keeps the type part of a column definition: everything up to the +// first attribute keyword. UNSIGNED and ZEROFILL belong to the type. +func typeWords(w []string) []string { + for i, x := range w { + switch strings.ToUpper(x) { + case "NOT", "NULL", "DEFAULT", "AUTO_INCREMENT", "COMMENT", "PRIMARY", "UNIQUE", "KEY", "REFERENCES", + "CHECK", "GENERATED", "AS", "COLLATE", "CHARACTER", "CHARSET", "ON", "INVISIBLE", "VISIBLE", + "STORED", "VIRTUAL", "SRID", "FIRST", "AFTER", "COLUMN_FORMAT", "STORAGE", "CONSTRAINT", + "ENGINE_ATTRIBUTE", "SECONDARY_ENGINE_ATTRIBUTE": + return w[:i] + } + } + return w +} + +func (r *reader) alterTable(w []string, at int, pos model.Pos) error { + if at+1 >= len(w) { + return nil + } + db, table, err := r.qualify(w[at+1]) + if err != nil { + return err + } + t := r.src.Table(db, table) + if t == nil { + return nil + } + // Actions are comma-separated after the name; each starts with a verb. + for _, a := range ddl.SplitTop(strings.Join(w[at+2:], " ")) { + aw := ddl.Words(a) + if len(aw) < 2 { + continue + } + verb := strings.ToUpper(aw[0]) + j := 1 + if strings.EqualFold(aw[1], "COLUMN") && verb != "RENAME" { + j = 2 + } + if j >= len(aw) { + continue + } + switch verb { + case "ADD": + if j == 1 && isKey(aw[1:]) || strings.EqualFold(aw[1], "PARTITION") || strings.EqualFold(aw[1], "PERIOD") { + continue + } + if ddl.HasPrefixFold(aw[j:], "IF", "NOT", "EXISTS") { + j += 3 + } + if j >= len(aw) { + continue + } + if strings.HasPrefix(aw[j], "(") { + // ADD [COLUMN] (a INT, b INT): several columns at once. + inner, _, _ := ddl.Body(aw[j]) + for _, def := range ddl.SplitTop(inner) { + if dw := ddl.Words(def); len(dw) > 0 && !isKey(dw) { + addColumn(t, dw, pos) + } + } + continue + } + addColumn(t, aw[j:], pos) + case "DROP": + if j == 1 && (isKey(aw[1:]) || strings.EqualFold(aw[1], "PARTITION") || strings.EqualFold(aw[1], "DEFAULT")) { + continue + } + if ddl.HasPrefixFold(aw[j:], "IF", "EXISTS") { + j += 2 + } + if j < len(aw) { + removeColumn(t, ddl.Unquote(aw[j])) + } + case "CHANGE": + // CHANGE [COLUMN] old new definition [FIRST | AFTER col] + if j+1 < len(aw) { + if c := t.Column(ddl.Unquote(aw[j])); c != nil { + c.Name = ddl.Unquote(aw[j+1]) + c.Type = strings.Join(typeWords(aw[j+2:]), " ") + c.Pos = pos + } + } + case "MODIFY": + if c := t.Column(ddl.Unquote(aw[j])); c != nil { + c.Type = strings.Join(typeWords(aw[j+1:]), " ") + } + case "RENAME": + switch { + case strings.EqualFold(aw[1], "COLUMN") && len(aw) >= 5 && strings.EqualFold(aw[3], "TO"): + if c := t.Column(ddl.Unquote(aw[2])); c != nil { + c.Name = ddl.Unquote(aw[4]) + c.Pos = pos + } + case strings.EqualFold(aw[1], "INDEX"), strings.EqualFold(aw[1], "KEY"): + default: + k := 1 + if strings.EqualFold(aw[1], "TO") || strings.EqualFold(aw[1], "AS") { + k = 2 + } + if k < len(aw) { + if err := r.renameTable(t, aw[k]); err != nil { + return err + } + } + } + } + } + return nil +} + +// addColumn adds one column definition, honouring FIRST and AFTER. +func addColumn(t *model.Table, dw []string, pos model.Pos) { + name := ddl.Unquote(dw[0]) + if t.Column(name) != nil { + return + } + col := model.Column{Name: name, Type: strings.Join(typeWords(dw[1:]), " "), Pos: pos} + at := len(t.Columns) + for k := 1; k < len(dw); k++ { + switch { + case strings.EqualFold(dw[k], "FIRST"): + at = 0 + case strings.EqualFold(dw[k], "AFTER") && k+1 < len(dw): + for i, c := range t.Columns { + if strings.EqualFold(c.Name, ddl.Unquote(dw[k+1])) { + at = i + 1 + } + } + } + } + t.Columns = append(t.Columns, model.Column{}) + copy(t.Columns[at+1:], t.Columns[at:]) + t.Columns[at] = col +} + +func removeColumn(t *model.Table, name string) { + for i := range t.Columns { + if strings.EqualFold(t.Columns[i].Name, name) { + t.Columns = append(t.Columns[:i], t.Columns[i+1:]...) + return + } + } +} + +// renameTables handles RENAME TABLE a TO b [, c TO d]. +func (r *reader) renameTables(w []string) error { + for _, pair := range ddl.SplitTop(strings.Join(w, " ")) { + pw := ddl.Words(pair) + if len(pw) < 3 || !strings.EqualFold(pw[1], "TO") { + continue + } + db, table, err := r.qualify(pw[0]) + if err != nil { + return err + } + if t := r.src.Table(db, table); t != nil { + if err := r.renameTable(t, pw[2]); err != nil { + return err + } + } + } + return nil +} + +// renameTable renames t, possibly moving it to another database. +func (r *reader) renameTable(t *model.Table, to string) error { + db, table, err := r.qualify(to) + if err != nil { + return err + } + t.Schema, t.Name = db, table + return nil +} + +func (r *reader) dropTables(w []string) error { + i := 2 + if strings.EqualFold(w[1], "TEMPORARY") { + return nil + } + if ddl.HasPrefixFold(w[i:], "IF", "EXISTS") { + i += 2 + } + for _, name := range ddl.SplitTop(strings.Join(w[i:], " ")) { + nw := ddl.Words(name) + if len(nw) == 0 || strings.EqualFold(nw[0], "RESTRICT") || strings.EqualFold(nw[0], "CASCADE") { + continue + } + db, table, err := r.qualify(nw[0]) + if err != nil { + return err + } + for j, t := range r.src.Tables { + if strings.EqualFold(t.Schema, db) && strings.EqualFold(t.Name, table) { + r.src.Tables = append(r.src.Tables[:j], r.src.Tables[j+1:]...) + break + } + } + } + return nil +} + +func containsFold(w []string, kw string) bool { + for _, x := range w { + if strings.EqualFold(x, kw) { + return true + } + } + return false +} + +// doubleEscapedQuotes rewrites \' and \" inside quoted strings as ” and "", +// which mean the same and which the shared word splitter understands. Both +// forms are two characters, so every offset and line stays where it was. +func doubleEscapedQuotes(s string) string { + b := []byte(s) + var q byte + for i := 0; i < len(b); i++ { + c := b[i] + switch { + case q == 0 && c == '`': + // An identifier: quotes inside it are just characters. + if end := strings.IndexByte(s[i+1:], '`'); end >= 0 { + i += end + 1 + } + case q == 0 && (c == '\'' || c == '"'): + q = c + case q != 0 && c == '\\' && i+1 < len(b): + if b[i+1] == q { + b[i] = q + } + i++ + case q != 0 && c == q: + q = 0 + } + } + return string(b) +} diff --git a/internal/source/mysql/mysql_test.go b/internal/source/mysql/mysql_test.go new file mode 100644 index 0000000..90b44f7 --- /dev/null +++ b/internal/source/mysql/mysql_test.go @@ -0,0 +1,128 @@ +package mysql + +import ( + "strings" + "testing" + + "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/source" +) + +func columns(t *model.Table) string { + var names []string + for _, c := range t.Columns { + names = append(names, c.Name) + } + return strings.Join(names, ",") +} + +func TestReadsTheSchemaMigrationsLeaveBehind(t *testing.T) { + files := []source.NamedFile{ + {Path: "V1__init.sql", Text: "# Flyway style\n" + + "CREATE TABLE `orders` (\n" + + " `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,\n" + + " `note` VARCHAR(40) DEFAULT 'it\\'s, fine' COMMENT 'a, b',\n" + + " spatial POINT NOT NULL SRID 4326,\n" + + " `status` ENUM('new','paid') NOT NULL DEFAULT 'new',\n" + + " PRIMARY KEY (`id`),\n" + + " KEY `idx_status` (`status`),\n" + + " SPATIAL INDEX (spatial),\n" + + " CONSTRAINT `fk` FOREIGN KEY (`id`) REFERENCES other (`id`)\n" + + ") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;\n" + + "CREATE TABLE IF NOT EXISTS orders (ignored INT);\n" + + "CREATE TEMPORARY TABLE scratch (x INT);\n"}, + {Path: "V2__alter.sql", Text: "ALTER TABLE orders\n" + + " ADD COLUMN currency CHAR(3) NOT NULL DEFAULT 'EUR' AFTER note,\n" + + " ADD (paid_at DATETIME(6), refunded TINYINT(1)),\n" + + " ADD INDEX idx_paid (paid_at),\n" + + " DROP INDEX idx_status,\n" + + " CHANGE COLUMN note memo TEXT,\n" + + " MODIFY status VARCHAR(10);\n" + + "ALTER TABLE orders ADD first_col INT FIRST, DROP COLUMN spatial;\n" + + "ALTER TABLE orders RENAME COLUMN refunded TO is_refunded;\n" + + "CREATE TABLE orders_archive LIKE orders;\n" + + "RENAME TABLE orders_archive TO archive.orders_2025;\n" + + "USE reporting;\n" + + "CREATE TABLE daily (day DATE PRIMARY KEY, total DECIMAL(12,2));\n" + + "DROP TABLE IF EXISTS shop.gone, daily;\n"}, + } + src, err := ReadFiles(files, "shop") + if err != nil { + t.Fatal(err) + } + orders := src.Table("shop", "orders") + if orders == nil { + t.Fatal("shop.orders missing") + } + if got, want := columns(orders), "first_col,id,memo,currency,status,paid_at,is_refunded"; got != want { + t.Errorf("shop.orders columns = %s, want %s", got, want) + } + if got := orders.Column("id").Type; got != "BIGINT UNSIGNED" { + t.Errorf("id type = %q", got) + } + if got := orders.Column("status").Type; got != "VARCHAR(10)" { + t.Errorf("status type after MODIFY = %q", got) + } + if got := strings.Join(orders.PrimaryKey, ","); got != "id" { + t.Errorf("primary key = %q", got) + } + if orders.ReplicaIdentity != "" { + t.Errorf("replica identity = %q, want empty for MySQL", orders.ReplicaIdentity) + } + // LIKE copied orders as it stood then, and RENAME TABLE moved it to + // another database. + moved := src.Table("archive", "orders_2025") + if moved == nil { + t.Fatal("archive.orders_2025 missing") + } + if got, want := columns(moved), "first_col,id,memo,currency,status,paid_at,is_refunded"; got != want { + t.Errorf("archive.orders_2025 columns = %s, want %s", got, want) + } + if src.Table("shop", "orders_archive") != nil || src.Table("shop", "scratch") != nil { + t.Error("renamed or temporary table still present") + } + if src.Table("reporting", "daily") != nil { + t.Error("reporting.daily was dropped by DROP TABLE after USE reporting") + } +} + +func TestANameGluedToAMultiLineColumnList(t *testing.T) { + // Mattermost's 000016_create_reactions: the paren opens a list that + // runs over several lines. + src, err := ReadFiles([]source.NamedFile{{Path: "V1.sql", + Text: "CREATE TABLE IF NOT EXISTS Reactions(\n UserId varchar(26) NOT NULL,\n PostId varchar(26) NOT NULL\n);"}}, "mm") + if err != nil { + t.Fatal(err) + } + if r := src.Table("mm", "Reactions"); r == nil || columns(r) != "UserId,PostId" { + t.Fatalf("mm.Reactions = %+v", r) + } +} + +func TestDownMigrationsAreLeftOut(t *testing.T) { + // Mattermost v10.11.0's 000092 down file drops a column of another + // table; golang-migrate never runs it forward, and neither does this. + files := []source.NamedFile{ + {Path: "000016_create_reactions.up.sql", Text: "CREATE TABLE Reactions (UserId varchar(26), CreateAt bigint);"}, + {Path: "000092_add_createat_to_teammembers.down.sql", Text: "ALTER TABLE Reactions DROP COLUMN CreateAt;"}, + {Path: "000093_goose.sql", Text: "-- +goose Up\nALTER TABLE Reactions ADD COLUMN EmojiName varchar(64);\n-- +goose Down\nDROP TABLE Reactions;\n"}, + } + src, err := ReadFiles(files, "mm") + if err != nil { + t.Fatal(err) + } + if r := src.Table("mm", "Reactions"); r == nil || columns(r) != "UserId,CreateAt,EmojiName" { + t.Fatalf("mm.Reactions = %+v", r) + } +} + +func TestAnUnqualifiedTableWithNoDatabaseIsAnError(t *testing.T) { + _, err := ReadFiles([]source.NamedFile{{Path: "V1.sql", Text: "CREATE TABLE t (id INT);"}}, "") + if err == nil || !strings.Contains(err.Error(), "database.include.list") { + t.Fatalf("err = %v, want one naming database.include.list", err) + } + src, err := ReadFiles([]source.NamedFile{{Path: "V1.sql", Text: "USE app;\nCREATE TABLE t (id INT);"}}, "") + if err != nil || src.Table("app", "t") == nil { + t.Fatalf("USE should place the table: %v %v", src, err) + } +} diff --git a/internal/source/postgres/postgres.go b/internal/source/postgres/postgres.go index 3978ebc..683e357 100644 --- a/internal/source/postgres/postgres.go +++ b/internal/source/postgres/postgres.go @@ -10,14 +10,12 @@ package postgres import ( "fmt" - "os" - "path/filepath" - "sort" "strings" "github.com/avison9/cdclint/internal/ddl" "github.com/avison9/cdclint/internal/migrate" "github.com/avison9/cdclint/internal/model" + "github.com/avison9/cdclint/internal/source" "github.com/avison9/cdclint/internal/sqlsplit" ) @@ -27,34 +25,16 @@ const DefaultSchema = "public" // ReadDir applies every *.sql file in dir, sorted by name. func ReadDir(dir string) (*model.Source, []string, error) { - entries, err := os.ReadDir(dir) + named, files, err := source.ReadDir(dir) if err != nil { return nil, nil, err } - var files []string - for _, e := range entries { - if !e.IsDir() && strings.HasSuffix(strings.ToLower(e.Name()), ".sql") { - files = append(files, filepath.Join(dir, e.Name())) - } - } - sort.Strings(files) - var named []NamedFile - for _, f := range files { - b, err := os.ReadFile(f) - if err != nil { - return nil, nil, err - } - named = append(named, NamedFile{Path: f, Text: string(b)}) - } src, err := ReadFiles(named) return src, files, err } // NamedFile is one migration's content, named by the path findings show. -type NamedFile struct { - Path string - Text string -} +type NamedFile = source.NamedFile // ReadFiles applies migrations already in memory, in the order given. // Down migrations are skipped: a forward migrate never runs them. diff --git a/internal/source/source.go b/internal/source/source.go new file mode 100644 index 0000000..f78c381 --- /dev/null +++ b/internal/source/source.go @@ -0,0 +1,41 @@ +// Package source holds what every migrations reader shares: a migration as +// a named text, and the directory listing every migration runner uses, all +// *.sql files in filename order. +package source + +import ( + "os" + "path/filepath" + "sort" + "strings" +) + +// NamedFile is one migration's content, named by the path findings show. +type NamedFile struct { + Path string + Text string +} + +// ReadDir returns every *.sql file in dir, sorted by name, with its text. +func ReadDir(dir string) ([]NamedFile, []string, error) { + entries, err := os.ReadDir(dir) + if err != nil { + return nil, nil, err + } + var files []string + for _, e := range entries { + if !e.IsDir() && strings.HasSuffix(strings.ToLower(e.Name()), ".sql") { + files = append(files, filepath.Join(dir, e.Name())) + } + } + sort.Strings(files) + var named []NamedFile + for _, f := range files { + b, err := os.ReadFile(f) + if err != nil { + return nil, nil, err + } + named = append(named, NamedFile{Path: f, Text: string(b)}) + } + return named, files, nil +} diff --git a/internal/sqlsplit/split.go b/internal/sqlsplit/split.go index 4eccdde..4948c7d 100644 --- a/internal/sqlsplit/split.go +++ b/internal/sqlsplit/split.go @@ -17,6 +17,19 @@ type Statement struct { // Split returns the statements in src. Comments are removed from the text so // readers never see them, but the line count they occupied is preserved. func Split(src string) []Statement { + return split(src, false) +} + +// SplitMySQL is Split for MySQL and MariaDB: # starts a line comment, a +// backslash escapes the next character inside '...' and "...", there is no +// dollar quoting ($ is an identifier character), and a DELIMITER line at +// the start of a line changes what ends a statement, the way the mysql +// client reads a file of stored routines. +func SplitMySQL(src string) []Statement { + return split(src, true) +} + +func split(src string, mysql bool) []Statement { var ( out []Statement buf strings.Builder @@ -25,6 +38,8 @@ func Split(src string) []Statement { i = 0 n = len(src) dollarT string // the $tag$ that opened a dollar quote, "" when not inside one + delim = ";" + bol = true // at the start of a line, ignoring leading blanks ) flush := func() { text := strings.TrimSpace(buf.String()) @@ -42,6 +57,21 @@ func Split(src string) []Statement { } for i < n { c := src[i] + if mysql && bol && c != ' ' && c != '\t' { + bol = false + if len(src)-i > 10 && strings.EqualFold(src[i:i+9], "DELIMITER") && (src[i+9] == ' ' || src[i+9] == '\t') { + end := strings.IndexByte(src[i:], '\n') + if end < 0 { + end = n - i + } + if d := strings.TrimSpace(src[i+10 : i+end]); d != "" { + flush() + delim = d + } + i += end + continue + } + } switch { case dollarT != "": // Inside $tag$ ... $tag$: copy through until the closing tag. @@ -56,6 +86,10 @@ func Split(src string) []Statement { } emit(string(c)) i++ + case mysql && c == '#': + for i < n && src[i] != '\n' { + i++ + } case c == '-' && i+1 < n && src[i+1] == '-': // Line comment: drop to end of line, keep the newline. for i < n && src[i] != '\n' { @@ -80,6 +114,13 @@ func Split(src string) []Statement { if src[j] == '\n' { line++ } + if mysql && q != '`' && src[j] == '\\' && j+1 < n { + if src[j+1] == '\n' { + line++ + } + j += 2 + continue + } if src[j] == q { if j+1 < n && src[j+1] == q { j += 2 @@ -94,7 +135,7 @@ func Split(src string) []Statement { } emit(src[i : j+1]) i = j + 1 - case c == '$': + case !mysql && c == '$': // $$ or $tag$ opens a dollar quote; a lone $ is just a character. if tag, ok := dollarTag(src[i:]); ok { dollarT = tag @@ -104,12 +145,13 @@ func Split(src string) []Statement { } emit("$") i++ - case c == ';': + case strings.HasPrefix(src[i:], delim): flush() - i++ + i += len(delim) default: if c == '\n' { line++ + bol = true } emit(string(c)) i++ diff --git a/internal/sqlsplit/split_test.go b/internal/sqlsplit/split_test.go index 4bb951d..e9ef195 100644 --- a/internal/sqlsplit/split_test.go +++ b/internal/sqlsplit/split_test.go @@ -32,3 +32,30 @@ ALTER TABLE a ADD COLUMN "x;y" INT; } } } + +func TestSplitMySQLHidesItsOwnSemicolons(t *testing.T) { + src := "# a hash comment; with a semicolon\n" + + "CREATE TABLE `a` (\n" + + " id INT, # trailing; comment\n" + + " note VARCHAR(20) DEFAULT 'it\\'s; fine'\n" + + ") ENGINE=InnoDB;\n" + + "/*!40101 SET NAMES utf8mb4 */;\n" + + "DELIMITER $$\n" + + "CREATE TRIGGER a_bi BEFORE INSERT ON a FOR EACH ROW BEGIN SET NEW.id = 1; END$$\n" + + "DELIMITER ;\n" + + "ALTER TABLE a ADD COLUMN price$ INT;\n" + got := SplitMySQL(src) + want := []Statement{ + {Text: "CREATE TABLE `a` (\n id INT, \n note VARCHAR(20) DEFAULT 'it\\'s; fine'\n) ENGINE=InnoDB", Line: 2}, + {Text: "CREATE TRIGGER a_bi BEFORE INSERT ON a FOR EACH ROW BEGIN SET NEW.id = 1; END", Line: 8}, + {Text: "ALTER TABLE a ADD COLUMN price$ INT", Line: 10}, + } + if len(got) != len(want) { + t.Fatalf("got %d statements, want %d: %#v", len(got), len(want), got) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("statement %d = %#v, want %#v", i, got[i], want[i]) + } + } +}