From 2eea384ae46fabcfa6098fd47fc458287cca0e9f Mon Sep 17 00:00:00 2001 From: Krzysztof Macewicz Date: Thu, 24 Sep 2026 19:15:34 +0200 Subject: [PATCH 1/3] feat: abandon a staging that cannot finish A staging could reach a state no resume repairs - another operation's barrier on a source table, a changed runtime login, a target that refuses the table, somebody else's object under a copy's name - and then blocked the customer's operator until its state was repaired by hand. Before its decision, LocalCutover.abandon / sde-operator abandon now ends it: the decision abandoned is recorded first, then only this staging's own tables are dropped - the ones under a copy's name that carry its exact creation marker and, once recorded, its native identity, including a table created just before a crash whose identity was never recorded. Anything else under the name is left alone. On ClickHouse the runtime logins' grants are revoked too: measured on 24.8.14.39, they outlive DROP TABLE. The map in force, the watermarks and every process stay as they were, and abandonment needs only the same native database - not the logins or a barrier-free source whose absence is why the staging cannot finish. The receipt gains the outcome abandoned (the map in force, per entity the identity of the dropped table or null). A state holding an abandoned staging is storage contract 4, which operators knowing contracts 1 to 3 refuse; the name-reuse check reads names from the stored authorization. After the decision prepared, abandonment is refused and resume publishes. The in-place index budget test gets a longer signed budget: its resume must fit a DROP and a CREATE INDEX CONCURRENTLY into it, which 2.5 s did not always leave on a loaded two-core machine. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/format-contract.md | 3 +- docs/local-cutover.md | 5 +- docs/qualification/staging-abandon/README.md | 15 + .../probe_grants_after_drop.out.json | 28 ++ .../probe_grants_after_drop.py | 30 ++ docs/staging.md | 32 +- python/src/sde/_cutover_project.py | 28 +- python/src/sde/engines/_staging.py | 53 +++ python/src/sde/local_cutover.py | 23 +- python/src/sde/staging_operator.py | 113 +++++- python/src/sde_operator/__main__.py | 2 +- python/tests/test_cutover_project.py | 41 +++ python/tests/test_index_operator_live.py | 7 +- python/tests/test_staging_abandon_live.py | 330 ++++++++++++++++++ 14 files changed, 694 insertions(+), 16 deletions(-) create mode 100644 docs/qualification/staging-abandon/README.md create mode 100644 docs/qualification/staging-abandon/probe_grants_after_drop.out.json create mode 100644 docs/qualification/staging-abandon/probe_grants_after_drop.py create mode 100644 python/tests/test_staging_abandon_live.py diff --git a/docs/format-contract.md b/docs/format-contract.md index 71be5f2..27575f2 100644 --- a/docs/format-contract.md +++ b/docs/format-contract.md @@ -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) diff --git a/docs/local-cutover.md b/docs/local-cutover.md index dd79c11..48caa66 100644 --- a/docs/local-cutover.md +++ b/docs/local-cutover.md @@ -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 diff --git a/docs/qualification/staging-abandon/README.md b/docs/qualification/staging-abandon/README.md new file mode 100644 index 0000000..4a092e4 --- /dev/null +++ b/docs/qualification/staging-abandon/README.md @@ -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. diff --git a/docs/qualification/staging-abandon/probe_grants_after_drop.out.json b/docs/qualification/staging-abandon/probe_grants_after_drop.out.json new file mode 100644 index 0000000..bb954f8 --- /dev/null +++ b/docs/qualification/staging-abandon/probe_grants_after_drop.out.json @@ -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": [] +} diff --git a/docs/qualification/staging-abandon/probe_grants_after_drop.py b/docs/qualification/staging-abandon/probe_grants_after_drop.py new file mode 100644 index 0000000..5e57075 --- /dev/null +++ b/docs/qualification/staging-abandon/probe_grants_after_drop.py @@ -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)) diff --git a/docs/staging.md b/docs/staging.md index 0fbb1da..1a075de 100644 --- a/docs/staging.md +++ b/docs/staging.md @@ -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 diff --git a/python/src/sde/_cutover_project.py b/python/src/sde/_cutover_project.py index 9cb23de..a160d8e 100644 --- a/python/src/sde/_cutover_project.py +++ b/python/src/sde/_cutover_project.py @@ -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.""" @@ -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 ( @@ -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 @@ -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(), } diff --git a/python/src/sde/engines/_staging.py b/python/src/sde/engines/_staging.py index 897895e..6070960 100644 --- a/python/src/sde/engines/_staging.py +++ b/python/src/sde/engines/_staging.py @@ -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 diff --git a/python/src/sde/local_cutover.py b/python/src/sde/local_cutover.py index 7964323..31fa675 100644 --- a/python/src/sde/local_cutover.py +++ b/python/src/sde/local_cutover.py @@ -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 diff --git a/python/src/sde/staging_operator.py b/python/src/sde/staging_operator.py index 792548b..67a4e36 100644 --- a/python/src/sde/staging_operator.py +++ b/python/src/sde/staging_operator.py @@ -68,7 +68,11 @@ def _snapshot(operator: LocalCutover, plan: StagingPlan, state: dict[str, Any]) if any(row[-1] in table_names for row in state["retired_names"]): raise MigrationRefused("staging cannot reuse a retired physical name") for previous in state.get("stages", {}).values(): - if any(row["identity"]["name"] in table_names for row in previous["receipt"]["tables"]): + # From the stored authorization, not the receipt: an abandoned staging's receipt names no + # identity for a table it never created, and its names are spent all the same. + record = previous["plan"] + used = record["prepared"]["groups"][record["group"]]["derived"][0]["layout"]["tables"] + if table_names & set(used.values()): raise MigrationRefused("staging cannot reuse a previously prepared physical name") return { "kind": "staging", @@ -117,6 +121,8 @@ def _finish( from .engines._staging import NativeStaging execution = state["execution"] + if execution["decision"] == "abandoned": + return _abandon(operator, state, plan, recovered=recovered) target = plan.prepared.groups[plan.group].derived[0] epoch = plan.prepared.groups[plan.group].write_epoch assert epoch is not None @@ -238,6 +244,111 @@ def publish() -> None: return StagingReceipt(receipt) +def _abandon( + operator: LocalCutover, state: dict[str, Any], plan: StagingPlan, *, recovered: bool +) -> StagingReceipt: + """Remove this staging's own tables and keep the map in force; the authorization is spent. + + It needs only the same native database: not the runtime logins, not a source free of another + barrier - the things whose absence is why a staging cannot finish - because it publishes + nothing. Anything under a staging name that is not this staging's is left where it is. + """ + from .engines._operator import TableIdentity + from .engines._staging import NativeStaging + + execution = state["execution"] + target = plan.prepared.groups[plan.group].derived[0] + binding = execution["bindings"][target.engine] + native = operator.native[target.engine] + if list(native.endpoint()) != binding["endpoint"]: + raise MigrationRefused("a staging binding names another native database") + creator = NativeStaging(native) + for row in execution["tables"]: + + def drop(row: dict[str, Any] = row) -> None: + if "removed" not in row: + recorded = None if row["identity"] is None else TableIdentity(**row["identity"]) + found = creator.owned(row["table"], row["marker"], recorded) + row["removed"] = None if found is None else found.as_record() + # Durable before the destructive statement: a crash after DROP must still know + # which object the receipt removed. + operator.store.write(state) + if row["removed"] is not None: + creator.drop_owned(TableIdentity(**row["removed"]), row["marker"]) + + operator._step(state, "stage_drop_" + row["entity"], drop) + removed = [ + TableIdentity(**row["removed"]) for row in execution["tables"] if row["removed"] is not None + ] + if native.dialect == "clickhouse" and removed: + operator._step( + state, + "stage_revoke", + lambda: creator.revoke_runtime(removed, sorted(binding["principals"])), + ) + if operator.active_map().fingerprint != plan.current.fingerprint: + raise MigrationRefused("an abandoned staging found another active map") + receipt = { + "protocol": 1, + "stage_id": plan.stage_id, + "stage_fingerprint": plan.fingerprint, + "project_id": operator.project_id, + "group": plan.group, + "outcome": "abandoned", + "map_version": plan.current.map_version, + "map_fingerprint": plan.current.fingerprint, + "tables": [ + {"engine": row["engine"], "entity": row["entity"], "identity": row["removed"]} + for row in execution["tables"] + ], + "elapsed_ms": operator._elapsed(), + "recovered": recovered, + } + state["stages"][plan.stage_id] = { + "plan_fingerprint": plan.fingerprint, + "plan": plan.as_record(), + "receipt": receipt, + } + state.setdefault("indexes", {}) # an abandoned staging record is storage contract 4 + state["execution"] = None + operator.store.write(state) + operator._after_step("stage_complete") + return StagingReceipt(receipt) + + +def abandon_stage(operator: LocalCutover) -> StagingReceipt: + """Abandon the unfinished staging before its decision: its own tables go, the map stays.""" + with operator.store.lock(): + state = operator.store.read() + execution = state["execution"] + if execution is None or execution.get("kind") != "staging": + raise MigrationRefused("there is no unfinished staging to abandon") + if execution["decision"] == "prepared": + raise MigrationRefused( + "a prepared staging cannot be abandoned: its next map is decided; resume to " + "publish it, or abort its cutover" + ) + plan = load_staging_plan( + execution["plan"], + model=operator.model, + project_id=operator.project_id, + public_key=operator.keys, + ) + if plan.fingerprint != execution["plan_fingerprint"]: + raise MigrationRefused("staging recovery names another signed authorization") + if execution["decision"] is None: + execution["decision"] = "abandoned" + operator.store.write(state) + operator._after_step("stage_abandoned") + operator._started, operator._enforce_budget = monotonic_ns(), False + try: + return _finish(operator, state, plan, recovered=True) + except Exception as exc: + raise CutoverRecoveryRequired( + "abandoning the staging remains incomplete; preserve its state" + ) from exc + + def execute_stage(operator: LocalCutover, plan: StagingPlan) -> StagingReceipt: with operator.store.lock(): state = operator.store.read() diff --git a/python/src/sde_operator/__main__.py b/python/src/sde_operator/__main__.py index c2e4986..7b32539 100644 --- a/python/src/sde_operator/__main__.py +++ b/python/src/sde_operator/__main__.py @@ -157,7 +157,7 @@ def main(argv: Sequence[str] | None = None) -> int: "message": ( "Inspect local status and resume with fresh connections; " "the durable decision must be preserved. An unfinished index build " - "may be abandoned instead." + "or staging may be abandoned instead." ), } ), diff --git a/python/tests/test_cutover_project.py b/python/tests/test_cutover_project.py index 0425af4..2622ec7 100644 --- a/python/tests/test_cutover_project.py +++ b/python/tests/test_cutover_project.py @@ -151,3 +151,44 @@ def test_the_index_build_history_belongs_to_contract_three_only( _rewrite(store, contract, payload) with pytest.raises(MigrationRefused, match="corrupt"): store.read() + + +def _stage_record(outcome: str) -> dict[str, Any]: + return {"plan_fingerprint": "f" * 64, "plan": {}, "receipt": {"outcome": outcome}} + + +def test_an_abandoned_staging_raises_the_state_contract_to_four(tmp_path: Path) -> None: + store = ProjectState(tmp_path, PROJECT, MODEL) + store.enroll({"map_version": 1}, b"map") + state = store.read() + state["stages"] = {"a" * 32: _stage_record("prepared")} + state["indexes"] = {} + store.write(state) + assert json.loads(store.path.read_bytes())["storage_contract"] == 3 + state["stages"]["b" * 32] = _stage_record("abandoned") + store.write(state) + assert json.loads(store.path.read_bytes())["storage_contract"] == 4 + assert store.read()["stages"]["b" * 32]["receipt"]["outcome"] == "abandoned" + + +@pytest.mark.parametrize("contract", [2, 3]) +def test_an_abandoned_staging_record_needs_contract_four(tmp_path: Path, contract: int) -> None: + store = ProjectState(tmp_path, PROJECT, MODEL) + store.enroll({"map_version": 1}, b"map") + payload = store.read() + payload["stages"] = {"b" * 32: _stage_record("abandoned")} + if contract == 3: + payload["indexes"] = {} + _rewrite(store, contract, payload) + with pytest.raises(MigrationRefused, match="corrupt"): + store.read() + + +def test_contract_four_still_carries_the_index_history(tmp_path: Path) -> None: + store = ProjectState(tmp_path, PROJECT, MODEL) + store.enroll({"map_version": 1}, b"map") + payload = store.read() + payload["stages"] = {"b" * 32: _stage_record("abandoned")} + _rewrite(store, 4, payload) # no "indexes" field + with pytest.raises(MigrationRefused, match="corrupt"): + store.read() diff --git a/python/tests/test_index_operator_live.py b/python/tests/test_index_operator_live.py index af7cfba..2c4f664 100644 --- a/python/tests/test_index_operator_live.py +++ b/python/tests/test_index_operator_live.py @@ -951,7 +951,10 @@ 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: - with initial(engine, tmp_path, budget_ms=2500) as build: + # 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. + with initial(engine, tmp_path, budget_ms=6000) as build: with admin(build.role) as run, held(build, run) as release: started = time.monotonic() with pytest.raises(CutoverRecoveryRequired, match="build budget"): @@ -1067,5 +1070,5 @@ def test_the_operator_cli_abandons_an_unfinished_build_once(tmp_path: Path) -> N code, refusal = run_cli([*base, "abandon"], environment) assert code == 2 and refusal == { "error": "refused", - "message": "there is no unfinished index build to abandon", + "message": "there is no unfinished index build or staging to abandon", } diff --git a/python/tests/test_staging_abandon_live.py b/python/tests/test_staging_abandon_live.py new file mode 100644 index 0000000..c293e53 --- /dev/null +++ b/python/tests/test_staging_abandon_live.py @@ -0,0 +1,330 @@ +"""A staging that cannot finish is abandoned: its own copy goes, the map in force stays. + +Before its decision a staging holds nothing the application relies on - the prepared map is not +published, no process writes to the copy - so abandoning it removes the tables it created and +keeps everything else as it was. After the decision the next map is decided and only a resume +remains; the copy then leaves through its cutover's abort. +""" + +from __future__ import annotations + +import json +from collections.abc import Callable +from pathlib import Path +from typing import Any + +import pytest +from test_staging_operator_live import PROJECT, initial + +import sde +from sde.local_cutover import CutoverRecoveryRequired + +BEFORE_DECISION = [ + "stage_prepared", + "stage_create_Event:intent", + "stage_create_Event:done", + "stage_generation_Event:intent", + "stage_generation_Event:done", + "stage_indexes:intent", + "stage_indexes:done", + "stage_qualify:intent", + "stage_qualify:done", + "stage_grants:intent", + "stage_grants:done", +] +AFTER_DECISION = [ + "stage_decision", + "stage_watermarks:intent", + "stage_watermarks:done", + "stage_publish:intent", + "stage_publish:done", +] + + +class Crash(BaseException): + pass + + +def crash_at(checkpoint: str) -> Callable[[str], None]: + def hook(step: str) -> None: + if step == checkpoint: + raise Crash(step) + + return hook + + +def quiet(_step: str) -> None: + return None + + +def copy_table(stage: Any) -> str: + return str(stage.prepared.groups["Event"].derived[0].layout.tables["Event"]) + + +def exists(role: Any, table: str) -> bool: + if role.operator.dialect == "postgres": + row = role.operator._cx.execute("SELECT to_regclass(%s)", (f'"{table}"',)).fetchone() + return row[0] is not None + return bool( + role.operator._cx.query( + "SELECT count() FROM system.tables WHERE database = {d:String} AND name = {t:String}", + parameters={"d": role.namespace, "t": table}, + ).result_rows[0][0] + ) + + +def runtime_grants(role: Any, table: str) -> int: + """ClickHouse grants of the fixture's runtime login on ``table`` - they survive DROP TABLE.""" + return int( + role.operator._cx.query( + "SELECT count() FROM system.grants WHERE user_name = {u:String} " + "AND database = {d:String} AND table = {t:String}", + parameters={"u": role.username, "d": role.namespace, "t": table}, + ).result_rows[0][0] + ) + + +def assert_abandoned( + operator: Any, stage: Any, old: Any, roles: Any, target: str, receipt: dict[str, Any] +) -> None: + assert receipt["outcome"] == "abandoned" + assert receipt["map_version"] == stage.current.map_version + assert receipt["map_fingerprint"] == stage.current.fingerprint + assert operator.active_map().fingerprint == stage.current.fingerprint + table = copy_table(stage) + assert not exists(roles[target], table) + if target == "clickhouse": + assert runtime_grants(roles[target], table) == 0 + status = operator.status() + assert (status["active_map_version"], status["recorded_map_version"]) == (1, 1) + assert status["plan_id"] is None + old.save("Event", {"id": 90, "value": 9}) # the source was never paused + assert old.get("Event", {"id": 90}) == {"id": 90, "value": 9} + envelope = json.loads(operator.store.path.read_text()) + assert envelope["storage_contract"] == 4 + assert envelope["payload"]["stages"][stage.stage_id]["receipt"] == receipt + + +@pytest.mark.parametrize("checkpoint", BEFORE_DECISION) +@pytest.mark.parametrize("source", ["postgres", "clickhouse"]) +def test_abandoning_before_the_decision_removes_only_the_copy( + source: str, checkpoint: str, tmp_path: Path +) -> None: + with initial(source, tmp_path) as (operator, stage, old, roles, *_, target): + operator._after_step = crash_at(checkpoint) + with pytest.raises(Crash): + operator.stage(stage) + operator._after_step = quiet + receipt = operator.abandon().as_record() + assert receipt["recovered"] is True + (row,) = receipt["tables"] + created = checkpoint not in ("stage_prepared", "stage_create_Event:intent") + assert (row["identity"] is not None) is created + if created: + assert row["identity"]["name"] == copy_table(stage) + assert_abandoned(operator, stage, old, roles, target, receipt) + # The authorization is spent: the same staging answers with its abandonment. + assert operator.stage(stage).as_record() == receipt + with pytest.raises(sde.MigrationRefused, match="no unfinished index build or staging"): + operator.abandon() + + +@pytest.mark.parametrize("checkpoint", AFTER_DECISION) +def test_a_prepared_staging_cannot_be_abandoned(checkpoint: str, tmp_path: Path) -> None: + with initial("postgres", tmp_path) as (operator, stage, _old, roles, *_, target): + operator._after_step = crash_at(checkpoint) + with pytest.raises(Crash): + operator.stage(stage) + operator._after_step = quiet + with pytest.raises(sde.MigrationRefused, match="cannot be abandoned"): + operator.abandon() + assert operator.store.read()["execution"]["decision"] == "prepared" + assert operator.resume().as_record()["outcome"] == "prepared" + assert exists(roles[target], copy_table(stage)) + + +@pytest.mark.parametrize( + ("checkpoint", "finish"), + [ + ("stage_abandoned", "resume"), + ("stage_abandoned", "abandon"), + ("stage_drop_Event:intent", "resume"), + ("stage_drop_Event:done", "abandon"), + ("stage_revoke:intent", "resume"), + ], +) +def test_an_interrupted_abandonment_completes_on_resume_or_retry( + checkpoint: str, finish: str, tmp_path: Path +) -> None: + with initial("postgres", tmp_path) as (operator, stage, old, roles, *_, target): + assert target == "clickhouse" # the revoke step exists only there + operator._after_step = crash_at("stage_grants:done") + with pytest.raises(Crash): + operator.stage(stage) + assert runtime_grants(roles[target], copy_table(stage)) == 2 # SELECT and INSERT + operator._after_step = crash_at(checkpoint) + with pytest.raises(Crash): + operator.abandon() + operator._after_step = quiet + # Once decided, an abandonment stays one: a recovery removes, it never builds again. + assert operator.store.read()["execution"]["decision"] == "abandoned" + finished = operator.resume() if finish == "resume" else operator.abandon() + receipt = finished.as_record() + assert receipt["tables"][0]["identity"]["name"] == copy_table(stage) + assert_abandoned(operator, stage, old, roles, target, receipt) + + +def test_a_staging_blocked_by_another_barrier_is_abandoned(tmp_path: Path) -> None: + """The reason abandonment exists: a staging whose every resume is refused.""" + with initial("postgres", tmp_path) as (operator, stage, old, roles, *_, target): + operator._after_step = crash_at("stage_create_Event:done") + with pytest.raises(Crash): + operator.stage(stage) + operator._after_step = quiet + fence = roles["postgres"].operator.write_fence("initial_events", project_id=PROJECT) + fence.freeze("e" * 32) + try: + for _ in range(2): + with pytest.raises(CutoverRecoveryRequired) as refused: + operator.resume() + assert isinstance(refused.value.__cause__, sde.MigrationRefused) + receipt = operator.abandon().as_record() + assert receipt["outcome"] == "abandoned" + assert "e" * 32 in fence.state().holds # the other operation's barrier stands + assert not exists(roles[target], copy_table(stage)) + assert operator.active_map().fingerprint == stage.current.fingerprint + finally: + fence.release("e" * 32) + old.save("Event", {"id": 91, "value": 1}) + + +def test_a_foreign_table_under_the_copy_name_is_left_alone(tmp_path: Path) -> None: + with initial("postgres", tmp_path) as (operator, stage, _old, roles, *_, target): + operator._after_step = crash_at("stage_prepared") + with pytest.raises(Crash): + operator.stage(stage) + operator._after_step = quiet + table = copy_table(stage) + roles[target].command(f"CREATE TABLE `{table}` (id Int64) ENGINE = MergeTree ORDER BY id") + receipt = operator.abandon().as_record() + assert receipt["outcome"] == "abandoned" + assert receipt["tables"][0]["identity"] is None # nothing of this staging's was there + assert exists(roles[target], table) # somebody else's table, untouched + assert operator.active_map().fingerprint == stage.current.fingerprint + + +@pytest.mark.parametrize("source", ["postgres", "clickhouse"]) +def test_a_killed_staging_leaves_a_table_the_abandonment_finds_by_its_marker( + source: str, tmp_path: Path +) -> None: + """Killed after CREATE, before the identity was recorded: only the marker says whose it is.""" + import os + import select + import signal + import subprocess + import sys + + from psycopg.conninfo import make_conninfo + + with initial(source, tmp_path) as (operator, stage, old, roles, model, public, *_, target): + payload = { + "directory": str(tmp_path), + "project_id": PROJECT, + "model": sde.neutral_declaration(model), + "public_key": public.hex(), + "plan": stage.as_record(), + "checkpoint": "native:created", + "operators": { + name: ( + make_conninfo(role.operator._dsn, options="-csearch_path=" + role.namespace) + if name == "postgres" + else role.operator._dsn + ) + for name, role in roles.items() + }, + "runtime": {name: role.runtime._dsn for name, role in roles.items()}, + } + worker = subprocess.Popen( + [sys.executable, str(Path(__file__).with_name("_staging_worker.py"))], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + try: + assert worker.stdin is not None and worker.stdout is not None + worker.stdin.write(json.dumps(payload)) + worker.stdin.close() + assert select.select([worker.stdout], [], [], 25)[0], "worker missed its checkpoint" + line = worker.stdout.readline() + if line != "READY\n": + assert worker.stderr is not None + raise AssertionError(worker.stderr.read()) + os.kill(worker.pid, signal.SIGKILL) + worker.wait(timeout=10) + finally: + if worker.poll() is None: + worker.kill() + worker.wait(timeout=10) + (row,) = operator.store.read()["execution"]["tables"] + assert row["identity"] is None and exists(roles[target], copy_table(stage)) + receipt = operator.abandon().as_record() + assert receipt["tables"][0]["identity"]["name"] == copy_table(stage) + assert_abandoned(operator, stage, old, roles, target, receipt) + + +def test_the_operator_cli_abandons_a_staging(tmp_path: Path) -> None: + import base64 + import os + import subprocess + import sys + + from psycopg.conninfo import make_conninfo + + root = tmp_path / "state" + with initial("postgres", root) as (operator, stage, old, roles, model, public, *_): + operator._after_step = crash_at("stage_qualify:done") + with pytest.raises(Crash): + operator.stage(stage) + config: dict[str, Any] = { + "protocol": 1, + "project_id": PROJECT, + "model": sde.neutral_declaration(model), + "public_keys": {"primary": base64.b64encode(public).decode()}, + "engines": {}, + } + environment = dict(os.environ) + for name, role in roles.items(): + op_name, run_name = "SDE_ABANDON_" + name.upper(), "SDE_ABANDON_APP_" + name.upper() + dsn = role.operator._dsn + if name == "postgres": + dsn = make_conninfo(dsn, options="-csearch_path=" + role.namespace) + environment[op_name], environment[run_name] = dsn, role.runtime._dsn + config["engines"][name] = { + "dialect": name, + "operator_dsn_env": op_name, + "runtime_dsn_envs": [run_name], + } + config_path = tmp_path / "config.json" + config_path.write_text(json.dumps(config)) + command = [ + sys.executable, + "-m", + "sde_operator", + "--project-dir", + str(operator.store.root), + "--config", + str(config_path), + "abandon", + ] + done = subprocess.run(command, env=environment, capture_output=True, text=True, timeout=90) + assert done.returncode == 0, done.stderr + receipt = json.loads(done.stdout) + assert receipt["outcome"] == "abandoned" + again = subprocess.run(command, env=environment, capture_output=True, text=True, timeout=90) + assert again.returncode == 2 + assert json.loads(again.stderr) == { + "error": "refused", + "message": "there is no unfinished index build or staging to abandon", + } + old.save("Event", {"id": 92, "value": 2}) From 17b747a2d05f7a5324ec5e992454fad1c2286610 Mon Sep 17 00:00:00 2001 From: Krzysztof Macewicz Date: Thu, 24 Sep 2026 19:16:46 +0200 Subject: [PATCH 2/3] test: an abandonment leaves a recreated object and the next staging prepares Co-Authored-By: Claude Opus 5.5 (1M context) --- python/tests/test_staging_abandon_live.py | 77 +++++++++++++++++++++++ 1 file changed, 77 insertions(+) diff --git a/python/tests/test_staging_abandon_live.py b/python/tests/test_staging_abandon_live.py index c293e53..41753ed 100644 --- a/python/tests/test_staging_abandon_live.py +++ b/python/tests/test_staging_abandon_live.py @@ -328,3 +328,80 @@ def test_the_operator_cli_abandons_a_staging(tmp_path: Path) -> None: "message": "there is no unfinished index build or staging to abandon", } old.save("Event", {"id": 92, "value": 2}) + + +def test_a_recreated_table_with_the_same_marker_is_not_this_stagings(tmp_path: Path) -> None: + """Ours is the marker *and* the object that was recorded, not the marker alone.""" + with initial("clickhouse", tmp_path) as (operator, stage, _old, roles, *_, target): + assert target == "postgres" + operator._after_step = crash_at("stage_create_Event:done") + with pytest.raises(Crash): + operator.stage(stage) + operator._after_step = quiet + (row,) = operator.store.read()["execution"]["tables"] + table, marker = copy_table(stage), row["marker"] + roles[target].command(f'DROP TABLE "{table}"') + roles[target].command(f'CREATE TABLE "{table}" (id bigint PRIMARY KEY)') + roles[target].command(f"COMMENT ON TABLE \"{table}\" IS '{marker}'") + receipt = operator.abandon().as_record() + assert receipt["tables"][0]["identity"] is None + assert exists(roles[target], table) # another object carrying the marker: left alone + + +@pytest.mark.parametrize("source", ["postgres", "clickhouse"]) +def test_after_an_abandonment_the_next_staging_prepares(source: str, tmp_path: Path) -> None: + """The customer is not stuck: a new authorization stages a new copy beside the history.""" + from copy import deepcopy + from uuid import uuid4 + + with initial(source, tmp_path) as ( + operator, + stage, + old, + roles, + model, + public, + signed, + material, + _source, + target, + ): + operator._after_step = crash_at("stage_grants:done") + with pytest.raises(Crash): + operator.stage(stage) + operator._after_step = quiet + assert operator.abandon().as_record()["outcome"] == "abandoned" + current = operator.store.read()["active_map"] + identity = uuid4().hex + prepared = deepcopy(current) + prepared["map_version"] = stage.prepared.map_version + 1 # the abandoned one is burned + fresh = { + **material(target, sde.staging_table_name(identity, 1), "next-copy"), + "lag_budget_ms": 30000, + } + prepared["groups"]["Event"].update(derived=[fresh], also_write=["next-copy"]) + next_stage = sde.load_staging_plan( + signed( + { + "kind": "sde-stage", + "protocol": 1, + "stage_id": identity, + "project_id": PROJECT, + "group": "Event", + "current": current, + "prepared": signed(prepared), + } + ), + model=model, + project_id=PROJECT, + public_key=public, + ) + receipt = operator.stage(next_stage).as_record() + assert receipt["outcome"] == "prepared" + assert exists(roles[target], sde.staging_table_name(identity, 1)) + history = operator.store.read()["stages"] + assert {record["receipt"]["outcome"] for record in history.values()} == { + "abandoned", + "prepared", + } + old.save("Event", {"id": 93, "value": 3}) From c4d5d38c04cab02fc7340ce4a8ce74c1ce06f47d Mon Sep 17 00:00:00 2001 From: Krzysztof Macewicz Date: Thu, 24 Sep 2026 19:18:04 +0200 Subject: [PATCH 3/3] test: the next staging after an abandonment before any table existed A mutation that read reused names from the receipt survived: the test abandoned only a staging whose table had been created, whose receipt names an identity. Co-Authored-By: Claude Opus 5.5 (1M context) --- python/tests/test_staging_abandon_live.py | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/python/tests/test_staging_abandon_live.py b/python/tests/test_staging_abandon_live.py index 41753ed..156d7a6 100644 --- a/python/tests/test_staging_abandon_live.py +++ b/python/tests/test_staging_abandon_live.py @@ -348,9 +348,16 @@ def test_a_recreated_table_with_the_same_marker_is_not_this_stagings(tmp_path: P assert exists(roles[target], table) # another object carrying the marker: left alone +@pytest.mark.parametrize("checkpoint", ["stage_prepared", "stage_grants:done"]) @pytest.mark.parametrize("source", ["postgres", "clickhouse"]) -def test_after_an_abandonment_the_next_staging_prepares(source: str, tmp_path: Path) -> None: - """The customer is not stuck: a new authorization stages a new copy beside the history.""" +def test_after_an_abandonment_the_next_staging_prepares( + source: str, checkpoint: str, tmp_path: Path +) -> None: + """The customer is not stuck: a new authorization stages a new copy beside the history. + + Abandoned before its table existed, a staging's receipt names no identity - and the next + staging's check for reused names reads that history. + """ from copy import deepcopy from uuid import uuid4 @@ -366,7 +373,7 @@ def test_after_an_abandonment_the_next_staging_prepares(source: str, tmp_path: P _source, target, ): - operator._after_step = crash_at("stage_grants:done") + operator._after_step = crash_at(checkpoint) with pytest.raises(Crash): operator.stage(stage) operator._after_step = quiet