Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions changelog.d/sr-7zxe.md
Original file line number Diff line number Diff line change
@@ -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.
41 changes: 41 additions & 0 deletions docs/adr/0002-addressing.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
55 changes: 42 additions & 13 deletions lib/statifier_router/addresses.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
131 changes: 131 additions & 0 deletions test/statifier_router/postgres_reap_test.exs
Original file line number Diff line number Diff line change
@@ -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
41 changes: 41 additions & 0 deletions test/statifier_router/sqlite_reap_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading