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
5 changes: 5 additions & 0 deletions changelog.d/sr-1rgb.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
### Fixed

- `StatifierRouter.Addresses.reap/3` runs on SQLite: its stamp and delete
no longer use Postgres's `= ANY(...)`, which failed every reap of a
SQLite host on 0.8.0 and later, under an integer or a text key.
21 changes: 21 additions & 0 deletions docs/adr/0002-addressing.md
Original file line number Diff line number Diff line change
Expand Up @@ -1433,3 +1433,24 @@ creates the table, and it is run outside the version walk.

The 0.9.0 section of `CHANGELOG.md` names the key, the processor, the
front, the two calls and the opt-in table.

## Note (2026-09-30, sr-1rgb): the address sweep's writes run on SQLite

A Note, not an amendment: it decides nothing and changes no decision,
amendment or Note above it. The Note of 2026-09-27 accepting the sr-w58a
Amendment says the private `stamp/3` and `delete/2` of
`StatifierRouter.Addresses` bind the ids in `fragment("? = ANY(?)",
a.id, ^ids)`. That form is Postgres's own; on SQLite the statement fails
with "no such function: ANY", so a host on SQLite could not reap its
address table from statifier_router 0.8.0 on.

Both functions now name their rows in `fragment("? IN (?)", a.id,
splice(^batch))`: an `IN` list with one bound parameter per id, which
Postgres and SQLite both take, written in batches of at most 500 ids a
statement (the private `in_batches/2` and `@ids_per_statement`), under
SQLite's smallest limit on bound parameters. Each id is still bound as
the table handed it over, so the Amendment's rule that the package
never casts an id it binds holds as before, and a reap answers the same
counts on Postgres as it did. The tests
`StatifierRouter.SQLiteReapTest` reap on SQLite under the default
integer key, under a text key of digits, and past one batch.
36 changes: 25 additions & 11 deletions lib/statifier_router/addresses.ex
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,7 @@ defmodule StatifierRouter.Addresses do

@terminal [:completed, :failed, :cancelled]
@default_limit 1_000
@ids_per_statement 500

@typedoc """
What one reap did: how many rows it stamped `terminal_seen_at` on, how
Expand Down Expand Up @@ -266,27 +267,40 @@ defmodule StatifierRouter.Addresses do
DateTime.compare(DateTime.add(seen_at, horizon_ms, :millisecond), now) != :gt
end

defp stamp(_config, [], _now), do: 0

# 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
# @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.
defp stamp(config, ids, now) do
{count, _} =
in_batches(ids, fn batch ->
from(a in Config.queryable(config, Address),
where: fragment("? = ANY(?)", a.id, ^ids) and is_nil(a.terminal_seen_at)
where: fragment("? IN (?)", a.id, splice(^batch)) and is_nil(a.terminal_seen_at)
)
|> config.repo.update_all(set: [terminal_seen_at: now])

count
end)
end

defp delete(_config, []), do: 0

defp delete(config, ids) do
{count, _} =
in_batches(ids, fn batch ->
config.repo.delete_all(
from(a in Config.queryable(config, Address), where: fragment("? = ANY(?)", a.id, ^ids))
from(a in Config.queryable(config, Address),
where: fragment("? IN (?)", a.id, splice(^batch))
)
)
end)
end

count
# The rows `write` counts over every batch of `ids`; no statement for
# no ids.
defp in_batches(ids, write) do
ids
|> Enum.chunk_every(@ids_per_statement)
|> Enum.reduce(0, fn batch, total ->
{count, _} = write.(batch)
total + count
end)
end

defp next(rows, limit) when length(rows) == limit, do: List.last(rows).id
Expand Down
7 changes: 4 additions & 3 deletions test/statifier_router/sqlite_migrations_test.exs
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
defmodule StatifierRouter.SQLiteMigrationsTest do
# The version walk on SQLite, through ecto_sqlite3, against a database
# file of each test's own: no Postgres, no SQL sandbox, nothing shared
# with the rest of the suite.
use ExUnit.Case, async: true
# file of each test's own: no Postgres, no SQL sandbox. The one repo
# process runs under the module's name, so every module that starts it
# is in the :sqlite_repo group, whose modules never run at once.
use ExUnit.Case, async: true, group: :sqlite_repo

alias Ecto.Adapters.SQL
alias Ecto.Migrator
Expand Down
234 changes: 234 additions & 0 deletions test/statifier_router/sqlite_reap_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,234 @@
defmodule StatifierRouter.SQLiteReapTest do
# StatifierRouter.Addresses.reap/3 on SQLite, through ecto_sqlite3,
# against a database file of each test's own, under the default key and
# under a text key built with the :primary_key option. The execution
# statuses come from StatifierRouter.StatusStore: the repo holds only the
# router's tables. In the :sqlite_repo group with
# StatifierRouter.SQLiteMigrationsTest: both start the one repo process.
use ExUnit.Case, async: true, group: :sqlite_repo

import Ecto.Query, only: [from: 2]

alias Ecto.Migrator
alias StatifierPersistence.Storage
alias StatifierRouter.Addresses
alias StatifierRouter.Binding
alias StatifierRouter.Config
alias StatifierRouter.Schema.Address
alias StatifierRouter.SQLiteRepo

# A host on SQLite under the repo's default key.
defmodule MigrateIntegerKey do
@moduledoc false
use Ecto.Migration

def up, do: StatifierRouter.Migrations.up()
def down, do: StatifierRouter.Migrations.down()
end

# A host on SQLite under a text key the database fills in. The ids are
# made of digits alone and lead with a zero, so a reap that cast one to
# an integer would lose the zero and match no row of the text column.
defmodule MigrateTextKey do
@moduledoc false
use Ecto.Migration

def up do
StatifierRouter.Migrations.up(
primary_key: [
type: :text,
default: fragment("(printf('0%011d', abs(random()) % 100000000000))")
]
)
end

def down, do: StatifierRouter.Migrations.down()
end

@version 20_260_930_000_501

@now ~U[2026-09-30 12:00:00.000000Z]
@hour 3_600_000

setup do
database =
Path.join(
System.tmp_dir!(),
"statifier_router_sqlite_#{System.unique_integer([:positive])}.db"
)

start_supervised!({SQLiteRepo, database: database, pool_size: 1})

on_exit(fn ->
for suffix <- ["", "-wal", "-shm"], do: File.rm(database <> suffix)
end)

:ok
end

# Five parcels, one address row each: two delivered, one still on the
# van, one whose execution the store no longer holds, and one delivered
# and already stamped a day ago.
@statuses %{
"ex_pcl_6001" => :completed,
"ex_pcl_6002" => :completed,
"ex_pcl_6003" => :active,
"ex_pcl_6005" => :completed
}

defp config do
{:ok, config} =
Config.new(
repo: SQLiteRepo,
delivery: StatifierRouter.RecordingDelivery,
store: %Storage{adapter: StatifierRouter.StatusStore, opts: @statuses}
)

config
end

defp seed do
rows =
for n <- 6001..6005 do
%{
scope: "depot_north",
document: "parcel_delivery",
key: "pcl_#{n}",
execution_id: "ex_pcl_#{n}",
inserted_at: @now,
terminal_seen_at: if(n == 6005, do: DateTime.add(@now, -24 * @hour, :millisecond))
}
end

{5, _} = SQLiteRepo.insert_all(Address, rows)
:ok
end

# The reap's answer, or the message of what it raised, so a statement
# the adapter refuses fails the assertion that names it.
defp reap(config, bindings, opts) do
Addresses.reap(config, bindings, opts)
rescue
error in Exqlite.Error -> {:raised, Exception.message(error)}
end

# key => terminal_seen_at, for every address row left.
defp addresses(config) do
config
|> Config.queryable(Address)
|> then(&from(a in &1, select: {a.key, a.terminal_seen_at}))
|> SQLiteRepo.all()
|> Map.new()
end

defp ids(config) do
SQLiteRepo.all(from(a in Config.queryable(config, Address), select: a.id))
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

# The same reap, twice: the first stamps the two parcels it is the first
# to see delivered and deletes the orphan and the row stamped a day ago;
# the second, an hour on, deletes the two it stamped.
defp sweep(config) do
assert reap(config, bindings(), now: @now) ==
{:ok, %{stamped: 2, deleted: 2, next: nil}}

assert addresses(config) == %{
"pcl_6001" => @now,
"pcl_6002" => @now,
"pcl_6003" => nil
}

an_hour_on = DateTime.add(@now, @hour, :millisecond)

assert reap(config, bindings(), now: an_hour_on) ==
{:ok, %{stamped: 0, deleted: 2, next: nil}}

assert addresses(config) == %{"pcl_6003" => nil}
end

describe "reap/3 on SQLite" do
# sabotage: put back stamp/3's and delete/2's `? = ANY(?)` fragment ->
# red, the first reap raised "no such function: ANY".
test "stamps and deletes under the default integer key" do
:ok = Migrator.up(SQLiteRepo, @version, MigrateIntegerKey, log: false)
config = config()
:ok = seed()

assert Enum.all?(ids(config), &is_integer/1)
sweep(config)
end

# sabotage: cast every id to an integer before it was bound in stamp/3
# and delete/2 -> red, the first reap stamped and deleted nothing.
test "stamps and deletes under a text key of digits, bound uncast" do
:ok = Migrator.up(SQLiteRepo, @version, MigrateTextKey, log: false)
config = config()
:ok = seed()

assert Enum.all?(ids(config), &(is_binary(&1) and &1 =~ ~r/^0\d{11}$/))
sweep(config)
end
end

describe "reap/3 on SQLite, past one statement's batch of ids" do
# sabotage: in_batches/2 answered the last batch's count alone -> red,
# the first reap answered stamped: 201 for 1201 rows stamped.
test "stamps and deletes every row of a reap that examines more ids than one statement binds" do
:ok = Migrator.up(SQLiteRepo, @version, MigrateIntegerKey, log: false)
numbers = 7001..8201
statuses = Map.new(numbers, &{"ex_pcl_#{&1}", :completed})

{:ok, config} =
Config.new(
repo: SQLiteRepo,
delivery: StatifierRouter.RecordingDelivery,
store: %Storage{adapter: StatifierRouter.StatusStore, opts: statuses}
)

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

SQLiteRepo.insert_all(Address, rows)
end

assert reap(config, bindings(), now: @now, limit: 2_000) ==
{:ok, %{stamped: 1201, deleted: 0, next: nil}}

assert addresses(config) |> Map.values() |> Enum.uniq() == [@now]

an_hour_on = DateTime.add(@now, @hour, :millisecond)

assert reap(config, bindings(), now: an_hour_on, limit: 2_000) ==
{:ok, %{stamped: 0, deleted: 1201, next: nil}}

assert addresses(config) == %{}
end
end
end
23 changes: 23 additions & 0 deletions test/support/status_store.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
defmodule StatifierRouter.StatusStore do
@moduledoc """
A storage adapter that answers one read, `fetch_execution/2`, from a map
of execution id to status handed to it when the store is built:

%StatifierPersistence.Storage{adapter: __MODULE__, opts: %{"ex_1" => :completed}}

An id the map does not name answers `{:error, :execution_not_found}`.
It lets a test reap address rows in a repo that holds no execution
table, the SQLite repo of `StatifierRouter.SQLiteReapTest`. Test-only
support code, not part of the package's public API.
"""

@doc "The execution's status as the map names it, or `:execution_not_found`."
@spec fetch_execution(%{String.t() => atom()}, String.t()) ::
{:ok, %{status: atom()}} | {:error, :execution_not_found}
def fetch_execution(statuses, execution_id) do
case Map.fetch(statuses, execution_id) do
{:ok, status} -> {:ok, %{status: status}}
:error -> {:error, :execution_not_found}
end
end
end
Loading