diff --git a/README.md b/README.md index 6ed9833..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,14 +51,16 @@ 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 | `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` | +| 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 @@ -70,7 +72,9 @@ 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 | +| `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 | @@ -198,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 | @@ -210,14 +231,25 @@ 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 +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, 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 c2dddf5..06a890a 100644 --- a/corpus/README.md +++ b/corpus/README.md @@ -28,6 +28,16 @@ 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) | +| `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/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/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/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 2ef7f03..ab9e0b8 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 @@ -28,12 +29,18 @@ 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 tableListKey string columnPatterns []string + tablePatterns []string routers []router + shape model.Shape Config map[string]string Class string Name string @@ -95,6 +102,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 } @@ -119,7 +131,18 @@ 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 } @@ -189,7 +212,79 @@ 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 +} + +// 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) @@ -232,4 +327,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/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 8370e92..c89acd5 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,18 @@ 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. +// 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 { ColumnPatterns() []string + TablePatterns() []string } // Read is one sink column resolved to the source column it expects. @@ -41,6 +51,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 @@ -61,6 +114,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 @@ -109,11 +166,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}) } } } @@ -128,12 +191,14 @@ 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) 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 @@ -197,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 @@ -212,8 +281,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) @@ -226,10 +295,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{ @@ -244,8 +347,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 { @@ -291,6 +394,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. 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) + } +} 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]) + } + } +}