diff --git a/CHANGELOG.md b/CHANGELOG.md index 9f7604b..8049290 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **Common values of declared columns.** The workload bundle now has `distributions.json`. It holds the most common values from `pg_stats` for each column under `database.workload.distributions`. Values ship verbatim unless `hash_values: true`. The capture role needs `SELECT` on each declared column. See [Recording the common values of chosen columns](docs/workload_capture.md#recording-the-common-values-of-chosen-columns). +### Changed + +- **Query log capture settings.** `capture_log_type` is required when `capture_log` is on, and is `statement` or `pgaudit`. `capture_log_source` is now `file` or `log_fdw`, and defaults to `file`. The tool reads whether an exported file is plain text, csvlog or a pgAudit JSON export. The old `capture_log_source` values `pgaudit`, `pgaudit-json`, `stderr` and `auto` are rejected, and the error names the settings that replace them. See [Capturing transaction shapes and values](docs/workload_capture.md#capturing-transaction-shapes-and-values). + ## [2.1.0] - 2026-10-05 ### Added diff --git a/README.md b/README.md index 9eac082..a6c1550 100644 --- a/README.md +++ b/README.md @@ -301,9 +301,12 @@ more accurate sharding plan. | 2 | Statement log or pgAudit records, read from the standard server log | Which tables each transaction writes together. How unevenly the values of a candidate shard key are accessed. | No, but recommended | Level 1 alone gives a complete bundle, and the bundle names what level 2 would -add. To add level 2, set `capture_log: true` under `database.workload`. pgAudit -is the default source. Platforms without pgAudit can use `log_fdw` (RDS and -Aurora) or `stderr` (self-managed hosts). See +add. To add level 2, set `capture_log: true` under `database.workload`, and set +`capture_log_type` to `statement` or `pgaudit`. `statement` needs no extension. +When a log holds both record types, the capture reads only the configured type. +`capture_log_source` is `file` by default, which reads a log that you export and +name in `capture_log_file`. On RDS and Aurora, `log_fdw` reads the log over the +tool's connection. See [Capturing transaction shapes and values](docs/workload_capture.md#capturing-transaction-shapes-and-values). Level 2 records literal values from your queries in `burst.csv`. Review that diff --git a/docs/providers/aws.md b/docs/providers/aws.md index 9d274d2..b59fa52 100644 --- a/docs/providers/aws.md +++ b/docs/providers/aws.md @@ -457,7 +457,7 @@ is a static parameter. For Aurora, set it on the cluster parameter group. so run it before you change anything. See [Workload Capture](../workload_capture.md#3-enable-pg_stat_statements). -### Capturing transaction shapes with pgAudit +### Capturing transaction shapes Beyond the query counts above, the capture can also read the server's own query log, to see which statements ran in the same transaction and the literal @@ -465,37 +465,46 @@ values they carried. See [Capturing transaction shapes and values](../workload_capture.md#capturing-transaction-shapes-and-values) for what it collects and why you might want it. -RDS and Aurora can do this two ways. **pgAudit is the default.** Read the -comparison below before you choose `log_fdw` instead. +Logging is configured using two settings. `capture_log_type` is required: +`statement` or `pgaudit`. There is no default. `statement` needs no extension +but requires additional configuration to emit statements into the log. `pgaudit` +requires an extension to be enabled and installed. +`capture_log_source` is where that log is read from. It defaults to `file`. +`log_fdw` reads the csv log over this connection instead. -#### log_fdw or pgAudit +#### Log type -| | pgAudit (`capture_log_source: pgaudit`) | log_fdw (`capture_log_source: log_fdw`) | +| | `statement` | `pgaudit` | +| --- | --- | --- | +| Extension | None | `shared_preload_libraries`, one reboot, `CREATE EXTENSION pgaudit` | +| Per-capture setup | `log_min_duration_statement = 0` for the window, then reset it | None if logging is preconfigured using server parameters, runtime per-role configuration possible | +| Scope of what is logged | The whole server, every role | Configurable at server level and per-role | +| What it records | Statement, session, transaction framing, per-statement duration, error text | Statement, session, transaction framing, bind values | +| Log volume | Every statement, at `log_min_duration_statement = 0` | Depends on configured logging scope, up to as much as "every statement" | +| Superuser traffic | Logged like any other | Depends on configured logging scope | + +Use `statement` unless you want granular configuration and can accept a reboot. +The reboot is once per instance when enabling the extension, not once per capture. +Either log type can be read from a file or over `log_fdw`. + +#### Log source + +| | `file` (the default) | `log_fdw` | | --- | --- | --- | | Writes to your database | Nothing | A work schema, a foreign server, and the extension if absent, per `collect` | -| One-time setup | `shared_preload_libraries` in the parameter group, one reboot, `CREATE EXTENSION pgaudit` | `CREATE EXTENSION log_fdw`, no reboot | -| Per-capture setup | Two `ALTER ROLE` statements and a reset | None | -| Privilege the capture role needs | `pg_monitor` | `rds_superuser` | -| Scope of what is logged | One role, the application's | The whole server, every role | | Hands the log over by | An export you download and name in `capture_log_file` | This tool's own connection, no export | -| What it records | Statement, session, transaction framing, bind values | The same, plus per-statement duration and error text | -| Log volume it produces | The audited role's traffic | Every statement on the instance, at `log_min_duration_statement = 0` | -| Superuser traffic | Not reliably audited | Logged like any other | -| Cleanup afterwards | Reset the role, delete the exported files | The same, plus `workload init --cleanup` if a `collect` was killed | +| Privilege the capture role needs | `pg_monitor` (for stats, not the log itself) | `rds_superuser` | +| What it can read | The exported log (plain text or csvlog) for either type | csvlog only, for either type | +| Cleanup afterwards | Delete the exported files | The same, plus `workload init --cleanup` if a `collect` was killed, disable extension if no longer needed | -- **Choose pgAudit** unless you cannot reboot. It logs one role and writes - nothing to your database, and the reboot is once per instance however many - captures you run. -- **Choose `log_fdw`** when a reboot on a production primary is not something - you can schedule, or when you need the durations and error text that only the - server log carries. It costs `rds_superuser` for the capture role, objects - created and dropped inside each `collect`, and a log holding every role's - statements rather than one. +Use `file` unless you cannot export files. `log_fdw` needs `csvlog` in +`log_destination`. It requires `rds_superuser` for the capture role, and +creates/drops objects inside each `collect`. It does not change what is logged. -Both produce the same bundle. +Both types and both sources produce the same bundle. -`ps-discovery workload init --check` reports which of the two this server can -supply today, and names any objects a killed `collect` left behind. +`ps-discovery workload init --check` reports what this server can supply, and +names any objects a killed `collect` left behind. #### Permissions for log capture @@ -572,7 +581,7 @@ Name the file in `config.yaml`, then run the usual `collect`. database: workload: capture_log: true - capture_log_source: pgaudit + capture_log_type: pgaudit capture_log_file: pgaudit-capture.log ``` @@ -592,30 +601,33 @@ Then delete the log files you downloaded. They hold literal values from your queries. Watch `FreeStorageSpace` in CloudWatch level off to confirm the logging stopped. -### Capturing over log_fdw instead +### Reading the log over log_fdw -Use this only after reading -[log_fdw or pgAudit](#log_fdw-or-pgaudit) above. +Use this only after reading [Log source](#log-source) above. `log_fdw` reads +the csv log. It only defines how the log is read, not what log type the server +generates (statement log or pgAudit). Set `capture_log_type` for that. **Set up:** `CREATE EXTENSION log_fdw;` as a member of `rds_superuser`, and grant `rds_superuser` to the capture role, which needs it to create the foreign server and read the log files. Reconnect after the grant, since it does not reach an open session. Set `log_destination` to include `csvlog` in the -parameter group, and `log_min_duration_statement = 0` for the window. Neither -needs a reboot, and both need the parameter-group permissions listed under +parameter group. That needs no reboot, and it needs the parameter-group +permissions listed under [Permissions for log capture](#permissions-for-log-capture). ```yaml database: workload: capture_log: true + capture_log_type: statement capture_log_source: log_fdw capture_log_seconds: 600 ``` Each `collect` then watches the log for `capture_log_seconds` and reads that window over its own connection. Keep the schedule interval longer than -`capture_log_seconds`. +`capture_log_seconds`. The same connection reads a pgAudit csv log when +`capture_log_type` is `pgaudit` instead of `statement`. **Turn it off afterwards:** return `log_min_duration_statement` to the value it had. RDS re-adds `stderr` alongside `csvlog`, so every event is written twice diff --git a/docs/providers/gcp.md b/docs/providers/gcp.md index d1cd348..50ae98b 100644 --- a/docs/providers/gcp.md +++ b/docs/providers/gcp.md @@ -558,7 +558,7 @@ yourself. database: workload: capture_log: true - capture_log_source: pgaudit-json + capture_log_type: pgaudit capture_log_file: pgaudit-capture.jsonl ``` @@ -584,8 +584,8 @@ timestamp>="..." timestamp<="..." The payload is the same `PgAuditEntry`, and chunking is reassembled by the tool the same way. -**Read it:** the same `capture_log_source: pgaudit-json` config block as -Cloud SQL, above. +**Read it:** the same `capture_log_type: pgaudit` config block as Cloud SQL, +above. The JSON encoding is read from the file. ## Additional Resources diff --git a/docs/providers/supabase.md b/docs/providers/supabase.md index 1811987..089cc49 100644 --- a/docs/providers/supabase.md +++ b/docs/providers/supabase.md @@ -422,7 +422,7 @@ Name the file in `config.yaml`, then run the usual `collect`. database: workload: capture_log: true - capture_log_source: pgaudit + capture_log_type: pgaudit capture_log_file: exported.log ``` diff --git a/docs/workload_capture.md b/docs/workload_capture.md index 25ba726..ffba13a 100644 --- a/docs/workload_capture.md +++ b/docs/workload_capture.md @@ -392,8 +392,8 @@ The directory is mode `0700`. Delete it when you have the bundle. To stop collection, remove the cron entry and delete the directory. Nothing persists on the database server with `capture_log` off, and nothing -persists with the default `pgaudit` source either, which reads a file you -exported. Only `capture_log_source: log_fdw` writes to the server: each collect +persists when `capture_log_source` is `file`, which reads a log you exported. +Only `capture_log_source: log_fdw` writes to the server: each collect creates a work schema, a foreign server and the `log_fdw` extension to read the log, and drops all three when it finishes. If a collect is killed between the two, the next one clears what was left behind, and @@ -687,13 +687,16 @@ This is a setting, not a command. Nothing about the flow changes: the same database: workload: capture_log: true - capture_log_source: pgaudit - capture_log_file: pgaudit-capture.log + capture_log_type: statement + capture_log_file: exported.log ``` -`pgaudit` is the default source. You turn statement logging on for one role, -export the window through your provider's own tooling, and name the file. Each -`collect` takes its snapshot as before and then reads the file. +`capture_log_type` is required. `statement` reads the server's statement log +and needs no extension. `pgaudit` reads pgAudit. `capture_log_source` defaults +to `file`: you export the window and name it in `capture_log_file`. Each +`collect` takes its snapshot as before and then reads the file. The tool reads +the file's encoding itself: plain text, csvlog, or a Cloud SQL or AlloyDB +pgAudit JSON export. With `capture_log_source: log_fdw`, the tool reads the server's log over its own connection instead, and each `collect` watches the log for @@ -712,42 +715,42 @@ reason as a warning and exits 0. A capture never fails because of this. ### Where the log comes from -`capture_log_source` selects how the log is read. +`capture_log_type` selects what was logged. `capture_log_source` selects where +the log is read. `capture_log_type` is required when `capture_log` is on. +`statement` needs no extension. `pgaudit` logs one role's traffic, it needs no +write of any kind to your database, and every setting it uses can be changed +without a restart after the extension is enabled. -**pgAudit is the default.** It logs one role's traffic, it needs no write of any -kind to your database, and every setting it uses can be changed without a -restart. Use another source only where pgAudit cannot run. +`capture_log_source` defaults to `file`. `log_fdw` reads the csv log the server +is already writing, over this connection, on RDS and Aurora. That read is +always csvlog. pgAudit over it is `capture_log_type: pgaudit` plus +`capture_log_source: log_fdw`. -| Platform | `capture_log_source` | Reads the log by | -| --- | --- | --- | -| RDS, Aurora | `pgaudit` | a file you export | -| RDS, Aurora, no pgAudit available | `log_fdw` | SQL, over `log_fdw` | -| Cloud SQL, AlloyDB | `pgaudit` or `pgaudit-json` | a file you export | -| Supabase | `pgaudit` | a file you export | -| Self-managed | `pgaudit`, or `stderr` with no extension | the log file on the host | -| Neon, Heroku Postgres, PlanetScale | cannot supply a log | see the provider's guide | +| `capture_log_type` | Supported `capture_log_source` values | +| --- | --- | +| `statement` | `file` (plain text or csvlog). `log_fdw` on RDS and Aurora, which reads csvlog over this connection | +| `pgaudit` | `file` (plain text, csvlog, or a Cloud SQL or AlloyDB JSON export). `log_fdw` on RDS and Aurora, which reads csvlog over this connection | -`ps-discovery workload init --check` reports which row your server is on. It -opens a connection, prints what the server can supply, creates no session and -changes nothing. Add `--json` to sweep an estate. +`ps-discovery workload init --check` reports what this server can supply. It +opens a connection, creates no session and changes nothing. Add `--json` to +sweep an estate. `log_fdw` reads the log the server is already writing, over this tool's own connection, so it needs no export and no restart. It needs the `log_fdw` extension, csvlog output and a role holding `rds_superuser`, and it is the only source that creates objects in your database. -[The AWS guide](providers/aws.md#log_fdw-or-pgaudit) compares it with pgAudit in -full. To get statements into the log in the first place, set +[The AWS guide](providers/aws.md#log-source) compares `file` and `log_fdw`. To get statements into the log in the first place, set `log_min_duration_statement = 0` for the window; see [Set the logging up yourself](#set-the-logging-up-yourself). -For the three file sources, export the log through the provider's own tooling -and name the file: +For a file source, export the log through the provider's own tooling and name +the file: ```yaml database: workload: capture_log: true - capture_log_source: pgaudit + capture_log_type: pgaudit capture_log_file: exported.log ``` @@ -762,18 +765,20 @@ read and the run says that it took no snapshot. A session of nothing but imports has no window to measure the log against, so run at least two `collect`s that can connect. -`pgaudit-json` reads Cloud SQL's and AlloyDB's JSON export, either one -`PgAuditEntry` payload per line or a `gcloud logging read --format=json` array. -`stderr` reads a plain server log, for hosts where pgAudit cannot be installed -at all. +A file is classified before it is parsed. A leading `{` or `[` is the Cloud +SQL and AlloyDB `PgAuditEntry` export, either one payload per line or a +`gcloud logging read --format=json` array, and it is read only when +`capture_log_type` is `pgaudit`. A csvlog row is csv. Anything else with a +severity marker (`LOG:`, `ERROR:`, and the rest) is plain text. +`capture_log_type: statement` pointed at a JSON export records no log window. Per-provider export steps: -[RDS and Aurora](providers/aws.md#capturing-transaction-shapes-with-pgaudit), +[RDS and Aurora](providers/aws.md#capturing-transaction-shapes), [Cloud SQL and AlloyDB](providers/gcp.md#capturing-transaction-shapes-with-pgaudit), [Supabase](providers/supabase.md#capturing-transaction-shapes-with-pgaudit). [Neon](providers/neon.md#capturing-transaction-shapes-with-pgaudit) and [Heroku](providers/heroku.md#capturing-transaction-shapes-with-pgaudit) cannot -produce this source; their guides say why. +produce a log this capture can read; their guides say why. ### Permissions for log capture @@ -856,8 +861,8 @@ statement and corrupts the record pairing. Keep `log_statement` at `'none'` during the window, or every statement is logged twice. Without pgAudit, set `log_min_duration_statement = 0` for the window instead, -leave `log_line_prefix` as it is (the tool reads it from the file), and use -`capture_log_source: stderr`. No extension is needed. The trade-off: the plain +leave `log_line_prefix` as it is (the tool reads it from the file), and set +`capture_log_type: statement`. No extension is needed. The trade-off: the plain log carries durations and pgAudit does not, but a busy server logs a great deal at `log_min_duration_statement = 0`. @@ -874,10 +879,10 @@ needs another. | AlloyDB | the `alloydb.enable_pgaudit` flag, restart, `CREATE EXTENSION pgaudit;` | | Supabase | enable pgAudit in the dashboard, or `CREATE EXTENSION pgaudit;`; no preload step | -Do this on RDS and Aurora too. A pgAudit capture writes nothing to your -database and logs one role rather than the whole server. Use `log_fdw` where -the reboot cannot be scheduled; see -[log_fdw or pgAudit](providers/aws.md#log_fdw-or-pgaudit). +Do this on RDS and Aurora too. A pgAudit capture logs one role. With +`capture_log_source: file` the tool writes nothing to the database. Without +the reboot, set `capture_log_type: statement`. See +[Log type](providers/aws.md#log-type). Some limits carry into every pgAudit capture. None of them stops a capture being useful, but each one qualifies what a transaction shape proves. diff --git a/planetscale_discovery/config/config_manager.py b/planetscale_discovery/config/config_manager.py index c207e98..8ed1589 100644 --- a/planetscale_discovery/config/config_manager.py +++ b/planetscale_discovery/config/config_manager.py @@ -49,7 +49,10 @@ class WorkloadConfig: statement_row_limit: int = 20000 capture_log: bool = False capture_log_seconds: int = 600 - capture_log_source: str = "pgaudit" + # None until capture_log is on, so a discovery run does not have to set it. + capture_log_type: Optional[str] = None + # file writes nothing; log_fdw is the only source that creates server objects. + capture_log_source: str = "file" capture_log_file: Optional[str] = None # The distribution tier, off unless declared: pg_stats for these columns, # verbatim unless hash_values is true, which hashes them at finalize. @@ -58,6 +61,34 @@ class WorkloadConfig: distributions_hash_values: bool = False +# Source is where the log is read. Type is what was logged. Encoding is sniffed. +_LOG_SOURCES: Tuple[str, ...] = ("file", "log_fdw") +_LOG_TYPES: Tuple[str, ...] = ("statement", "pgaudit") + + +def workload_log_errors(workload: WorkloadConfig) -> List[str]: + """Bad source or type. handle_workload must call this: it exits before validate_config().""" + errors: List[str] = [] + source = workload.capture_log_source + if source not in _LOG_SOURCES: + errors.append( + "database.workload.capture_log_source must be one of " + f"{', '.join(_LOG_SOURCES)}, got {source!r}" + ) + log_type = workload.capture_log_type + if log_type is not None and log_type not in _LOG_TYPES: + errors.append( + "database.workload.capture_log_type must be one of " + f"{', '.join(_LOG_TYPES)}, got {log_type!r}" + ) + if workload.capture_log and not log_type: + errors.append( + "database.workload.capture_log_type is required when " + "capture_log is true. Set it to statement or pgaudit." + ) + return errors + + @dataclass class DatabaseConfig: """Database connection configuration.""" @@ -824,15 +855,8 @@ def _load_from_environment(self) -> DiscoveryConfig: } # Values restricted to a fixed set, as dotted paths. - _ENUMS: Dict[str, Tuple[str, ...]] = { - "database.workload.capture_log_source": ( - "pgaudit", - "pgaudit-json", - "stderr", - "log_fdw", - "auto", - ), - } + # Source and type are not here: a workload command never reaches validate_config(). + _ENUMS: Dict[str, Tuple[str, ...]] = {} def _resolve_path(self, path: str) -> Any: """Walk a dotted path into the loaded config. Returns None if absent.""" @@ -857,6 +881,13 @@ def _validate_numeric_ranges(self, errors: List[str]) -> None: if not (low <= value <= high): errors.append(f"{path} must be between {low} and {high}, got {value}") + def _validate_workload_log(self, errors: List[str]) -> None: + """Require a log type when capture is on, and reject a bad source or type.""" + config = self.config + if config is None: + return + errors.extend(workload_log_errors(config.database.workload)) + def _validate_enums(self, errors: List[str]) -> None: """Check every entry in _ENUMS, appending to errors.""" for path, allowed in self._ENUMS.items(): @@ -877,6 +908,8 @@ def _validate_config(self) -> None: self._validate_numeric_ranges(errors) self._validate_enums(errors) + # Required only when capture_log is on; an explicit bad type still fails when it is off. + self._validate_workload_log(errors) # Validate engine valid_engines = {"postgres", "mysql"} @@ -1004,9 +1037,12 @@ def save_config_template( # max_session_mb: 512 # Hard stop; the session refuses to grow past it # capture_log: false # Also read the query log. Puts literal values # # in the bundle; see docs/workload_capture.md + # capture_log_type: # Required when capture_log is true. + # # statement (no extension) or pgaudit # capture_log_seconds: 600 # log_fdw only: each collect watches this long # # (10 to 3600). Keep the schedule longer - # capture_log_source: pgaudit # pgaudit, pgaudit-json, stderr, log_fdw + # capture_log_source: file # file or log_fdw. file reads capture_log_file; + # # the file's encoding is read from the file # capture_log_file: # The exported log, for every source but log_fdw # distributions: # The columns your migration engineer names # by_column: # Every table carrying this column name diff --git a/planetscale_discovery/workload/burst/__init__.py b/planetscale_discovery/workload/burst/__init__.py index 76cb15c..9db0868 100644 --- a/planetscale_discovery/workload/burst/__init__.py +++ b/planetscale_discovery/workload/burst/__init__.py @@ -7,8 +7,10 @@ ) from planetscale_discovery.workload.burst.coverage import coverage, representativeness from planetscale_discovery.workload.burst.importer import ( + collect_csv_file, collect_pgaudit_file, collect_stderr_file, + sniff_packaging, ) from planetscale_discovery.workload.burst.readiness import LogCaptureProbe from planetscale_discovery.workload.logs.pgaudit import ObjectLoggingError @@ -22,8 +24,10 @@ "find_leftovers", "coverage", "representativeness", + "collect_csv_file", "collect_pgaudit_file", "collect_stderr_file", + "sniff_packaging", "LogCaptureProbe", "ObjectLoggingError", ] diff --git a/planetscale_discovery/workload/burst/collector.py b/planetscale_discovery/workload/burst/collector.py index a28c19b..8dbb531 100644 --- a/planetscale_discovery/workload/burst/collector.py +++ b/planetscale_discovery/workload/burst/collector.py @@ -1,10 +1,17 @@ import re -from typing import Any, Dict, Iterator, List, Optional +from typing import Any, Dict, Iterator, List, Optional, Tuple from planetscale_discovery.common.base_analyzer import DatabaseAnalyzer from planetscale_discovery.common.utils import generate_timestamp from planetscale_discovery.workload.burst.sql import SafeCursor from planetscale_discovery.workload.logs.csvlog import to_statement +from planetscale_discovery.workload.logs.pgaudit import ( + ObjectLoggingError, + loss_summary, + new_summary, + records_from_csvlog, + summary_warnings, +) from planetscale_discovery.workload.logs.record import COLUMNS, KNOWN_WIDTHS from planetscale_discovery.workload.logs.timestamps import instant_utc from planetscale_discovery.workload.logs.transactions import ( @@ -50,6 +57,11 @@ def __init__( self._skipped_before_watermark = 0 self._skipped_after_watermark = 0 self._log_tz: Optional[str] = None + # pgAudit csv rows are AUDIT: lines; statement rows stay on to_statement. + self._audit = False + self._audit_summary: Optional[Dict[str, Any]] = None + # Persists across files so a repeated statement id is still object logging. + self._audit_last: List[Tuple[str, str]] = [] @property def errors_seen(self) -> List[str]: @@ -98,6 +110,11 @@ def collect(self) -> Dict[str, Any]: return result statements: List[Dict[str, Any]] = [] + # log_fdw is always csv; the type selects which parser reads message. + self._audit = (self.config or {}).get("capture_log_type") == "pgaudit" + if self._audit: + self._audit_summary = new_summary() + self._audit_last = [] log_timezone = self._log_timezone() self._log_tz = log_timezone self._since_instant = instant_utc(self.since, log_timezone) @@ -118,6 +135,9 @@ def collect(self) -> Dict[str, Any]: continue try: read = list(self._read_file(name)) + # Object logging corrupts the whole window, so this file is not skipped. + except ObjectLoggingError: + raise except Exception as exc: self.add_warning(f"could not read {name}: {exc}") result["files_skipped"].append({"file": name, "why": str(exc)}) @@ -137,6 +157,13 @@ def collect(self) -> Dict[str, Any]: result["sessions"] = _session_summary(statements) result["status"] = STATUS_OK if result["files_read"] else STATUS_DEGRADED result["warnings"] = list(self.errors_seen) + _skip_warnings(result) + # Same loss warnings a pgAudit file import reports. + if self._audit and self._audit_summary is not None: + result["warnings"] = list(result["warnings"]) + summary_warnings( + self._audit_summary + ) + if result["files_read"]: + result["loss_summary"] = loss_summary(self._audit_summary) result["window_start"], result["window_end"] = _event_window( statements, result["captured_at_server"], log_timezone ) @@ -239,7 +266,11 @@ def _read_file(self, name: str) -> Iterator[Dict[str, Any]]: f"{name} produced {width} columns, which is not a csvlog " "shape (23, 24 or 26)" ) - for row in self._stream(table, width): + rows: Iterator[Dict[str, Any]] = self._stream(table, width) + # to_statement drops AUDIT: lines, so parse the audit payload first. + if self._audit: + rows = records_from_csvlog(rows, self._audit_summary, self._audit_last) + for row in rows: statement = to_statement(row) if statement is None: continue diff --git a/planetscale_discovery/workload/burst/importer.py b/planetscale_discovery/workload/burst/importer.py index c7b00ef..9d36cd8 100644 --- a/planetscale_discovery/workload/burst/importer.py +++ b/planetscale_discovery/workload/burst/importer.py @@ -1,3 +1,5 @@ +import csv +import json import os import re from typing import Any, Callable, Dict, Iterator, List, Optional, TextIO @@ -10,13 +12,19 @@ _event_window, _session_summary, ) -from planetscale_discovery.workload.logs.csvlog import to_statement +from planetscale_discovery.workload.logs.csvlog import ( + is_csvlog_row, + read_records, + to_statement, +) from planetscale_discovery.workload.logs.pgaudit import ( loss_summary, new_summary, read_jsonl_records, read_log_records, + records_from_csvlog, sniff_prefix, + summary_warnings, ) from planetscale_discovery.workload.logs.stderrlog import ( read_records as stderr_records, @@ -26,6 +34,11 @@ MAX_EXPORT_BYTES = 2 * 1024 * 1024 * 1024 +_NO_PGAUDIT_STATEMENTS = ( + "no pgAudit statements found in it; a capture needs pgaudit.log " + "enabled for the window (the burst recipe is 'read,write,misc')" +) + def collect_pgaudit_file( path: str, *, json_export: bool, already_read=(), logger=None @@ -40,8 +53,7 @@ def records(handle: TextIO) -> Iterator[Dict[str, Any]]: burst = _collect_file( path, records, - "no pgAudit statements found in it; a capture needs pgaudit.log " - "enabled for the window (the burst recipe is 'read,write,misc')", + _NO_PGAUDIT_STATEMENTS, warnings_from=lambda: _pgaudit_warnings(summary), already_read=already_read, ) @@ -50,6 +62,82 @@ def records(handle: TextIO) -> Iterator[Dict[str, Any]]: return burst +def collect_csv_file(path: str, *, audit: bool, already_read=()) -> Dict[str, Any]: + """Read an exported csvlog. ``audit`` selects pgAudit rows over statement rows.""" + summary = new_summary() + + def records(handle: TextIO) -> Iterator[Dict[str, Any]]: + rows = read_records(handle) + if audit: + return records_from_csvlog(rows, summary) + return rows + + if audit: + empty_error = _NO_PGAUDIT_STATEMENTS + warnings_from: Optional[Callable[[], List[str]]] = lambda: _pgaudit_warnings( + summary + ) + else: + empty_error = ( + "no statements found in it; expected log_statement or " + "log_min_duration_statement rows ('statement: ...' / " + "'duration: ... ms') for the window" + ) + warnings_from = None + + burst = _collect_file( + path, + records, + empty_error, + warnings_from=warnings_from, + already_read=already_read, + ) + if audit and burst.get("files_read"): + burst["loss_summary"] = loss_summary(summary) + return burst + + +def sniff_packaging(path: str) -> str: + """json, csv, or text. An empty file or a ``[%p]`` log line is not json.""" + with open(path, encoding="utf-8", errors="replace", newline="") as handle: + if _json_export(handle.read(8192)): + return "json" + handle.seek(0) + try: + row = next(csv.reader(handle), None) + except csv.Error: + row = None + if row and is_csvlog_row(row): + return "csv" + handle.seek(0) + if sniff_prefix(handle) is not None: + return "text" + raise ValueError( + "no PostgreSQL log lines found in it; expected server-log " + "text with a severity marker (LOG:, ERROR:, ...)" + ) + + +def _json_export(sample: str) -> bool: + """A leading ``{`` is JSON. A leading ``[`` needs one value, then ``,`` or ``]``.""" + text = sample.lstrip(" \t\r\n\ufeff") + if not text: + return False + if text[0] == "{": + return True + if text[0] != "[": + return False + body = text[1:].lstrip() + try: + _, end = json.JSONDecoder().raw_decode(body) + except json.JSONDecodeError: + return False + rest = body[end:].lstrip() + if not rest or rest[0] == ",": + return True + return rest[0] == "]" and not rest[1:].strip() + + def collect_stderr_file(path: str, already_read=(), logger=None) -> Dict[str, Any]: summary: Dict[str, int] = {} @@ -128,7 +216,8 @@ def _collect_file( statements: List[Dict[str, Any]] = [] newest: Optional[str] = None - with open(path, encoding="utf-8", errors="replace") as handle: + # newline="" keeps a \r\n inside a quoted csv field. The text readers share this open. + with open(path, encoding="utf-8", errors="replace", newline="") as handle: for record in records_from(handle): statement = to_statement(record) if statement is None: @@ -187,30 +276,8 @@ def _clock_warnings( def _pgaudit_warnings(summary: Dict[str, Any]) -> List[str]: - warnings = [] - dropped = summary["dropped"] - if dropped: - tags = ", ".join(f"{tag} x{count}" for tag, count in dropped.most_common()) - warnings.append( - f"{sum(dropped.values())} audit record(s) outside the shape " - f"commands were dropped ({tags})" - ) - errors = sum(summary["errors"].values()) - if errors: - warnings.append( - f"{errors} ERROR line(s) were counted but not attached: pgAudit " - "logs no statement for one that never executed" - ) - for key, wording in ( - ("malformed", "audit record(s) did not parse and were dropped"), - ("substatements", "function-body record(s) were dropped"), - ("incomplete_chunks", "chunked statement(s) never completed"), - ("skipped", "non-audit line(s) in the export were skipped"), - ("unmatched", "log line(s) did not match the sniffed prefix"), - ): - count = summary[key] - if count: - warnings.append(f"{count} {wording}") + warnings = summary_warnings(summary) + # A JSON array loaded whole is a file-export problem; log_fdw never hits it. if summary.get("array_bytes", 0) > _ARRAY_WARN_BYTES: warnings.append( "the export was one JSON array over 100 MB and was parsed whole " diff --git a/planetscale_discovery/workload/burst/readiness.py b/planetscale_discovery/workload/burst/readiness.py index 1509739..fb29b33 100644 --- a/planetscale_discovery/workload/burst/readiness.py +++ b/planetscale_discovery/workload/burst/readiness.py @@ -52,7 +52,8 @@ "pgaudit.log_relation off (object logging corrupts transaction-shape " "assembly) and log_statement at 'none' during the window, or every " "statement is logged twice. Then export the log from the provider's " - "console or logging API and name the file in capture_log_file." + "console or logging API, set capture_log_type: pgaudit, and name the " + "file in capture_log_file." ) PGAUDIT_RECOMMENDATION_AVAILABLE = ( "pgAudit is available on this server but not installed. The one-time " @@ -62,20 +63,14 @@ "can need a restart; after it every pgAudit setting is role-scoped and " "changeable without one, and a capture is two ALTER ROLEs." ) -PGAUDIT_RECOMMENDATION_ABSENT = { - "rds": ( - "pgAudit is not available on this server. Set capture_log_source: " - "log_fdw to read the log over this connection instead. It creates a " - "work schema and a foreign server for the length of each collect, and " - "it reads the whole log rather than one role's traffic; " - "docs/providers/aws.md compares the two." - ), - "default": ( - "pgAudit is not available on this server. Set " - "log_min_duration_statement = 0 for the window, export the log, and " - "read it with capture_log_source: stderr." - ), -} +PGAUDIT_RECOMMENDATION_ABSENT = ( + "pgAudit is not available on this server. Use a statement log instead: " + "set log_min_duration_statement = 0 for the window, export the log, and " + "set capture_log_type: statement. The file's encoding is read from the " + "file. Or enable pgAudit first (a parameter group on RDS and Aurora, the " + "cloudsql.enable_pgaudit flag on Cloud SQL, alloydb.enable_pgaudit on " + "AlloyDB), then CREATE EXTENSION pgaudit and set capture_log_type: pgaudit." +) class LogCaptureProbe(DatabaseAnalyzer): @@ -102,7 +97,7 @@ def run(self) -> Dict[str, Any]: requirements = self._requirements( settings, version, provider, destinations, prefix, log_fdw ) - pgaudit = self._pgaudit(provider) + pgaudit = self._pgaudit() notes = self._notes(settings, version, provider, destinations, prefix) notes.extend(self._pgaudit_notes(pgaudit)) return { @@ -210,9 +205,10 @@ def _requirements( "Aurora, CREATE EXTENSION log_fdw makes the server's own " "log readable through a foreign table over this same " "connection, at the cost of a work schema and a foreign " - "server for the length of each collect. The default " - "source, pgAudit, reads a log you export and writes " - "nothing to the database." + "server for the length of each collect. Set " + "capture_log_type to statement or pgaudit for that read. " + "A file source reads a log you export and writes nothing " + "to the database." ), where=provider, ), @@ -377,7 +373,7 @@ def _log_fdw(self) -> str: return "installed" return "available" if row.get("available") else "absent" - def _pgaudit(self, provider: str = "unknown") -> Dict[str, Any]: + def _pgaudit(self) -> Dict[str, Any]: row = self._sql.one_or_none( "SELECT (SELECT extversion FROM pg_extension" " WHERE extname = 'pgaudit') AS installed_version," @@ -394,9 +390,7 @@ def _pgaudit(self, provider: str = "unknown") -> Dict[str, Any]: recommendation = PGAUDIT_RECOMMENDATION_AVAILABLE else: state = "absent" - recommendation = PGAUDIT_RECOMMENDATION_ABSENT.get( - provider, PGAUDIT_RECOMMENDATION_ABSENT["default"] - ) + recommendation = PGAUDIT_RECOMMENDATION_ABSENT return { "state": state, "installed": bool(installed_version), diff --git a/planetscale_discovery/workload/cli_workload.py b/planetscale_discovery/workload/cli_workload.py index 533fa72..da5e684 100644 --- a/planetscale_discovery/workload/cli_workload.py +++ b/planetscale_discovery/workload/cli_workload.py @@ -13,7 +13,10 @@ from typing import Any, Dict from planetscale_discovery import __version__ -from planetscale_discovery.config.config_manager import WorkloadConfig +from planetscale_discovery.config.config_manager import ( + WorkloadConfig, + workload_log_errors, +) from planetscale_discovery.workload.bundle import bundle_dir_name, write_bundle from planetscale_discovery.workload.codes import label from planetscale_discovery.workload.collect import ( @@ -41,9 +44,8 @@ WORKLOAD_COMMANDS = ("init", "collect", "finalize", "status") SOURCE_LIVE = "log_fdw" -SOURCE_LIVE_ALIAS = "auto" -SOURCE_PGAUDIT_JSON = "pgaudit-json" -SOURCE_STDERR = "stderr" +# Packaging is sniffed from the file; this only selects the pgAudit parser. +TYPE_PGAUDIT = "pgaudit" # The reuse boundary: the schema comes from the existing analyzer, not a new query. CATALOG_MODULES = ["schema"] @@ -83,11 +85,14 @@ Set 'capture_log: true' under database.workload to also read the server's query log, which carries the transaction shapes and the literal values that -pg_stat_statements does not record. The default source is pgAudit: you export -the window and name the file in capture_log_file. On RDS and Aurora, -capture_log_source: log_fdw reads the log over this connection instead, and -each collect then watches it for capture_log_seconds. 'init --check' reports -whether this server can supply it. +pg_stat_statements does not record. capture_log_type is required: 'statement' +needs no extension and reads statement: records, and 'pgaudit' reads pgaudit +records. If a log contains both record types, the capture consumes only the +configured record type and not both. The default source is a file you export and +name in capture_log_file. The tool reads whether that file is plain text, csvlog, +or a pgAudit JSON export. On RDS and Aurora, capture_log_source: log_fdw reads +the csv log over this connection instead, and each collect then watches it for +capture_log_seconds. 'init --check' reports whether this server can supply it. Declare the columns your migration engineer names under 'distributions:' in database.workload. Each collect records their most common values from pg_stats, @@ -246,13 +251,22 @@ def handle_workload(args, config, logger) -> int: ) return EXIT_USAGE - if not _workload_config(config).enabled: + workload = _workload_config(config) + if not workload.enabled: logger.error( "workload capture is disabled. Set database.workload.enabled: " "true in the config file to run this command." ) return EXIT_USAGE + # Discovery validation runs after this command has already exited, so a + # bad source or type has to be rejected here or collect treats it as a file. + log_errors = workload_log_errors(workload) + if log_errors: + for error in log_errors: + logger.error(error) + return EXIT_USAGE + if command == "init": return _init(args, config, logger) if command == "collect": @@ -268,8 +282,6 @@ def _workload_config(config) -> WorkloadConfig: # Type-checked, or a stub config answers getattr for every field. candidate = getattr(getattr(config, "database", None), "workload", None) workload = candidate if isinstance(candidate, WorkloadConfig) else WorkloadConfig() - if workload.capture_log_source == SOURCE_LIVE_ALIAS: - workload.capture_log_source = SOURCE_LIVE return workload @@ -466,7 +478,7 @@ def _init(args, config, logger) -> int: def _report_log_readiness(connection, workload, logger) -> None: - """Say at init whether this server can supply what capture_log asks for.""" + """Init report for capture_log. A missing type is already rejected.""" if workload.capture_log_source != SOURCE_LIVE: if not workload.capture_log_file: logger.warning( @@ -478,7 +490,7 @@ def _report_log_readiness(connection, workload, logger) -> None: ) return logger.info( - f"capture_log is on, reading {workload.capture_log_source} records " + f"capture_log is on, reading {workload.capture_log_type} records " f"from {workload.capture_log_file}" ) return @@ -497,7 +509,7 @@ def _report_log_readiness(connection, workload, logger) -> None: "capture_log is on, but this server's log cannot be read over this " "connection, so every collect will record a snapshot only. Run " "'workload init --check' for what would have to change, or export the " - "log and set capture_log_source and capture_log_file." + "log and set capture_log_type, capture_log_source, and capture_log_file." ) @@ -690,6 +702,9 @@ def _capture_log_window(store, config, workload, logger) -> None: for warning in burst.get("warnings") or []: logger.warning(f" {warning}") return + # A collect that did not fail still carries warnings, including skipped files. + for warning in burst.get("warnings") or []: + logger.warning(warning) _report_burst(burst, store.append_burst(burst), logger) @@ -697,8 +712,10 @@ def _collect_exported_log(store, workload, logger) -> int: """Read a log the operator exported. Opens no database connection.""" from planetscale_discovery.workload.burst import ( ObjectLoggingError, + collect_csv_file, collect_pgaudit_file, collect_stderr_file, + sniff_packaging, ) path = workload.capture_log_file @@ -709,18 +726,43 @@ def _collect_exported_log(store, workload, logger) -> int: ) return EXIT_USAGE + # Encoding is a property of the file, not a config key. try: - if workload.capture_log_source == SOURCE_STDERR: - burst = collect_stderr_file( - path, already_read=store.log_files_read(), logger=logger - ) - else: + packaging = sniff_packaging(path) + except (OSError, ValueError) as exc: + logger.warning( + f"could not read {path}, so this collect recorded a snapshot " + f"only: {exc}" + ) + return EXIT_OK + logger.info(f"log packaging read from the file as {packaging}") + # The JSON reader only understands a PgAuditEntry export. + if packaging == "json" and workload.capture_log_type != TYPE_PGAUDIT: + logger.warning( + f"{path} is a JSON log export, which is read only when " + "capture_log_type is pgaudit. No log window was recorded." + ) + return EXIT_OK + + already_read = store.log_files_read() + # Type selects the logger. Packaging selects the reader. + try: + if packaging == "json": burst = collect_pgaudit_file( + path, json_export=True, already_read=already_read, logger=logger + ) + elif packaging == "csv": + burst = collect_csv_file( path, - json_export=workload.capture_log_source == SOURCE_PGAUDIT_JSON, - already_read=store.log_files_read(), - logger=logger, + audit=workload.capture_log_type == TYPE_PGAUDIT, + already_read=already_read, ) + elif workload.capture_log_type == TYPE_PGAUDIT: + burst = collect_pgaudit_file( + path, json_export=False, already_read=already_read, logger=logger + ) + else: + burst = collect_stderr_file(path, already_read=already_read, logger=logger) except ObjectLoggingError as exc: logger.warning(f"{path} cannot be used, so no log window was recorded: {exc}") return EXIT_OK @@ -1006,7 +1048,7 @@ def _check(args, config, logger) -> int: pgaudit = log["pgaudit"] print("") - print(f"pgAudit: {pgaudit['state']} (the default source)") + print(f"pgAudit: {pgaudit['state']} (set capture_log_type: pgaudit to read it)") for name, value in pgaudit["settings"].items(): if value is not None: print(f" {name} = {value}") diff --git a/planetscale_discovery/workload/logs/csvlog.py b/planetscale_discovery/workload/logs/csvlog.py index e7790f6..ba6caca 100644 --- a/planetscale_discovery/workload/logs/csvlog.py +++ b/planetscale_discovery/workload/logs/csvlog.py @@ -21,6 +21,26 @@ ERROR_SEVERITIES = ("ERROR", "FATAL", "PANIC") +# csvlog log_time is a timestamp alone; a text line's first field also holds the severity. +_TIMESTAMP_ONLY = re.compile( + r"\d{4}-\d{2}-\d{2}[ T]\d{2}:\d{2}:\d{2}(?:\.\d+)?" + r"(?:[ ]?[A-Za-z0-9+\-]+(?::\d{2})?)?" +) +_SEVERITY_ONLY = re.compile( + r"LOG|DETAIL|STATEMENT|ERROR|FATAL|PANIC|WARNING|NOTICE|INFO|HINT|CONTEXT" + r"|DEBUG\d?" +) + + +def is_csvlog_row(row: List[str]) -> bool: + """True when a parsed row has a csvlog shape, not a plain text line.""" + if len(row) not in KNOWN_WIDTHS: + return False + if _TIMESTAMP_ONLY.fullmatch(row[0].strip()) is None: + return False + severity = row[COLUMNS.index("error_severity")].strip() + return _SEVERITY_ONLY.fullmatch(severity) is not None + def read_records(source: Iterable[str]) -> Iterator[Dict[str, Any]]: for row in csv.reader(source): diff --git a/planetscale_discovery/workload/logs/pgaudit.py b/planetscale_discovery/workload/logs/pgaudit.py index 4033d69..113fd40 100644 --- a/planetscale_discovery/workload/logs/pgaudit.py +++ b/planetscale_discovery/workload/logs/pgaudit.py @@ -157,6 +157,63 @@ def new_summary() -> Dict[str, Any]: } +def records_from_csvlog( + rows: Iterable[Dict[str, Any]], + summary: Optional[Dict[str, Any]] = None, + last: Optional[List[Tuple[str, str]]] = None, +) -> Iterator[Dict[str, Any]]: + """AUDIT: rows become statements. An ERROR row is counted, then skipped.""" + counts = summary if summary is not None else new_summary() + seen: List[Tuple[str, str]] = last if last is not None else [] + for row in rows: + message = str(row.get("message") or "") + audit = AUDIT_RE.match(message) + if audit is None: + if str(row.get("error_severity") or "").strip() == "ERROR": + session = row.get("session_id") + if session in (None, ""): + session = row.get("process_id") + counts["errors"]["" if session is None else str(session)] += 1 + continue + fields = dict(row) + # log_fdw can return non-strings; session identity is compared as text. + for key in ("session_id", "process_id"): + if fields.get(key) is not None: + fields[key] = str(fields[key]) + record = _csv_record(fields, audit.group(1), counts, seen) + if record is not None: + yield record + + +def summary_warnings(summary: Dict[str, Any]) -> List[str]: + """Warnings a pgAudit summary can carry, apart from a huge JSON array.""" + warnings = [] + dropped = summary["dropped"] + if dropped: + tags = ", ".join(f"{tag} x{count}" for tag, count in dropped.most_common()) + warnings.append( + f"{sum(dropped.values())} audit record(s) outside the shape " + f"commands were dropped ({tags})" + ) + errors = sum(summary["errors"].values()) + if errors: + warnings.append( + f"{errors} ERROR line(s) were counted but not attached: pgAudit " + "logs no statement for one that never executed" + ) + for key, wording in ( + ("malformed", "audit record(s) did not parse and were dropped"), + ("substatements", "function-body record(s) were dropped"), + ("incomplete_chunks", "chunked statement(s) never completed"), + ("skipped", "non-audit line(s) in the export were skipped"), + ("unmatched", "log line(s) did not match the sniffed prefix"), + ): + count = summary[key] + if count: + warnings.append(f"{count} {wording}") + return warnings + + def loss_summary(summary: Dict[str, Any]) -> Dict[str, Any]: out = {} for key, value in summary.items(): @@ -308,8 +365,9 @@ def _json_entries(lines: Iterable[str], counts: Dict[str, Any]) -> Iterator[Any] counts["skipped"] += 1 +# A log_fdw row carries datetimes and ints, not only CSV strings. def _csv_record( - fields: Dict[str, str], + fields: Dict[str, Any], text: str, counts: Dict[str, Any], last: List[Tuple[str, str]], diff --git a/tests/unit/test_workload_burst_collector.py b/tests/unit/test_workload_burst_collector.py index 3bd3864..c68c1d1 100644 --- a/tests/unit/test_workload_burst_collector.py +++ b/tests/unit/test_workload_burst_collector.py @@ -2,6 +2,8 @@ from unittest.mock import MagicMock +import pytest + from planetscale_discovery.workload.burst.collector import ( STATUS_FAILED, BurstCollector, @@ -9,6 +11,7 @@ _session_summary, _table_name, ) +from planetscale_discovery.workload.logs.pgaudit import ObjectLoggingError from planetscale_discovery.workload.logs.timestamps import instant_utc @@ -235,6 +238,79 @@ def test_drops_a_row_before_the_watermark_and_tags_kept_rows(self): assert collector._skipped_before_watermark == 1 +def _audit_row(message): + return { + "log_time": "2024-01-01 00:00:00.000 UTC", + "session_id": "5f1.3", + "error_severity": "LOG", + "message": message, + } + + +def _collect_live_rows(rows, log_type="pgaudit"): + """Run collect() through the mocked log_fdw read, with no database.""" + connection = MagicMock() + cursor = MagicMock() + cursor.__iter__.return_value = iter(rows) + connection.cursor.return_value.__enter__.return_value = cursor + collector = BurstCollector(connection, config={"capture_log_type": log_type}) + collector._sql = FakeSql( + execute_ok=True, + one_map={ + "pg_foreign_server": None, + "pg_extension": {"schema": "public"}, + "information_schema.columns": {"columns": 24}, + }, + all_map={ + "list_postgres_log_files": [ + {"file_name": "postgresql.csv", "file_size": 12} + ], + }, + ) + return collector.collect() + + +class TestStatementRowsStayOnTheStatementParser: + """The bug: a statement capture on log_fdw parsed AUDIT lines and dropped + statement lines.""" + + def test_a_statement_row_is_kept_and_an_audit_row_is_dropped(self): + result = _collect_live_rows( + [ + { + "log_time": "2024-01-01 00:00:00.000 UTC", + "session_id": "5f1.3", + "error_severity": "LOG", + "message": "statement: select from_statement", + }, + _audit_row( + 'AUDIT: SESSION,1,1,READ,SELECT,,,"select from_audit",' + ), + ], + log_type="statement", + ) + assert [row["sql"] for row in result["statements"]] == ["select from_statement"] + assert "loss_summary" not in result + + +class TestPgauditRowsAreParsedOnTheLiveRead: + """The bug: log_fdw handed AUDIT lines to to_statement, which drops them, + so a pgAudit capture stored an empty window.""" + + def test_an_audit_row_becomes_a_statement(self): + result = _collect_live_rows( + [_audit_row('AUDIT: SESSION,1,1,READ,SELECT,,,"select 1",')] + ) + assert result["statements"][0]["sql"] == "select 1" + assert "malformed" in result["loss_summary"] + + def test_an_object_audit_row_aborts_the_window(self): + with pytest.raises(ObjectLoggingError): + _collect_live_rows( + [_audit_row('AUDIT: OBJECT,1,1,READ,SELECT,,,"select 1",')] + ) + + class TestEventWindow: def test_min_and_max_across_statements(self): statements = [ diff --git a/tests/unit/test_workload_burst_importer.py b/tests/unit/test_workload_burst_importer.py index a7820c7..4019220 100644 --- a/tests/unit/test_workload_burst_importer.py +++ b/tests/unit/test_workload_burst_importer.py @@ -1,5 +1,7 @@ """Tests for collect_pgaudit_file and collect_stderr_file.""" +import csv +import io import json import pytest @@ -7,9 +9,12 @@ from planetscale_discovery.workload.burst.importer import ( MAX_EXPORT_BYTES, _clock_warnings, + collect_csv_file, collect_pgaudit_file, collect_stderr_file, + sniff_packaging, ) +from planetscale_discovery.workload.logs.record import COLUMNS, WIDTH_PG13 def write(tmp_path, name, content): @@ -113,6 +118,155 @@ def test_a_file_over_the_ceiling_names_the_limit(self, tmp_path, mocker): collect_stderr_file(path) +def _csvlog_line(**overrides): + values = {name: "" for name in COLUMNS[:WIDTH_PG13]} + values["log_time"] = "2024-01-01 00:00:00.000 UTC" + values["error_severity"] = "LOG" + values.update(overrides) + buffer = io.StringIO() + csv.writer(buffer).writerow([values[name] for name in COLUMNS[:WIDTH_PG13]]) + return buffer.getvalue() + + +class TestSniffPackaging: + def test_a_leading_brace_is_json(self, tmp_path): + path = write( + tmp_path, "export.json", '{"jsonPayload": {"command": "SELECT"}}\n' + ) + assert sniff_packaging(path) == "json" + + def test_a_csvlog_row_is_csv(self, tmp_path): + path = write( + tmp_path, + "postgresql.csv", + _csvlog_line(message="statement: select 1"), + ) + assert sniff_packaging(path) == "csv" + + def test_a_csv_message_containing_a_severity_marker_stays_csv(self, tmp_path): + """The bug: a csv message containing LOG: was read as plain text.""" + path = write( + tmp_path, + "postgresql.csv", + _csvlog_line(message="statement: select 1 LOG: extra"), + ) + assert sniff_packaging(path) == "csv" + + def test_a_leading_bracket_is_json(self, tmp_path): + """The bug: a JSON array export was not recognized as json.""" + path = write( + tmp_path, "export.json", '[{"jsonPayload": {"command": "SELECT"}}]\n' + ) + assert sniff_packaging(path) == "json" + + def test_a_pretty_printed_array_is_json(self, tmp_path): + path = write( + tmp_path, + "export.json", + '[\n {"jsonPayload": {"command": "SELECT"}}\n]\n', + ) + assert sniff_packaging(path) == "json" + + def test_a_compact_array_longer_than_the_sample_is_json(self, tmp_path): + """The bug: a one-line JSON array over 8192 bytes was not recognized.""" + records = [{"jsonPayload": {"command": "SELECT", "statement": "select 1"}}] + records.extend({"n": i, "pad": "x" * 40} for i in range(200)) + text = json.dumps(records) + assert len(text) > 8192 + path = write(tmp_path, "export.json", text) + assert sniff_packaging(path) == "json" + + def test_a_pretty_object_longer_than_the_sample_is_json(self, tmp_path): + """The bug: a pretty-printed object over 8192 bytes was not recognized.""" + payload = {"entries": [{"n": i, "statement": "select 1"} for i in range(200)]} + text = json.dumps(payload, indent=2) + assert len(text) > 8192 + path = write(tmp_path, "export.json", text) + assert sniff_packaging(path) == "json" + + def test_a_bracket_prefix_is_text(self, tmp_path): + """The bug: a log_line_prefix of [%p] was read as a JSON export.""" + path = write( + tmp_path, + "postgresql.log", + "[4242] LOG: statement: select 1\n", + ) + assert sniff_packaging(path) == "text" + + def test_an_empty_file_is_not_json(self, tmp_path): + """The bug: a blank file was reported as a JSON log export.""" + path = write(tmp_path, "empty.log", "") + with pytest.raises(ValueError, match="no PostgreSQL log lines"): + sniff_packaging(path) + + def test_a_known_width_row_without_a_timestamp_is_text(self, tmp_path): + """The bug: a known-width row whose first field is not a timestamp was read as csv.""" + path = write( + tmp_path, + "postgresql.csv", + _csvlog_line( + log_time="not-a-timestamp", + message="statement: select 1 LOG: extra", + ), + ) + assert sniff_packaging(path) == "text" + + def test_a_csv_error_falls_through_to_text(self, tmp_path): + """The bug: a field over the csv limit escaped sniff as csv.Error.""" + path = write( + tmp_path, + "postgresql.log", + "2024-01-01 00:00:00.000 UTC [111] LOG: statement: select 1\n", + ) + limit = csv.field_size_limit() + csv.field_size_limit(8) + try: + assert sniff_packaging(path) == "text" + finally: + csv.field_size_limit(limit) + + def test_a_text_audit_line_with_commas_is_text(self, tmp_path): + path = write(tmp_path, "postgresql.log", PGAUDIT_LOG) + assert sniff_packaging(path) == "text" + + def test_a_file_with_no_log_line_raises(self, tmp_path): + path = write(tmp_path, "notes.txt", "not a log\n") + with pytest.raises(ValueError, match="no PostgreSQL log lines"): + sniff_packaging(path) + + +class TestCollectCsvFile: + def test_an_audit_csv_is_read_as_pgaudit(self, tmp_path): + path = write( + tmp_path, + "postgresql.csv", + _csvlog_line( + session_id="5f1.3", + message='AUDIT: SESSION,1,1,READ,SELECT,,,"select 1",', + ), + ) + result = collect_csv_file(path, audit=True) + assert result["statements"][0]["sql"] == "select 1" + + def test_a_statement_csv_is_read_as_statements(self, tmp_path): + path = write( + tmp_path, + "postgresql.csv", + _csvlog_line(message="statement: select 1"), + ) + result = collect_csv_file(path, audit=False) + assert result["statements"][0]["sql"] == "select 1" + + def test_a_crlf_inside_a_quoted_field_is_kept(self, tmp_path): + """The bug: a CR LF inside a quoted csvlog field was stored as LF.""" + path = tmp_path / "postgresql.csv" + path.write_bytes( + _csvlog_line(message="statement: select 1\r\nfrom t").encode("utf-8") + ) + result = collect_csv_file(str(path), audit=False) + assert result["statements"][0]["sql"] == "select 1\r\nfrom t" + + class TestClockWarnings: def test_unresolved_zone_abbreviation_warns(self): statements = [{"log_time": "2024-01-01 00:00:00 EST"}] diff --git a/tests/unit/test_workload_cli_log_capture.py b/tests/unit/test_workload_cli_log_capture.py index 659995f..7f15f8d 100644 --- a/tests/unit/test_workload_cli_log_capture.py +++ b/tests/unit/test_workload_cli_log_capture.py @@ -1,5 +1,9 @@ """collect reads the query log when the config asks for it, and never fails on it.""" +import csv +import io +import json + import pytest from planetscale_discovery.config.config_manager import WorkloadConfig @@ -9,6 +13,7 @@ EXIT_USAGE, _collect, ) +from planetscale_discovery.workload.logs.record import COLUMNS, WIDTH_PG13 from planetscale_discovery.workload.store import WorkloadStore @@ -83,7 +88,10 @@ def test_collect_still_exits_zero_and_keeps_the_snapshot(self, session, mocker): ) logger = mocker.Mock() workload = WorkloadConfig( - capture_log=True, capture_log_source="log_fdw", capture_log_seconds=10 + capture_log=True, + capture_log_type="statement", + capture_log_source="log_fdw", + capture_log_seconds=10, ) assert _collect(_Args(session.directory), _Config(workload), logger) == EXIT_OK @@ -111,7 +119,10 @@ def test_a_failed_burst_is_not_stored_as_a_used_window(self, session, mocker): ) logger = mocker.Mock() workload = WorkloadConfig( - capture_log=True, capture_log_source="log_fdw", capture_log_seconds=10 + capture_log=True, + capture_log_type="statement", + capture_log_source="log_fdw", + capture_log_seconds=10, ) assert _collect(_Args(session.directory), _Config(workload), logger) == EXIT_OK @@ -122,7 +133,7 @@ def test_a_failed_burst_is_not_stored_as_a_used_window(self, session, mocker): class TestAnExportedLogNeedsItsFile: def test_a_file_source_with_no_file_is_a_usage_error(self, session, mocker): - workload = WorkloadConfig(capture_log=True, capture_log_source="pgaudit") + workload = WorkloadConfig(capture_log=True, capture_log_type="statement") logger = mocker.Mock() assert ( @@ -137,7 +148,7 @@ def test_a_file_source_still_takes_its_snapshot(self, session, mocker, tmp_path) export.write_text("") workload = WorkloadConfig( capture_log=True, - capture_log_source="stderr", + capture_log_type="statement", capture_log_file=str(export), ) @@ -152,7 +163,7 @@ def test_an_unreadable_file_is_a_warning_not_a_non_zero_exit( _no_snapshot_connection(mocker) workload = WorkloadConfig( capture_log=True, - capture_log_source="stderr", + capture_log_type="statement", capture_log_file=str(tmp_path / "missing.log"), ) logger = mocker.Mock() @@ -161,6 +172,29 @@ def test_an_unreadable_file_is_a_warning_not_a_non_zero_exit( assert len(session.snapshot_paths()) == 1 assert logger.warning.called + def test_a_sniffed_log_with_no_statements_keeps_the_snapshot( + self, session, mocker, tmp_path + ): + """The bug: a sniffed log with no statements failed collect after the snapshot was stored.""" + _no_snapshot_connection(mocker) + export = tmp_path / "exported.log" + export.write_text( + "2024-01-01 00:00:00.000 UTC [111] FATAL: terminating connection " + "due to administrator command\n" + ) + workload = WorkloadConfig( + capture_log=True, + capture_log_type="statement", + capture_log_file=str(export), + ) + + assert ( + _collect(_Args(session.directory), _Config(workload), mocker.Mock()) + == EXIT_OK + ) + assert len(session.snapshot_paths()) == 1 + assert not session.burst_paths() + def test_a_file_already_read_stores_no_second_burst( self, session, mocker, tmp_path ): @@ -173,7 +207,7 @@ def test_a_file_already_read_stores_no_second_burst( ) workload = WorkloadConfig( capture_log=True, - capture_log_source="stderr", + capture_log_type="statement", capture_log_file=str(export), ) args, config = _Args(session.directory), _Config(workload) @@ -192,7 +226,7 @@ def test_a_file_source_survives_a_dead_connection(self, session, mocker, tmp_pat export.write_text("") workload = WorkloadConfig( capture_log=True, - capture_log_source="stderr", + capture_log_type="statement", capture_log_file=str(export), ) logger = mocker.Mock() @@ -202,6 +236,157 @@ def test_a_file_source_survives_a_dead_connection(self, session, mocker, tmp_pat assert logger.warning.called +class TestAStatementJsonExportIsNotRead: + def test_no_window_is_stored(self, session, mocker, tmp_path): + """The bug: a statement capture stored a burst from a JSON export that held a statement.""" + _no_snapshot_connection(mocker) + export = tmp_path / "export.json" + export.write_text( + json.dumps( + { + "timestamp": "2024-01-01T00:00:00Z", + "jsonPayload": { + "command": "SELECT", + "statement": "select 1", + "databaseSessionId": "abc", + "statementId": "1", + "auditType": "SESSION", + }, + } + ) + + "\n" + ) + workload = WorkloadConfig( + capture_log=True, + capture_log_type="statement", + capture_log_file=str(export), + ) + logger = mocker.Mock() + + assert _collect(_Args(session.directory), _Config(workload), logger) == EXIT_OK + assert len(session.snapshot_paths()) == 1 + assert not session.burst_paths() + assert ( + "read only when capture_log_type is pgaudit" + in logger.warning.call_args[0][0] + ) + + +def _csvlog_line(**overrides): + values = {name: "" for name in COLUMNS[:WIDTH_PG13]} + values["log_time"] = "2024-01-01 00:00:00.000 UTC" + values["error_severity"] = "LOG" + values.update(overrides) + buffer = io.StringIO() + csv.writer(buffer).writerow([values[name] for name in COLUMNS[:WIDTH_PG13]]) + return buffer.getvalue() + + +def _collect_export(session, mocker, tmp_path, name, content, log_type, logger=None): + _no_snapshot_connection(mocker) + export = tmp_path / name + export.write_text(content, encoding="utf-8") + workload = WorkloadConfig( + capture_log=True, + capture_log_type=log_type, + capture_log_file=str(export), + ) + return _collect( + _Args(session.directory), + _Config(workload), + logger if logger is not None else mocker.Mock(), + ) + + +class TestTheReaderFollowsTheSniffedPackaging: + """The bug: the file was classified, then collect called the reader for a + different packaging or type, and the window was empty.""" + + def test_a_pgaudit_text_file_is_read(self, session, mocker, tmp_path): + content = ( + "2024-01-01 00:00:00.000 UTC [111] LOG: AUDIT: " + 'SESSION,1,1,READ,SELECT,,,"select 1",\n' + ) + code = _collect_export( + session, mocker, tmp_path, "postgresql.log", content, "pgaudit" + ) + assert code == EXIT_OK + assert session.read_bursts()[0]["statements"][0]["sql"] == "select 1" + + def test_a_statement_csv_is_read(self, session, mocker, tmp_path): + code = _collect_export( + session, + mocker, + tmp_path, + "postgresql.csv", + _csvlog_line(message="statement: select 1"), + "statement", + ) + assert code == EXIT_OK + assert session.read_bursts()[0]["statements"][0]["sql"] == "select 1" + + def test_a_pgaudit_csv_is_read(self, session, mocker, tmp_path): + code = _collect_export( + session, + mocker, + tmp_path, + "postgresql.csv", + _csvlog_line( + session_id="5f1.3", + message='AUDIT: SESSION,1,1,READ,SELECT,,,"select 1",', + ), + "pgaudit", + ) + assert code == EXIT_OK + assert session.read_bursts()[0]["statements"][0]["sql"] == "select 1" + + def test_a_pgaudit_json_export_is_read(self, session, mocker, tmp_path): + content = json.dumps( + { + "timestamp": "2024-01-01T00:00:00Z", + "jsonPayload": { + "command": "SELECT", + "statement": "select 1", + "databaseSessionId": "abc", + "statementId": "1", + "auditType": "SESSION", + }, + } + ) + code = _collect_export( + session, mocker, tmp_path, "pgaudit.jsonl", content + "\n", "pgaudit" + ) + assert code == EXIT_OK + assert session.read_bursts()[0]["statements"][0]["sql"] == "select 1" + + +class TestAnObjectAuditRecordStoresNoWindow: + """The bug: object logging escaped collect, so cron read a finished + snapshot as a fault.""" + + def test_a_pgaudit_text_file_exits_zero_and_keeps_the_snapshot( + self, session, mocker, tmp_path + ): + logger = mocker.Mock() + content = ( + "2024-01-01 00:00:00.000 UTC [111] LOG: AUDIT: " + 'OBJECT,1,1,READ,SELECT,,,"select 1",\n' + ) + code = _collect_export( + session, + mocker, + tmp_path, + "postgresql.log", + content, + "pgaudit", + logger=logger, + ) + assert code == EXIT_OK + assert len(session.snapshot_paths()) == 1 + assert not session.burst_paths() + assert "cannot be used" in logger.warning.call_args[0][0] + + class TestTheWindowIsBounded: def test_the_burst_is_read_over_the_window_the_wait_held(self, session, mocker): _no_snapshot_connection(mocker) @@ -221,7 +406,10 @@ def test_the_burst_is_read_over_the_window_the_wait_held(self, session, mocker): return_value=collector, ) workload = WorkloadConfig( - capture_log=True, capture_log_source="log_fdw", capture_log_seconds=10 + capture_log=True, + capture_log_type="statement", + capture_log_source="log_fdw", + capture_log_seconds=10, ) _collect(_Args(session.directory), _Config(workload), mocker.Mock()) @@ -231,19 +419,51 @@ def test_the_burst_is_read_over_the_window_the_wait_held(self, session, mocker): assert len(session.burst_paths()) == 1 -class TestPgauditIsTheDefaultSource: +class TestASuccessfulBurstReportsItsLoss: + def test_warnings_are_logged_and_the_burst_is_stored(self, session, mocker): + """The bug: a successful pgAudit read stored the window and logged no loss warning.""" + _no_snapshot_connection(mocker) + mocker.patch.object(cli_workload, "_wait", return_value=False) + collector = mocker.Mock() + collector.collect.return_value = { + "status": "ok", + "window_start": "2026-09-14T10:00:00Z", + "window_end": "2026-09-14T10:10:00Z", + "statements": [], + "sessions": {}, + "files_read": ["a"], + "warnings": ["1 audit record(s) did not parse and were dropped"], + } + mocker.patch( + "planetscale_discovery.workload.burst.BurstCollector", + return_value=collector, + ) + logger = mocker.Mock() + workload = WorkloadConfig( + capture_log=True, + capture_log_type="pgaudit", + capture_log_source="log_fdw", + capture_log_seconds=10, + ) + + assert _collect(_Args(session.directory), _Config(workload), logger) == EXIT_OK + assert len(session.burst_paths()) == 1 + logger.warning.assert_any_call( + "1 audit record(s) did not parse and were dropped" + ) + + +class TestTheDefaultReadsAFile: """The bug: turning capture_log on wrote to the customer's database.""" def test_the_default_source_reads_a_file(self): - assert WorkloadConfig().capture_log_source == "pgaudit" + assert WorkloadConfig().capture_log_source == "file" - def test_auto_still_means_the_live_log_fdw_read(self): - workload = WorkloadConfig(capture_log=True, capture_log_source="auto") - config = _Config(workload) - assert cli_workload._workload_config(config).capture_log_source == "log_fdw" + def test_the_type_has_no_default(self): + assert WorkloadConfig().capture_log_type is None def test_a_file_source_with_no_file_warns_at_init(self, mocker): - workload = WorkloadConfig(capture_log=True) + workload = WorkloadConfig(capture_log=True, capture_log_type="statement") logger = mocker.Mock() cli_workload._report_log_readiness(mocker.Mock(), workload, logger) @@ -309,6 +529,12 @@ class TestTheNewSettingsAreValidated: ("capture_log_seconds: 4", "between 10 and 3600"), ("capture_log_source: pgaudit_json", "must be one of"), ("capture_log_source: logfdw", "must be one of"), + ("capture_log_source: pgaudit", "must be one of"), + ("capture_log_source: pgaudit-json", "must be one of"), + ("capture_log_source: stderr", "must be one of"), + ("capture_log_source: auto", "must be one of"), + ("capture_log_type: csv", "must be one of"), + ("capture_log: true", "capture_log_type is required"), ), ) def test_a_bad_value_is_rejected_at_load(self, tmp_path, yaml_body, expected): @@ -321,6 +547,38 @@ def test_a_bad_value_is_rejected_at_load(self, tmp_path, yaml_body, expected): with pytest.raises(ValueError, match=expected): ConfigManager(str(path)).load_config() + @pytest.mark.parametrize( + "yaml_body, log_type, source", + ( + ( + "capture_log: true\n" + " capture_log_type: statement\n" + " capture_log_source: file", + "statement", + "file", + ), + ( + "capture_log: true\n" + " capture_log_type: pgaudit\n" + " capture_log_source: log_fdw", + "pgaudit", + "log_fdw", + ), + ), + ) + def test_a_legal_combination_loads(self, tmp_path, yaml_body, log_type, source): + """The bug: a legal type and source were rejected at load.""" + from planetscale_discovery.config.config_manager import ConfigManager + + path = tmp_path / "config.yaml" + path.write_text( + "database:\n host: localhost\n workload:\n " + yaml_body + "\n" + ) + workload = ConfigManager(str(path)).load_config().database.workload + assert workload.capture_log is True + assert workload.capture_log_type == log_type + assert workload.capture_log_source == source + class TestAnInterruptedWaitStoresWhatItHeld: def test_it_is_a_shorter_window_not_an_error(self, session, mocker): @@ -340,7 +598,10 @@ def test_it_is_a_shorter_window_not_an_error(self, session, mocker): return_value=collector, ) workload = WorkloadConfig( - capture_log=True, capture_log_source="log_fdw", capture_log_seconds=600 + capture_log=True, + capture_log_type="statement", + capture_log_source="log_fdw", + capture_log_seconds=600, ) assert ( diff --git a/tests/unit/test_workload_cli_surface.py b/tests/unit/test_workload_cli_surface.py index a50b216..3d14013 100644 --- a/tests/unit/test_workload_cli_surface.py +++ b/tests/unit/test_workload_cli_surface.py @@ -124,6 +124,42 @@ class Args: assert handle_workload(Args(), self._Config("mysql"), logger) != EXIT_USAGE +class TestABadLogSourceIsRejected: + """The bug: collect treated logfdw and stderr as a file and stored a burst.""" + + class _Args: + workload_command = "collect" + session = "./wl" + + def _config(self, source): + class Database: + workload = WorkloadConfig( + enabled=True, + capture_log=True, + capture_log_type="statement", + capture_log_source=source, + ) + + class Config: + engine = "postgres" + database = Database() + + return Config() + + @pytest.mark.parametrize("source", ("logfdw", "stderr")) + def test_collect_exits_before_a_session(self, mocker, source): + store = mocker.patch( + "planetscale_discovery.workload.cli_workload.WorkloadStore" + ) + logger = mocker.Mock() + + code = handle_workload(self._Args(), self._config(source), logger) + + assert code == EXIT_USAGE + assert "capture_log_source" in logger.error.call_args[0][0] + store.assert_not_called() + + class _FinalizeDatabase: workload = WorkloadConfig(schemas=None) schemas = ["public"] diff --git a/tests/unit/test_workload_pgaudit.py b/tests/unit/test_workload_pgaudit.py index f505185..17c3e61 100644 --- a/tests/unit/test_workload_pgaudit.py +++ b/tests/unit/test_workload_pgaudit.py @@ -2,6 +2,7 @@ import pytest +from planetscale_discovery.workload.logs.csvlog import to_statement from planetscale_discovery.workload.logs.pgaudit import ( DEFAULT_PREFIX, ObjectLoggingError, @@ -9,8 +10,10 @@ new_summary, read_jsonl_records, read_log_records, + records_from_csvlog, sniff_prefix, ) +from planetscale_discovery.workload.logs.record import COLUMNS, WIDTH_PG13 def _audit_line(pid, statement_id, sql, audit_type="SESSION", command="SELECT"): @@ -71,6 +74,55 @@ def _chunk_payload(session, statement_id, chunk_count, chunk_index, statement_pa } +class TestRecordsFromCsvlog: + def test_an_audit_message_becomes_a_statement(self): + row = {name: "" for name in COLUMNS[:WIDTH_PG13]} + row["log_time"] = "2024-01-01 00:00:00.000 UTC" + row["session_id"] = "5f1.3" + row["error_severity"] = "LOG" + row["message"] = 'AUDIT: SESSION,1,1,READ,SELECT,,,"select 1",' + + records = list(records_from_csvlog([row])) + + assert len(records) == 1 + statement = to_statement(records[0]) + assert statement["sql"] == "select 1" + assert statement["session_id"] == "5f1.3" + + def test_a_non_audit_row_is_ignored(self): + row = {name: "" for name in COLUMNS[:WIDTH_PG13]} + row["message"] = "connection received: host=127.0.0.1" + assert list(records_from_csvlog([row])) == [] + + def test_an_error_row_is_counted_and_not_yielded(self): + """The bug: a csvlog ERROR row was ignored, so the unattached-error count stayed zero.""" + error = {name: "" for name in COLUMNS[:WIDTH_PG13]} + error["session_id"] = 12345 + error["error_severity"] = "ERROR" + error["message"] = "division by zero" + by_process = dict(error) + by_process["session_id"] = "" + by_process["process_id"] = 99 + fatal = dict(error) + fatal["error_severity"] = "FATAL" + fatal["session_id"] = 7 + summary = new_summary() + + assert list(records_from_csvlog([error, by_process, fatal], summary)) == [] + assert summary["errors"] == {"12345": 1, "99": 1} + + def test_an_integer_process_id_still_raises_on_a_repeated_statement(self): + """The bug: an integer process_id raised AttributeError, so collect skipped the file instead of rejecting object logging.""" + row = {name: "" for name in COLUMNS[:WIDTH_PG13]} + row["session_id"] = None + row["process_id"] = 12345 + row["error_severity"] = "LOG" + row["message"] = 'AUDIT: SESSION,1,1,READ,SELECT,,,"select 1",' + + with pytest.raises(ObjectLoggingError): + list(records_from_csvlog([dict(row), dict(row)])) + + class TestReadJsonlRecords: def test_a_chunked_record_numbered_from_one_is_reassembled(self): lines = [