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
3 changes: 2 additions & 1 deletion docs/format-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -1578,7 +1578,8 @@ the canonical payload fingerprint.
exact source, write generation and routing - in another engine binding (protocol 1, a move) or in
the source's own (protocol 2, a relayout). Both SDKs validate the signed packet, map binding and
portable physical names. `migration/099`–`121`, `138` and `139` pin these shared rules; the Python
local operator executes creation and recovery. Loading an authorization does not activate its prepared map.
local operator executes creation, recovery and - before the decision `prepared` - abandonment.
Loading an authorization does not activate its prepared map.


## 7i. Physical design (map contract 5)
Expand Down
5 changes: 4 additions & 1 deletion docs/local-cutover.md
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,10 @@ reattached only with the matching stored UUID and native drain intent. Partial g
are completed without lowering an epoch. A replaced table/login or another active map is refused.
Repeated `execute` for an already completed identical packet returns its saved receipt. `resume`
requires an unfinished execution; a response lost after completion can be recovered with `execute`.
Do not drop retired tables or remove the state on the basis of a missing response.
Do not drop retired tables or remove the state on the basis of a missing response. A cutover has no
`abandon`: without a durable decision its recovery aborts it by itself. `abandon` ends an unfinished
[staging](staging.md#abandoning-a-staging-that-cannot-finish) or
[in-place index build](in-place-index.md#abandonment) instead.

## Pause budget and receipt

Expand Down
15 changes: 15 additions & 0 deletions docs/qualification/staging-abandon/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
# Abandoning a staging: evidence

[Abandonment](../../staging.md#abandoning-a-staging-that-cannot-finish) removes a staging's own
copy before its decision. One engine behaviour decides what it has to do beyond `DROP TABLE`, and it
was measured before the code was written - ClickHouse 24.8.14.39 in the SDK's test container,
24 September 2026:

`probe_grants_after_drop.py` -> `probe_grants_after_drop.out.json`: a runtime login's `SELECT` and
`INSERT` grants on a table are still listed in `system.grants` after `DROP TABLE ... SYNC`, and
`REVOKE` on the dropped table is accepted and removes them. So an abandonment on ClickHouse revokes
the runtime logins' grants on the copy after dropping it, and reads `system.grants` back; PostgreSQL
removes a table's privileges with the table.

The behaviour of the abandonment itself is pinned by `python/tests/test_staging_abandon_live.py` on
both engines.
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
{
"after_grant": [
[
"SELECT",
"probe_7e56929922a5",
"t"
],
[
"INSERT",
"probe_7e56929922a5",
"t"
]
],
"after_drop_table": [
[
"SELECT",
"probe_7e56929922a5",
"t"
],
[
"INSERT",
"probe_7e56929922a5",
"t"
]
],
"revoke_on_dropped_table": "accepted",
"after_revoke": []
}
30 changes: 30 additions & 0 deletions docs/qualification/staging-abandon/probe_grants_after_drop.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
"""ClickHouse 24.8: do table grants outlive DROP TABLE, and does REVOKE accept a dropped table?

SDE_CLICKHOUSE_DSN=clickhouse://default:sde@127.0.0.1:58123/sde python probe_grants_after_drop.py
"""
import json, os, uuid
from urllib.parse import urlsplit
import clickhouse_connect
p = urlsplit(os.environ["SDE_CLICKHOUSE_DSN"])
root = clickhouse_connect.get_client(host=p.hostname, port=p.port, username=p.username, password=p.password)
db, user = "probe_" + uuid.uuid4().hex[:12], "probe_user_" + uuid.uuid4().hex[:8]
out = {}
root.command(f"CREATE DATABASE {db} ENGINE = Atomic")
root.command(f"CREATE USER {user} IDENTIFIED WITH sha256_password BY 'x{uuid.uuid4().hex}'")
try:
root.command(f"CREATE TABLE {db}.t (id Int64) ENGINE = MergeTree ORDER BY id")
root.command(f"GRANT SELECT, INSERT ON {db}.t TO {user}")
grants = lambda: [list(map(str, r)) for r in root.query(f"SELECT access_type, database, table FROM system.grants WHERE user_name = '{user}' ORDER BY access_type").result_rows]
out["after_grant"] = grants()
root.command(f"DROP TABLE {db}.t SYNC")
out["after_drop_table"] = grants()
try:
root.command(f"REVOKE SELECT, INSERT ON {db}.t FROM {user}")
out["revoke_on_dropped_table"] = "accepted"
except Exception as exc:
out["revoke_on_dropped_table"] = "refused: " + str(exc)[:160]
out["after_revoke"] = grants()
finally:
root.command(f"DROP USER IF EXISTS {user}")
root.command(f"DROP DATABASE IF EXISTS {db} SYNC")
print(json.dumps(out, indent=1))
32 changes: 30 additions & 2 deletions docs/staging.md
Original file line number Diff line number Diff line change
Expand Up @@ -84,14 +84,42 @@ connections and inspection of persisted state. It is not a workload latency guar

Do not delete the project directory or manually replace `active-map.json` for the next migration.
The same directory retains stage receipts, cutover decisions and retired names through subsequent
successes and aborts. After the first staging intent, `project.json` uses storage contract 2; older
successes and aborts. After the first staging intent, `project.json` uses storage contract 2; after
an abandoned staging, storage contract 4, which operators that know contracts 1 to 3 refuse; older
operators must refuse it. Completed retries reconfirm filesystem durability before returning the
stored receipt, including a retry after an uncertain final directory fsync.

## Abandoning a staging that cannot finish

A staging can reach a state no resume repairs: another operation's barrier appeared on a source
table, a runtime login or its grants changed, the target engine refuses the table or does not
answer within the operator's watchdog, or somebody else's object took a copy's name. Before its
decision `prepared`, such a staging is abandoned:

```sh
sde-operator --config local-operator.json --project-dir ./client-state abandon
```

`LocalCutover.abandon()` records the decision `abandoned` first, then removes **only this
staging's own tables**: a table under a copy's name is removed when it carries this staging's exact
creation marker and, once its identity was recorded, that native object - also a table created
just before a crash, whose identity was never recorded. Anything else under the name is left
alone. Indexes and generation constraints go with the table. On ClickHouse the runtime logins'
grants on the copy are revoked as well: measured on 24.8.14.39, table grants outlive `DROP TABLE`
(`qualification/staging-abandon/`). The map in force, the watermarks and every running process
stay as they were. Abandonment needs only the same native database - not the runtime logins, not a
source free of another barrier - which is what keeps it available when a staging cannot finish.
An interrupted abandonment is finished by `resume` or by `abandon` again, and the authorization is
spent: `stage` with the same packet returns the abandonment.

After the decision `prepared` the next map is decided: `abandon` is refused, `resume` publishes the
prepared map, and the copy leaves through its cutover's abort. The same command also abandons an
unfinished [in-place index build](in-place-index.md).

## Receipt and limits

The metadata-only receipt contains `protocol: 1`, `stage_id`, `stage_fingerprint`, `project_id`,
`group`, `outcome: "prepared"`, `map_version`, `map_fingerprint`, `tables`, `elapsed_ms`, and
`group`, `outcome` (`prepared`, or `abandoned`), `map_version`, `map_fingerprint`, `tables`, `elapsed_ms`, and
`recovered`. Each table identifies its engine binding, logical entity and native `identity`
(`dialect`, `server`, `database`, `namespace`, `object`, `name`). No rows or credentials are included.
`StagingReceipt.as_record()` returns an independent snapshot. The controller must validate the
Expand Down
28 changes: 25 additions & 3 deletions python/src/sde/_cutover_project.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,18 @@ def encode(value: Any) -> bytes:
).encode("utf-8")


def _abandoned_stages(payload: Any) -> bool:
stages = payload.get("stages") if isinstance(payload, dict) else None
if not isinstance(stages, dict):
return False
return any(
isinstance(record, dict)
and isinstance(record.get("receipt"), dict)
and record["receipt"].get("outcome") == "abandoned"
for record in stages.values()
)


class ProjectState:
"""One POSIX directory and lock, shared by all local executors for this project."""

Expand All @@ -43,7 +55,7 @@ def read(self) -> dict[str, Any]:
raise ValueError("invalid state envelope")
if type(envelope["storage_contract"]) is not int or envelope[
"storage_contract"
] not in (1, 2, 3):
] not in (1, 2, 3, 4):
raise ValueError("unsupported state storage contract")
payload = envelope["payload"]
if (
Expand All @@ -68,12 +80,16 @@ def read(self) -> dict[str, Any]:
fields.add("stages")
if not isinstance(payload.get("stages"), dict):
raise ValueError("staging history must be an object")
if envelope["storage_contract"] == 3:
if envelope["storage_contract"] >= 3:
# In-place index builds: a reader of contracts 1 and 2 refuses this state rather
# than ignoring a history it would not know to keep.
fields.add("indexes")
if not isinstance(payload.get("indexes"), dict):
raise ValueError("index build history must be an object")
if envelope["storage_contract"] < 4 and _abandoned_stages(payload):
# Contract 4 is what makes an older operator refuse an abandoned staging's record
# instead of reading a receipt whose tables may name no identity.
raise ValueError("an abandoned staging needs state storage contract 4")
if set(payload) != fields:
raise ValueError("unknown or missing project state fields")
return payload
Expand All @@ -92,7 +108,13 @@ def confirm(self) -> None:
def write(self, payload: dict[str, Any]) -> None:
body = encode(payload)
envelope = {
"storage_contract": 3 if "indexes" in payload else 2 if "stages" in payload else 1,
"storage_contract": 4
if _abandoned_stages(payload)
else 3
if "indexes" in payload
else 2
if "stages" in payload
else 1,
"payload": payload,
"sha256": hashlib.sha256(body).hexdigest(),
}
Expand Down
53 changes: 53 additions & 0 deletions python/src/sde/engines/_staging.py
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,59 @@ def create_table(
raise MigrationRefused("staging creation marker was not established with the table")
return self.native.identity(table)

def owned(
self, table: str, marker: str, identity: TableIdentity | None
) -> TableIdentity | None:
"""This staging's own table under ``table``, or ``None`` when it is absent or not ours.

Ours means the exact creation marker and, once the identity was recorded, that native
object. Anything else under the name - another marker, none, another object - belongs to
somebody else and is left alone by an abandonment.
"""
present = self.marker(table)
if present is None or present[1] != marker:
return None
if identity is not None:
return identity if present[0] == identity.object else None
return self.native.identity(table)

def drop_owned(self, table: TableIdentity, marker: str) -> None:
"""Drop this staging's own table; its indexes and generation constraints go with it."""
present = self.marker(table.name)
if present is None:
return # dropped before a crash; nothing is left to remove
if present[1] != marker or present[0] != table.object:
raise MigrationRefused("the table to abandon is no longer this staging's own")
# SYNC: in an Atomic database a dropped table otherwise lingers until the server removes
# it, and its name and data with it.
suffix = " SYNC" if self.dialect == "clickhouse" else ""
self.native.command(f"DROP TABLE {self.quote(table.name)}{suffix}")
if self.marker(table.name) is not None:
raise MigrationRefused("an abandoned staging table is still in the catalogue")

def revoke_runtime(self, tables: Sequence[TableIdentity], principals: Sequence[str]) -> None:
"""Take back the runtime grants on dropped tables; ClickHouse keeps them after DROP TABLE.

Measured on 24.8.14.39: ``system.grants`` still lists a table's grants once the table is
gone, and ``REVOKE`` on the dropped table is accepted and removes them. PostgreSQL removes a
table's privileges with the table, so there is nothing to take back there.
"""
if self.dialect != "clickhouse":
return
for table in tables:
qualified = self.quote(table.namespace) + "." + self.quote(table.name)
for name in principals:
self.native.command(f"REVOKE SELECT, INSERT ON {qualified} FROM {self.quote(name)}")
for table in tables:
for name in principals:
rows = self.native.rows(
"SELECT count() FROM system.grants WHERE user_name={user:String} "
"AND database={database:String} AND table={table:String}",
{"user": name, "database": table.namespace, "table": table.name},
)
if rows[0][0]:
raise MigrationRefused("a runtime grant on an abandoned staging table remains")

def create_indexes(self, layout: PhysicalLayout) -> None:
if self.dialect != "postgres":
# ClickHouse data-skipping indexes were declared inside `CREATE TABLE`; qualification
Expand Down
23 changes: 18 additions & 5 deletions python/src/sde/local_cutover.py
Original file line number Diff line number Diff line change
Expand Up @@ -679,21 +679,34 @@ def index(self, plan: IndexPlan) -> IndexReceipt:
finally:
self._alarm = None

def abandon(self) -> IndexReceipt:
"""Abandon the unfinished index build: drop its own indexes, keep the map in force."""
def abandon(self) -> IndexReceipt | StagingReceipt:
"""Abandon the unfinished index build or staging before its decision.

Only this operation's own objects are removed - the indexes a build created, the tables a
staging created - and the map in force stays. A cutover has no abandonment: without a
decision its recovery aborts it by itself.
"""
from ._operator_deadline import DeadlineInterrupt, OperatorDeadline
from .index_operator import abandon_index
from .staging_operator import abandon_stage

self._owner()
try:
with OperatorDeadline() as deadline:
self._alarm = deadline
deadline.arm(30000) # re-armed from the stored authorization's build budget
return abandon_index(self)
deadline.arm(30000) # an index build re-arms from its signed build budget
with self.store.lock():
execution = self.store.read()["execution"]
kind = None if execution is None else execution.get("kind")
if kind == "index":
return abandon_index(self)
if kind == "staging":
return abandon_stage(self)
raise MigrationRefused("there is no unfinished index build or staging to abandon")
except DeadlineInterrupt as exc:
self._interrupted()
raise CutoverRecoveryRequired(
"abandoning the index build was interrupted; reconnect and abandon again"
"abandoning was interrupted; reconnect and abandon again"
) from exc
finally:
self._alarm = None
Expand Down
Loading
Loading