From 2df4fc69dadcb1e18e10866ed227088c9d2acb2c Mon Sep 17 00:00:00 2001 From: JohnnyT Date: Fri, 2 Oct 2026 04:01:06 -0600 Subject: [PATCH] Binds the reap's ids as one array on Postgres StatifierRouter.Addresses.reap/3 named its rows in an IN list of one bound parameter per id on every adapter, so on Postgres the statement text varied with the batch's length and Postgres prepared one statement per distinct length where the array form had one. The private stamp/3 and delete/2 now branch on the repo's adapter: on Postgres each batch is one bound array, fragment("? = ANY(?)", a.id, ^batch), the same statement for every batch length; every other adapter, and a repo module that defines no __adapter__/0, keeps the spliced IN list, byte-identical to before. Ids are still bound uncast; Postgres types the array from the id column. Batches of 500 and every count and cursor are unchanged. Adds StatifierRouter.PostgresReapTest and a SQLiteReapTest case, each reading the statements from the repo's query telemetry event. ADR-0002 gets a dated foot Note; a changelog fragment records the change. Refs: sr-7zxe --- changelog.d/sr-7zxe.md | 6 + docs/adr/0002-addressing.md | 41 ++++++ lib/statifier_router/addresses.ex | 55 ++++++-- test/statifier_router/postgres_reap_test.exs | 131 +++++++++++++++++++ test/statifier_router/sqlite_reap_test.exs | 41 ++++++ 5 files changed, 261 insertions(+), 13 deletions(-) create mode 100644 changelog.d/sr-7zxe.md create mode 100644 test/statifier_router/postgres_reap_test.exs diff --git a/changelog.d/sr-7zxe.md b/changelog.d/sr-7zxe.md new file mode 100644 index 0000000..d09a2a9 --- /dev/null +++ b/changelog.d/sr-7zxe.md @@ -0,0 +1,6 @@ +### Changed + +- `StatifierRouter.Addresses.reap/3` binds each batch of ids as one array + on Postgres again (`= ANY(...)`), so Postgres prepares one statement per + write rather than one per batch length; SQLite and every other adapter + keep the `IN (...)` list, and every reap answers the same counts. diff --git a/docs/adr/0002-addressing.md b/docs/adr/0002-addressing.md index 3fe26c5..32b1671 100644 --- a/docs/adr/0002-addressing.md +++ b/docs/adr/0002-addressing.md @@ -1739,3 +1739,44 @@ and nothing else changes. - **The mitigation.** Rotate a location with `rotate_location/2` when it may have leaked, and never run a production repo or logger at `:debug`. + +## Note (2026-10-02, sr-7zxe): the address sweep binds one array per batch on Postgres again + +A Note, not an amendment: it decides nothing and changes no decision, +amendment or Note above it. The sr-1rgb Note above moved the private +`stamp/3` and `delete/2` of `StatifierRouter.Addresses` to an `IN` list +of one bound parameter per id on every adapter. On Postgres that +statement's text varies with the batch's length, so Postgres prepared a +statement for each distinct length where the array form of 0.8.0 had +one. Ruled by the operator, 2026-10-01: the array form on Postgres, the +`IN` list on every other adapter. + +- **The adapter branch.** The private `postgres?/1` reads the repo's + adapter from its `__adapter__/0`. Only `Ecto.Adapters.Postgres` takes + the array form; every other adapter takes the `IN` list, and so does a + repo module that defines no `__adapter__/0` (one that delegates to an + Ecto repo rather than being one). +- **The two forms.** On Postgres the private `stamp_query/2` and + `delete_query/2` name the rows in `fragment("? = ANY(?)", a.id, + ^batch)`: one bound parameter for the whole batch, so a write is the + same statement whatever the batch's length. Elsewhere they name them in + `fragment("? IN (?)", a.id, splice(^batch))`, the sr-1rgb form, + unchanged. +- **The array's type.** The ids are bound as the table handed them over, + never cast. Postgres types the parameter as an array of the id + column's own type (`bigint[]` under the default key, `text[]` under a + text key), and Postgrex encodes the list as that type, so a text id + made of digits alone stays a string, as the sr-w58a Amendment's rule + requires. +- **What does not change.** The batches stay at most 500 ids a + statement on every adapter (the private `in_batches/2` and + `@ids_per_statement`), and `reap/3` answers the same `stamped`, + `deleted` and `next` as before. No option, table or answer changes. +- **The tests.** On Postgres, "binds each batch of ids as one array, so + every batch of a write is one statement" in + `StatifierRouter.PostgresReapTest` reads the statements from the + repo's query telemetry event over a reap of three batches; on SQLite, + "binds each batch of ids as a spliced IN list" in + `StatifierRouter.SQLiteReapTest` does the same. "never casts a text id + made of digits alone" in `StatifierRouter.PrimaryKeyTest` reaps under a + text key of digits on Postgres. diff --git a/lib/statifier_router/addresses.ex b/lib/statifier_router/addresses.ex index cfc2043..a6dca4f 100644 --- a/lib/statifier_router/addresses.ex +++ b/lib/statifier_router/addresses.ex @@ -267,29 +267,58 @@ defmodule StatifierRouter.Addresses do DateTime.compare(DateTime.add(seen_at, horizon_ms, :millisecond), now) != :gt end - # Both writes name their rows in an IN list of one bound parameter per - # id, which Postgres and SQLite both take, each id bound uncast as the - # rest of this module binds them. A list longer than + # Both writes name their rows by id, each id bound uncast as the rest of + # this module binds them, in one of two forms chosen by the repo's + # adapter. On Postgres the batch is one bound array, `? = ANY(?)`, so a + # write is the same statement whatever the batch's length and Postgres + # prepares it once; the server types the array from the id column + # (bigint[] under the default key, text[] under a text key), and Postgrex + # encodes the ids as that. Every other adapter takes an IN list of one + # bound parameter per id, `? IN (?)` with the batch spliced, since `ANY` + # is Postgres's own (SQLite has no such function). A list longer than # @ids_per_statement is written in batches of that many, one statement - # each, which keeps every statement under SQLite's smallest limit on - # bound parameters (999) with room for the stamp's own time. + # each, on every adapter, which keeps every IN list under SQLite's + # smallest limit on bound parameters (999) with room for the stamp's own + # time. defp stamp(config, ids, now) do in_batches(ids, fn batch -> + config + |> stamp_query(batch) + |> config.repo.update_all(set: [terminal_seen_at: now]) + end) + end + + defp stamp_query(config, batch) do + if postgres?(config) do + from(a in Config.queryable(config, Address), + where: fragment("? = ANY(?)", a.id, ^batch) and is_nil(a.terminal_seen_at) + ) + else from(a in Config.queryable(config, Address), where: fragment("? IN (?)", a.id, splice(^batch)) and is_nil(a.terminal_seen_at) ) - |> config.repo.update_all(set: [terminal_seen_at: now]) - end) + end end defp delete(config, ids) do - in_batches(ids, fn batch -> - config.repo.delete_all( - from(a in Config.queryable(config, Address), - where: fragment("? IN (?)", a.id, splice(^batch)) - ) + in_batches(ids, fn batch -> config.repo.delete_all(delete_query(config, batch)) end) + end + + defp delete_query(config, batch) do + if postgres?(config) do + from(a in Config.queryable(config, Address), where: fragment("? = ANY(?)", a.id, ^batch)) + else + from(a in Config.queryable(config, Address), + where: fragment("? IN (?)", a.id, splice(^batch)) ) - end) + end + end + + # A repo module that names no adapter (one that delegates to an Ecto + # repo rather than being one) takes the IN list, as before the array + # form came back. The module is loaded: examine/3 has called it. + defp postgres?(%Config{repo: repo}) do + function_exported?(repo, :__adapter__, 0) and repo.__adapter__() == Ecto.Adapters.Postgres end # The rows `write` counts over every batch of `ids`; no statement for diff --git a/test/statifier_router/postgres_reap_test.exs b/test/statifier_router/postgres_reap_test.exs new file mode 100644 index 0000000..1af3e44 --- /dev/null +++ b/test/statifier_router/postgres_reap_test.exs @@ -0,0 +1,131 @@ +defmodule StatifierRouter.PostgresReapTest do + # The statements StatifierRouter.Addresses.reap/3 sends on Postgres, read + # from the repo's query telemetry event: each batch of ids is bound as + # one array, so every batch of a write is the same statement whatever + # its length. The SQLite half of the same branch is in + # StatifierRouter.SQLiteReapTest. The execution statuses come from + # StatifierRouter.StatusStore, so the reap reads no execution table. + use ExUnit.Case, async: true, group: :database + + alias Ecto.Adapters.SQL.Sandbox + alias StatifierPersistence.Storage + alias StatifierRouter.Addresses + alias StatifierRouter.Binding + alias StatifierRouter.Config + alias StatifierRouter.Schema.Address + alias StatifierRouter.TestRepo + + @now ~U[2026-10-02 12:00:00.000000Z] + @hour 3_600_000 + @numbers 7001..8201 + + setup do + :ok = Sandbox.checkout(TestRepo) + + handler = "postgres-reap-#{System.unique_integer([:positive])}" + query = TestRepo.config()[:telemetry_prefix] ++ [:query] + :ok = :telemetry.attach(handler, query, &__MODULE__.handle_query/4, self()) + on_exit(fn -> :telemetry.detach(handler) end) + + :ok + end + + @doc false + # Runs in the process that sent the query, so a statement this test's + # own process sent is the only one it is told about. + def handle_query(_event, _measurements, %{query: query}, test) do + if self() == test, do: send(test, {:query, query}) + end + + # Every UPDATE or DELETE statement this test sent, in order. + defp writes do + receive do + {:query, "UPDATE " <> _ = query} -> [query | writes()] + {:query, "DELETE " <> _ = query} -> [query | writes()] + {:query, _other} -> writes() + after + 0 -> [] + end + end + + # The delivered-scan binding, whose horizon keeps a delivered parcel's + # row for an hour after a reap first sees it delivered. + defp bindings do + {:ok, binding} = + Binding.new(%{ + id: "delivered_scans", + source: "depot_scans", + match: "event.kind == 'delivered'", + key: "event.parcel_id", + document: "parcel_delivery", + event: "delivered", + data: ["parcel_id"], + dedupe: %{by: :message_id, horizon_ms: @hour} + }) + + [binding] + end + + # One address row per delivered parcel, more rows than one statement + # binds, so a reap writes them in batches of two lengths. + defp seed do + for chunk <- Enum.chunk_every(@numbers, 100) do + rows = + for n <- chunk do + %{ + scope: "depot_north", + document: "parcel_delivery", + key: "pcl_#{n}", + execution_id: "ex_pcl_#{n}", + inserted_at: @now + } + end + + TestRepo.insert_all(Address, rows) + end + + {:ok, config} = + Config.new( + repo: TestRepo, + delivery: StatifierRouter.RecordingDelivery, + store: %Storage{ + adapter: StatifierRouter.StatusStore, + opts: Map.new(@numbers, &{"ex_pcl_#{&1}", :completed}) + } + ) + + config + end + + describe "reap/3 on Postgres" do + # sabotage: made postgres?/1 answer false, so stamp/3 and delete/2 took + # the spliced IN form on Postgres -> red, the three stamp batches sent + # two statement texts (a 500-parameter IN list and a 201-parameter + # one), not one; restored, green. + test "binds each batch of ids as one array, so every batch of a write is one statement" do + config = seed() + + assert Addresses.reap(config, bindings(), now: @now, limit: 2_000) == + {:ok, %{stamped: 1201, deleted: 0, next: nil}} + + stamps = writes() + assert length(stamps) == 3 + assert [stamp] = Enum.uniq(stamps) + assert stamp =~ "= ANY(" + refute stamp =~ " IN (" + + an_hour_on = DateTime.add(@now, @hour, :millisecond) + + assert Addresses.reap(config, bindings(), now: an_hour_on, limit: 2_000) == + {:ok, %{stamped: 0, deleted: 1201, next: nil}} + + deletes = writes() + assert length(deletes) == 3 + assert [delete] = Enum.uniq(deletes) + assert delete =~ "= ANY(" + refute delete =~ " IN (" + + assert TestRepo.aggregate(Config.queryable(config, Address), :count) == 0 + end + end +end diff --git a/test/statifier_router/sqlite_reap_test.exs b/test/statifier_router/sqlite_reap_test.exs index 491a9a5..9ffa407 100644 --- a/test/statifier_router/sqlite_reap_test.exs +++ b/test/statifier_router/sqlite_reap_test.exs @@ -121,6 +121,24 @@ defmodule StatifierRouter.SQLiteReapTest do |> Map.new() end + @doc false + # Runs in the process that sent the query, so a statement this test's + # own process sent is the only one it is told about. + def handle_query(_event, _measurements, %{query: query}, test) do + if self() == test, do: send(test, {:query, query}) + end + + # Every UPDATE or DELETE statement this test sent, in order. + defp writes do + receive do + {:query, "UPDATE " <> _ = query} -> [query | writes()] + {:query, "DELETE " <> _ = query} -> [query | writes()] + {:query, _other} -> writes() + after + 0 -> [] + end + end + defp ids(config) do SQLiteRepo.all(from(a in Config.queryable(config, Address), select: a.id)) end @@ -186,6 +204,29 @@ defmodule StatifierRouter.SQLiteReapTest do assert Enum.all?(ids(config), &(is_binary(&1) and &1 =~ ~r/^0\d{11}$/)) sweep(config) end + + # sabotage: made postgres?/1 answer true, so stamp/3 and delete/2 took + # the `? = ANY(?)` array form on SQLite -> red, the first reap raised + # "no such function: ANY" instead of stamping; restored, green. + test "binds each batch of ids as a spliced IN list" do + :ok = Migrator.up(SQLiteRepo, @version, MigrateIntegerKey, log: false) + config = config() + :ok = seed() + + handler = "sqlite-reap-#{System.unique_integer([:positive])}" + query = SQLiteRepo.config()[:telemetry_prefix] ++ [:query] + :ok = :telemetry.attach(handler, query, &__MODULE__.handle_query/4, self()) + on_exit(fn -> :telemetry.detach(handler) end) + + sweep(config) + + assert [_stamp, _delete_now, _delete_later] = statements = writes() + + for statement <- statements do + assert statement =~ " IN (" + refute statement =~ "ANY(" + end + end end describe "reap/3 on SQLite, past one statement's batch of ids" do