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
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,15 @@ The two libraries are released separately, one tag per language: `python-v*` to
not make them agree. What does is the conformance suite and
[`conformance/contract-version.txt`](conformance/contract-version.txt).

## Unreleased

**Fixed in `smart-data-engine-sdk`.** Resume and abandonment of a PostgreSQL in-place index build
used to read the catalogue while the interrupted statement was still running. A deadline or a
killed operator ends the client, not its `CREATE INDEX CONCURRENTLY`. Recovery then either raced
that statement's commit and failed with "tuple concurrently updated", or dropped the index it was
finishing. Both now wait for the statement to end, then keep what it finished
([`docs/in-place-index.md`](docs/in-place-index.md)).

## `smart-data-engine-sdk` 0.1.0 and `@smart-data-engines/sde` 0.1.0

Published on 27 September 2026, each tag on its first run of the release workflow, with attestations
Expand Down
29 changes: 25 additions & 4 deletions docs/in-place-index.md
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,8 @@ declared would publish a map about another table. A refusal at this point leaves
index of the bound name on the bound table, of the declared method and columns, not unique, without
predicate, expression, INCLUDE columns or sort options, is this build's own: if it is valid and
ready it is done; if it is not - what an interrupted concurrent build leaves - it is dropped with
`DROP INDEX CONCURRENTLY` and built again. Anything else under the name is somebody else's object
`DROP INDEX CONCURRENTLY` and built again, once no statement that names it is still running (see
the deadline below). Anything else under the name is somebody else's object
and the build refuses it. A concurrent build waits for every transaction with an older snapshot, so
a long transaction anywhere in the database - an idle session left in a transaction included -
holds it until that transaction ends or the build budget does.
Expand Down Expand Up @@ -166,14 +167,32 @@ removal.
The deadline is the signed build budget, not the 30-second watchdog of staging and cutover. Nothing
is paused while an index builds, so the budget is not a pause budget; it bounds a build on a server
that stopped answering. When it ends a build the operator closes its connections and stops; resume
or abandon with fresh ones. A PostgreSQL build interrupted this way may run on in the server until
it notices; recovery waits for it, and keeps what it finished or drops what it left unfinished.
or abandon with fresh ones.

**A PostgreSQL build interrupted this way runs on in the server**, and so does the build of an
operator that was killed. The server notices a vanished client only when it next talks to it. For
a concurrent build that is once the index is valid (measured, PostgreSQL 15.19). So resume and
abandonment first wait until no active statement names the index in `pg_stat_activity`, within
the same budget, and only then read the catalogue. What the interrupted statement finished is kept,
and what it left unfinished is dropped and built again.

Reading first was not only wasteful. The build releases the table's session lock before its last
transaction commits, and a `DROP INDEX CONCURRENTLY` waiting for that lock started in the window:
"tuple concurrently updated", in the SDK's CI on 27 September 2026.

`pg_stat_activity` shows a statement's text to its own login, which is the one an earlier run of
the operator used. If recovery runs as another login, it cannot see the statement and cannot wait
for it.

A removal needs no such wait. Two concurrent drops of one index are serialized by the exclusive
lock the first holds on the index until it commits, and the second ends on `IF EXISTS`.

## Abandonment

`LocalCutover.abandon()` / `sde-operator abandon` ends an unfinished build without publishing
anything, as long as its decision is not yet `built`. It records the decision `abandoned` first,
then removes this build's own indexes: PostgreSQL with `DROP INDEX CONCURRENTLY`; ClickHouse by
then removes this build's own indexes: PostgreSQL with `DROP INDEX CONCURRENTLY`, once any
interrupted build of the index has ended; ClickHouse by
killing the materialization if it is still running and dropping the index with
`SETTINGS alter_sync = 0`, which leaves the catalogue at once instead of waiting for the merge pool
(measured: the default form waits for as long as merges are stopped). Objects under the bound names
Expand Down Expand Up @@ -231,6 +250,8 @@ abandoned one returns the abandonment.
engine itself while the application writes, then published; kept indexes and several new ones;
recovery after every durable step and after a killed operator process mid-build (PostgreSQL leaves
a leftover that is dropped and rebuilt, ClickHouse's mutation is found again, not repeated);
recovery and abandonment beside a killed operator's PostgreSQL build that still runs in the server,
which send no drop until it has ended and keep the index it finished;
abandonment before the decision and its refusal after; an interrupted abandonment finished by
resume or by a second abandonment; foreign objects under the bound name refused before any DDL and
left alone by abandonment; another operation's barrier refused before and during a build; the
Expand Down
32 changes: 30 additions & 2 deletions python/src/sde/engines/_index_build.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,8 +66,33 @@ def inspect(self, table: TableIdentity, index: Mapping[str, Any]) -> tuple[Statu
def status(self, table: TableIdentity, index: Mapping[str, Any]) -> Status:
return self.inspect(table, index)[0]

def settle(self, table: TableIdentity, index: Mapping[str, Any]) -> None:
"""Wait until no statement still running in the server names this build's index.

On PostgreSQL a deadline or a killed operator ends the client, not its statement: the
server notices a vanished client only when it next talks to it, which for a concurrent
build is once the index is valid (measured, PostgreSQL 15.19). The build also releases the
table's session lock before its last transaction commits, so a drop waiting for that lock
started in the window and failed with "tuple concurrently updated". Read after the wait,
the catalogue says what that statement finished, which is kept, or left, which is dropped.

pg_stat_activity shows a statement's text to its own login, which is the one an earlier
run of this operator used. The wait is bounded by the operator's deadline. ClickHouse
needs none: its materialization is found again in ``system.mutations`` and waited for.
"""
name = self.quote(self._check(table, index))
if self.dialect != "postgres":
return
while self.native.rows(
"SELECT 1 FROM pg_stat_activity WHERE pid <> pg_backend_pid() AND state = 'active' "
"AND strpos(query, %s) > 0",
[name],
):
time.sleep(POLL_SECONDS)

def build(self, table: TableIdentity, index: Mapping[str, Any]) -> None:
"""Bring the index to ``ready`` without pausing the table's writers."""
self.settle(table, index)
status, reason = self.inspect(table, index)
if status == "foreign":
raise MigrationRefused(reason)
Expand Down Expand Up @@ -105,7 +130,9 @@ def remove(self, table: TableIdentity, index: Mapping[str, Any]) -> None:
if self.dialect == "postgres":
# IF EXISTS only for an index that goes while this drop waits for its lock: an earlier
# drop whose client the budget closed runs on in the server until the transaction it
# waits for ends. What is left afterwards is read back below, as always.
# waits for ends. What is left afterwards is read back below, as always. No settle
# here, unlike a build: two drops of one index are serialized by the exclusive lock
# the first holds on it until it commits, so this one never meets a half-done drop.
self.native.command(f"DROP INDEX CONCURRENTLY IF EXISTS {self.quote(name)}")
else:
for mutation_id, is_done, _failure in self._ch_mutations(table, name):
Expand All @@ -127,7 +154,8 @@ def drop(self, table: TableIdentity, index: Mapping[str, Any]) -> None:

A foreign object under the name means ours is not there: it is left as it is.
"""
status, _ = self.inspect(table, index) # checks the name and the table first
self.settle(table, index) # checks the name and the table first
status, _ = self.inspect(table, index)
if status == "foreign":
return
name = str(index["name"])
Expand Down
103 changes: 100 additions & 3 deletions python/tests/test_index_operator_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -939,6 +939,94 @@ def release_during_the_build(step: str) -> None:
assert build.session().get("Event", {"id": 5000}) == {"id": 5000, "value": 5}


def killed_mid_build(build: Build, run: Callable[..., Any]) -> tuple[bool, bool, str, str]:
"""Kill the operator while its concurrent build waits for the held snapshot; leave the build.

Nothing here stands in for the server noticing the vanished client: it notices when it next
talks to it, which for a concurrent build is once the index is valid (measured, PostgreSQL
15.19). What is returned is the index as the killed build left it so far.
"""
name = build.names[0]
worker = spawn(build, "native")
try:
expect_line(worker, "STARTED")
pid = wait_until_held(build, run, name)
kill(worker)
finally:
kill(worker)
assert run("SELECT state FROM pg_stat_activity WHERE pid = %s", [pid]) == [("active",)]
interrupted = pg_indexes(build.role)[name]
assert interrupted[0] is False
return interrupted


def release_a_second_into(
step: str, release: Callable[[], None], run: Callable[..., Any], drops: list[bool]
) -> Callable[[str], None]:
"""At ``step``'s intent, release the snapshot a second later, noting any drop issued by then."""

def hook(seen: str) -> None:
if matches(seen, step):

def look_then_release() -> None:
try:
drops.append(
bool(
run(
"SELECT 1 FROM pg_stat_activity WHERE state = 'active' "
"AND query LIKE 'DROP INDEX CONCURRENTLY%'"
)
)
)
finally:
release()

threading.Timer(1.0, look_then_release).start()

return hook


def test_recovery_waits_for_a_killed_operators_build_and_keeps_what_it_finished(
tmp_path: Path,
) -> None:
"""The killed operator's build runs on in the server, and recovery waits for it to end.

PostgreSQL releases the table's session lock before a concurrent build's last transaction
commits. A drop waiting for that lock started in the window and failed with "tuple
concurrently updated" (SDK CI, 27 September 2026). Recovery that waits for the statement sends
no drop while it runs, and then keeps the index it finished: the same object, not a rebuild.
"""
with initial("postgres", tmp_path) as build:
drops: list[bool] = []
with admin(build.role) as run, held(build, run) as release:
interrupted = killed_mid_build(build, run)
build.operator._after_step = release_a_second_into(
"index_build:intent", release, run, drops
)
receipt = build.operator.resume().as_record()
build.operator._after_step = quiet
assert drops == [False]
assert receipt["outcome"] == "built" and receipt["recovered"] is True
assert_built(build)
assert pg_indexes(build.role)[build.names[0]][3] == interrupted[3]


def test_abandoning_waits_for_a_killed_operators_build_before_dropping_it(tmp_path: Path) -> None:
"""Abandoned beside the killed operator's build that still runs: the drop waits for its end."""
with initial("postgres", tmp_path) as build:
drops: list[bool] = []
with admin(build.role) as run, held(build, run) as release:
killed_mid_build(build, run)
build.operator._after_step = release_a_second_into(
"index_drop:intent", release, run, drops
)
receipt = build.operator.abandon().as_record()
build.operator._after_step = quiet
assert drops == [False]
assert receipt["outcome"] == "abandoned"
assert_absent(build)


@pytest.mark.parametrize("engine", ["postgres", "clickhouse"])
def test_a_killed_operator_between_watermark_and_publication_publishes_on_resume(
engine: str, tmp_path: Path
Expand All @@ -961,16 +1049,23 @@ def test_a_killed_operator_between_watermark_and_publication_publishes_on_resume
def test_the_build_budget_bounds_a_held_build_and_recovery_finishes_it(
engine: str, tmp_path: Path
) -> None:
# The first build is held for good, so any budget ends it; the resume after the release has
# to fit a DROP and a CREATE INDEX CONCURRENTLY into the same signed budget, which 2.5 s did
# not always leave on a loaded two-core machine.
# The first build is held for good, so any budget ends it. On PostgreSQL the budget closes the
# operator's connection, not its CREATE INDEX CONCURRENTLY: that finishes in the server once
# the snapshot is released, and the resume waits for it and keeps the index. Resuming at once
# used to race that statement's commit: "tuple concurrently updated" (SDK CI, 27 September).
# The budget is still 6 s, not 2.5 s, which a rebuild did not always fit on a loaded two-core
# machine.
with initial(engine, tmp_path, budget_ms=6000) as build:
interrupted = None
with admin(build.role) as run, held(build, run) as release:
started = time.monotonic()
with pytest.raises(CutoverRecoveryRequired, match="build budget"):
build.operator.index(build.plan)
assert time.monotonic() - started < 20
operator = build.reconnect()
if engine == "postgres":
interrupted = pg_indexes(build.role)[build.names[0]]
assert interrupted[0] is False
if engine == "clickhouse":
# Recovery is bounded by the same signed budget, not by the 30 s watchdog.
started = time.monotonic()
Expand All @@ -988,6 +1083,8 @@ def test_the_build_budget_bounds_a_held_build_and_recovery_finishes_it(
receipt = operator.resume().as_record()
assert receipt["outcome"] == "built" and receipt["recovered"] is True
assert_built(build)
assert interrupted is not None
assert pg_indexes(build.role)[build.names[0]][3] == interrupted[3]


def test_the_build_adds_its_own_history_and_leaves_the_rest_untouched(tmp_path: Path) -> None:
Expand Down
Loading