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
12 changes: 12 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1017,9 +1017,11 @@ jobs:
"

- name: Setup database
id: scale_setup_db
run: mix ecto.create && mix ecto.migrate

- name: Grant RLS roles access to tables
id: scale_grant
run: |
PGPASSWORD=postgres psql -h localhost -U postgres -d loopctl_test -c "
GRANT ALL PRIVILEGES ON ALL TABLES IN SCHEMA public TO loopctl_app;
Expand All @@ -1029,6 +1031,15 @@ jobs:
- name: Run scale tests (TC-27.1.1 through TC-27.1.4, AC-27.1.7)
run: SCALE_TESTS=true mix test --only scale test/loopctl/knowledge/scale_seed_test.exs

# US-32.1 AC-32.1.2: the revoke sweep's natural plan over a committed, ANALYZEd,
# production-shaped dispatches table. Here and not in the default suite because an
# ANALYZE writes reltuples in place and would skew every other test's plans there.
- name: Run revoke-sweep plan guard (AC-32.1.2)
# Independent of the gates around it, so one gate's failure does not hide another's,
# but never on a database that did not finish migrating and granting.
if: ${{ !cancelled() && steps.scale_setup_db.outcome == 'success' && steps.scale_grant.outcome == 'success' }}
run: SCALE_TESTS=true mix test --only scale test/loopctl/workers/revoke_expired_dispatches_plan_scale_test.exs

# US-41.1 AC-41.1.12(i) — the per-dimension ANN plan GATE. It EXPLAINs the real
# request-path inner ANN (BOTH article search and agent-memory recall — each has
# its own per-dimension index and its own verbatim `::vector(N)` cast) over a
Expand All @@ -1044,6 +1055,7 @@ jobs:
# This step is the gate; AC-41.1.12(ii)'s production EXPLAIN is an artifact report,
# not CI.
- name: Run per-dimension ANN plan gate (TC-41.1.2, AC-41.1.12(i))
if: ${{ !cancelled() && steps.scale_setup_db.outcome == 'success' && steps.scale_grant.outcome == 'success' }}
run: |
SCALE_TESTS=true mix test --only scale \
test/loopctl/knowledge/embedding_dimension_plan_scale_test.exs
Expand Down
8 changes: 4 additions & 4 deletions docs/user_stories/epic_32_scale_quick_wins/us_32.1.json
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,10 @@
},
"rationale": "The revoke sweep is unavoidably tenant-agnostic (it must find expired dispatches across all tenants), so it can never SEEK the tenant-leading composite index (at best the planner reads it whole). As the dispatches table grows with every agent dispatch, that full read, or a seq scan, becomes a steadily worsening fixed cost on the shared primary. A partial index matching the exact predicate makes the sweep O(expired rows), not O(table).",
"acceptance_criteria": [
{ "id": "AC-32.1.1", "description": "A new migration adds `create index(:dispatches, [:expires_at], where: \"revoked_at IS NULL\")`. Because dispatches is a hot, actively-written table, the index is built online: `@disable_ddl_transaction true`, `@disable_migration_lock true`, and `create index(..., concurrently: true)`, preceded by `drop index(..., if_exists: true, concurrently: true)` so a prior failed/invalid index is rebuilt rather than skipped-by-name.", "status": "pending" },
{ "id": "AC-32.1.2", "description": "After migration, `EXPLAIN` of the worker's query (`WHERE revoked_at IS NULL AND expires_at < now()`) shows an Index Scan / Index-Only Scan using the new partial index, not a Seq Scan.", "status": "pending" },
{ "id": "AC-32.1.3", "description": "RevokeExpiredDispatchesWorker behavior is UNCHANGED — the story adds only an index; the query, revoke semantics, and audit/cascade to api_keys are untouched.", "status": "pending" },
{ "id": "AC-32.1.4", "description": "The migration follows loopctl RLS conventions: dispatches already has RLS enabled; adding an index does not alter RLS. No BYPASSRLS/ownership change.", "status": "pending" }
{ "id": "AC-32.1.1", "description": "A new migration adds `create index(:dispatches, [:expires_at], where: \"revoked_at IS NULL\")`. Because dispatches is a hot, actively-written table, the index is built online: `@disable_ddl_transaction true`, `@disable_migration_lock true`, and `create index(..., concurrently: true)`, preceded by `drop index(..., if_exists: true, concurrently: true)` so a prior failed/invalid index is rebuilt rather than skipped-by-name.", "status": "complete" },
{ "id": "AC-32.1.2", "description": "After migration, `EXPLAIN` of the worker's query (`WHERE revoked_at IS NULL AND expires_at < now()`) uses the new partial index (an Index Scan or a Bitmap Index Scan on it) and not the tenant-leading composite index, and no Seq Scan. Asserted in CI's scale job by `test/loopctl/workers/revoke_expired_dispatches_plan_scale_test.exs`, on a committed, ANALYZEd, production-shaped table, against the query the worker issues.", "status": "complete" },
{ "id": "AC-32.1.3", "description": "RevokeExpiredDispatchesWorker behavior is UNCHANGED — the story adds only an index; the query, revoke semantics, and audit/cascade to api_keys are untouched.", "status": "complete" },
{ "id": "AC-32.1.4", "description": "The migration follows loopctl RLS conventions: dispatches already has RLS enabled; adding an index does not alter RLS. No BYPASSRLS/ownership change.", "status": "complete" }
],
"test_cases": [
{ "id": "TC-32.1.1", "name": "revoke sweep still revokes expired non-revoked dispatches", "type": "integration",
Expand Down
2 changes: 1 addition & 1 deletion docs/user_stories/epic_35_sth_cron_redesign/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ Issue #350 predates recent work; **3 of its 6 findings are already done** and ar
|--------------|--------|----------|
| Fanout uses N single-row inserts instead of `insert_all` | ✅ done | US-32.5 — `compute_sth_worker.ex:59-68` uses `Oban.insert_all/1`, chunked at 5000 |
| Zero-jitter thundering herd at the `:00` tick | ✅ done (jitter part) | US-32.5 — `schedule_in: jitter(tenant_id)`, deterministic `phash2` over a 56s window (`compute_sth_worker.ex:45,57,113`) |
| `RevokeExpiredDispatchesWorker` query can't use its index | ✅ done | US-32.1 — `20260713000000_add_dispatches_expires_at_active_index.exs` adds `CREATE INDEX CONCURRENTLY … ON dispatches (expires_at) WHERE revoked_at IS NULL`; `revoke_expired_dispatches_worker_test.exs` asserts the index's shape and predicate from `pg_index`; the planner's choice is not asserted in the async suite |
| `RevokeExpiredDispatchesWorker` query can't use its index | ✅ done | US-32.1 — `20260713000000_add_dispatches_expires_at_active_index.exs` adds `CREATE INDEX CONCURRENTLY … ON dispatches (expires_at) WHERE revoked_at IS NULL`; `revoke_expired_dispatches_worker_test.exs` asserts the index's shape and predicate from `pg_index`; `revoke_expired_dispatches_plan_scale_test.exs` (CI scale job) asserts the worker's query is planned through it (Index Scan or Bitmap Index Scan) and not through the tenant-leading composite |
| `dispatches` missing partial index `(expires_at) WHERE revoked_at IS NULL` | ✅ done | Same migration — exact predicate match, not a different one |

The **genuinely-remaining** work is the three findings the audit rated as the
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,13 @@ defmodule Loopctl.Repo.Migrations.AddDispatchesExpiresAtActiveIndex do
# (other per-tenant lookups still use it).
# VERIFICATION (AC-32.1.2) — EXPLAIN of the worker's exact predicate
# (`WHERE revoked_at IS NULL AND expires_at < now()`) uses this partial index,
# not a Seq Scan. Two captures, each taken once by hand when this migration was
# written; neither is asserted in the suite. The ExUnit guard
# `RevokeExpiredDispatchesWorkerTest` asserts the index's shape and predicate from
# `pg_index`, because the planner's choice in the shared test table moves with the rows
# concurrent tests have in flight.
# not a Seq Scan. Two captures, both taken by hand when this migration was written and
# neither re-run by CI. What CI does assert, in `RevokeExpiredDispatchesPlanScaleTest`
# (scale job), is its own shape: the query the worker issues, with a bound timestamp
# parameter, over a committed, ANALYZEd table of mostly revoked history with a small live
# set and expired backlog, is planned through this index (an Index Scan, or a Bitmap
# Index Scan on it). The default suite's `RevokeExpiredDispatchesWorkerTest` asserts only
# the index's shape and predicate from `pg_index`.
#
# * Empty table, planner forced to reveal usability (`SET LOCAL
# enable_seqscan = off`), which proves eligibility rather than choice:
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
defmodule Loopctl.Repo.Migrations.IndexForeignKeysIntoDispatches do
use Ecto.Migration

# Two foreign keys into `dispatches` had no index on the referencing column:
# `dispatches.parent_dispatch_id` (its own self-reference) and
# `story_acceptance_criteria.verified_by_dispatch_id`. Deleting a dispatch runs one
# referential check per deleted row against each referencing table, and without an index
# each check is a scan of that table, so deleting N dispatches costs N scans. Measured
# 2026-09-28: deleting 20k seeded dispatches did not finish inside a test's exit. Every
# path that deletes dispatches in bulk pays it, tenant teardown included.
#
# Partial on IS NOT NULL: the checks look up a non-null id, most dispatches have no parent,
# and most criteria were never verified by a dispatch, so the NULL rows would only make
# the indexes bigger.
#
# CONCURRENTLY, with the DDL transaction and migration lock off, because `dispatches` is a
# hot write path. An interrupted CONCURRENTLY build leaves an INVALID index that
# `IF NOT EXISTS` would match by name and skip, so each index is dropped first, but ONLY
# when it exists invalid or in another shape (`stale?/2`, the same reconciliation as
# 20260919100000_add_runner_unsupported_kind_index). Dropping unconditionally would, on a
# re-run after the second build failed, drop the first index while it was valid and leave
# the hot table without it for a whole rebuild.
@disable_ddl_transaction true
@disable_migration_lock true

@indexes [
{"dispatches_parent_dispatch_id_index", "dispatches", "parent_dispatch_id"},
{"story_acceptance_criteria_verified_by_dispatch_id_index", "story_acceptance_criteria",
"verified_by_dispatch_id"}
]

def up do
for {name, table, column} <- @indexes do
if stale?(name, shape(table, column)),
do: execute("DROP INDEX CONCURRENTLY IF EXISTS #{name}")

execute("""
CREATE INDEX CONCURRENTLY IF NOT EXISTS #{name}
ON #{table} (#{column})
WHERE #{column} IS NOT NULL
""")
end
end

def down do
for {name, _table, _column} <- @indexes do
execute("DROP INDEX CONCURRENTLY IF EXISTS #{name}")
end
end

defp shape(table, column) do
~r/^CREATE INDEX \S+ ON public\.#{table} USING btree \(#{column}\) WHERE \(#{column} IS NOT NULL\)$/
end

# True when an index of this name exists but is INVALID or not the shape above; false
# when it is absent (the CREATE builds it) or already valid in this shape (nothing to do).
defp stale?(name, shape) do
sql = """
SELECT pg_get_indexdef(c.oid), x.indisvalid
FROM pg_class c
JOIN pg_index x ON x.indexrelid = c.oid
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relname = $1 AND c.relkind = 'i' AND n.nspname = 'public'
"""

case repo().query!(sql, [name]).rows do
[[indexdef, true]] -> not (indexdef =~ shape)
rows -> rows != []
end
end
end
48 changes: 48 additions & 0 deletions test/loopctl/dispatches/foreign_key_indexes_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
defmodule Loopctl.Dispatches.ForeignKeyIndexesTest do
@moduledoc """
Every foreign key that references `dispatches` has a valid btree index whose leading KEY
columns (INCLUDE columns do not count) are exactly the key's referencing columns, and
which is either unconditional or conditioned only on one of those columns being NOT NULL,
the one predicate every referential lookup on that key satisfies.

Deleting a dispatch runs one referential check per deleted row against each referencing
table; without that index each check scans the table, so a bulk delete of N dispatches
costs N scans. 20260928140000_index_foreign_keys_into_dispatches added the two that were
missing. This reads the catalog, so a new foreign key into `dispatches` without its index
fails here rather than in the first bulk delete that meets it.
"""

use Loopctl.DataCase, async: true

test "every foreign key referencing dispatches is indexed on its referencing column" do
%{rows: rows} =
Loopctl.AdminRepo.query!("""
SELECT c.conrelid::regclass::text, c.conname,
EXISTS (
SELECT 1
FROM pg_index i
JOIN pg_class ic ON ic.oid = i.indexrelid
JOIN pg_am am ON am.oid = ic.relam,
LATERAL (SELECT (string_to_array(i.indkey::text, ' ')::int2[])
[1:least(i.indnkeyatts, array_length(c.conkey, 1))] AS lead) k
WHERE i.indrelid = c.conrelid
AND i.indisvalid
AND am.amname = 'btree'
AND k.lead @> c.conkey AND k.lead <@ c.conkey
AND (i.indpred IS NULL
OR pg_get_expr(i.indpred, i.indrelid) IN (
SELECT '(' || quote_ident(a.attname) || ' IS NOT NULL)'
FROM pg_attribute a
WHERE a.attrelid = c.conrelid AND a.attnum = ANY (c.conkey)
))
)
FROM pg_constraint c
WHERE c.confrelid = 'dispatches'::regclass AND c.contype = 'f'
""")

assert rows != [], "found no foreign keys into dispatches; the catalog query is wrong"

unindexed = for [table, constraint, false] <- rows, do: "#{table} #{constraint}"
assert unindexed == []
end
end
130 changes: 130 additions & 0 deletions test/loopctl/workers/revoke_expired_dispatches_plan_scale_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
defmodule Loopctl.Workers.RevokeExpiredDispatchesPlanScaleTest do
@moduledoc """
Epic 32, US-32.1, AC-32.1.2: the query `RevokeExpiredDispatchesWorker` actually issues is
served by the partial index `dispatches_expires_at_active_index`, chosen by the DEFAULT
planner (no `enable_seqscan = off`), and is the ONLY index in the plan, so a bitmap that
also reads the tenant-leading composite or any other index whole fails. The query reads one
table, so a plan that uses the index cannot also Seq Scan it. With a small live set the
planner reaches the index by an Index Scan, with a larger one by a Bitmap Index Scan; both
seek it.

The seed is the shape a table that has run for a while settles into, since every swept
dispatch stays in it past expiry: a large revoked history, a small live set and a small
expired backlog the next sweep will take. Measured 2026-09-28 on Postgres 16, the index is
used at every live set tried, from 50 rows to half the table.

The query is captured from the worker's own `perform/1` through repo telemetry, not
rebuilt here, so a WHERE clause that stops implying the index predicate
`revoked_at IS NULL`, or stops being a range on `expires_at`, turns this red. The sweep
runs inside a transaction this test rolls back: it is cross-tenant, and in a database other
tests have used it would otherwise revoke their expired dispatches too.

Synchronous and outside the sandbox like every `:scale` test, which is this repo's
standing exception to its async rule: the seed has to be committed for a separate VACUUM
ANALYZE to count it. The ANALYZE writes `pg_class.reltuples`/`relpages` in place, so run it
in a database no other suite is using at the same time; CI runs it in the scale job's own:

SCALE_TESTS=true mix test --only scale test/loopctl/workers/revoke_expired_dispatches_plan_scale_test.exs
"""

use ExUnit.Case, async: false

import Ecto.Query
import Loopctl.Fixtures
import Mox

alias Ecto.Adapters.SQL.Sandbox
alias Loopctl.AdminRepo
alias Loopctl.Dispatches.Dispatch
alias Loopctl.PlanAssertions
alias Loopctl.Tenants.Tenant
alias Loopctl.Workers.RevokeExpiredDispatchesWorker

@moduletag :scale

@index "dispatches_expires_at_active_index"
@revoked_history 20_000
@live 50
@backlog 20

setup :verify_on_exit!

defp unboxed(fun), do: Sandbox.unboxed_run(AdminRepo, fun)

setup do
Loopctl.DataCase.stub_all_defaults()

tenant =
unboxed(fn -> fixture(:tenant, slug: "revoke-plan-scale-#{Ecto.UUID.generate()}") end)

# Registered before anything is seeded, so a seed that fails part-way is still removed.
on_exit(fn ->
unboxed(fn ->
AdminRepo.delete_all(from(d in Dispatch, where: d.tenant_id == ^tenant.id))
AdminRepo.delete_all(from(t in Tenant, where: t.id == ^tenant.id))
AdminRepo.query!("VACUUM ANALYZE dispatches")
end)
end)

unboxed(fn -> seed!(tenant.id) end)
{:ok, tenant: tenant}
end

defp seed!(tenant_id) do
fixture(:dispatch_sweep_history,
tenant_id: tenant_id,
revoked: @revoked_history,
live: @live,
backlog: @backlog
)

# Plan against the table's clean state, which is what CI's fresh database has. A
# database this test has already run in carries its earlier seeds' deleted rows: as dead
# heap tuples until a VACUUM, and as index pages a VACUUM empties but never returns.
# Measured 2026-09-28 after repeated local runs: the partial index held 70 entries in
# 131 pages, which priced its scan at 540 against about 590 for a Seq Scan, and the
# planner went either way from run to run. REINDEX rebuilds it at its real size;
# CONCURRENTLY, so writers elsewhere in the database are not blocked behind it. An
# interrupted concurrent rebuild leaves an INVALID `<index>_ccnew*` behind that every
# later insert keeps maintaining, so any such leftover is dropped first.
%{rows: leftovers} =
AdminRepo.query!(
"SELECT c.relname FROM pg_class c JOIN pg_index i ON i.indexrelid = c.oid " <>
"WHERE i.indrelid = 'public.dispatches'::regclass AND NOT i.indisvalid " <>
"AND c.relname ~ ('^' || $1 || '_cc(new|old)[0-9]*$')",
[@index]
)

for [name] <- leftovers, do: AdminRepo.query!("DROP INDEX CONCURRENTLY IF EXISTS #{name}")

AdminRepo.query!("REINDEX INDEX CONCURRENTLY #{@index}")
AdminRepo.query!("VACUUM ANALYZE dispatches")
end

test "the worker's sweep is served by the partial index", %{tenant: tenant} do
unboxed(fn ->
captured =
PlanAssertions.capture_repo_queries(fn ->
{:error, :rolled_back} =
AdminRepo.transaction(fn ->
assert :ok = RevokeExpiredDispatchesWorker.perform(%Oban.Job{args: %{}})

# perform/1 answers :ok whatever its write did, so check the write: the backlog,
# and only the backlog, of this tenant's dispatches is revoked.
# Raw, unquoted SQL so the capture below cannot mistake it for the sweep.
assert %{rows: [[@live]]} =
AdminRepo.query!(
"select count(*) from dispatches where tenant_id = $1 and revoked_at is null",
[Ecto.UUID.dump!(tenant.id)]
)

AdminRepo.rollback(:rolled_back)
end)
end)

sweep = PlanAssertions.only_query_matching(captured, ~r/\ASELECT .* FROM "dispatches"/s)

PlanAssertions.assert_only_index_used(sweep, "dispatches", @index, "revoked_at")
end)
end
end
Loading
Loading