From a49da0b9a930da6367f777c73247e3edffc000e6 Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:07:08 +0000 Subject: [PATCH 01/11] fix(datastore): auto-recover malformed peewee SQLite on startup Preserve a corrupt peewee-sqlite.v2.db as .corrupt- and replace it with a recovered copy so aw-server can start instead of restart-looping. Uses sqlite3 .recover when sqlite_dbpage is available; otherwise a sanitized .bail-off .dump (Debian/Ubuntu sqlite3 is built without dbpage). Reconstructs any eventmodel.bucket_id rows missing from bucketmodel. Disable with AW_SQLITE_AUTO_RECOVER=0. Git-Session-Id: 08b7 --- aw_datastore/storages/peewee.py | 2 + aw_datastore/storages/sqlite_recover.py | 357 ++++++++++++++++++++++++ tests/test_sqlite_recover.py | 142 ++++++++++ 3 files changed, 501 insertions(+) create mode 100644 aw_datastore/storages/sqlite_recover.py create mode 100644 tests/test_sqlite_recover.py diff --git a/aw_datastore/storages/peewee.py b/aw_datastore/storages/peewee.py index b003cc9..349362c 100644 --- a/aw_datastore/storages/peewee.py +++ b/aw_datastore/storages/peewee.py @@ -32,6 +32,7 @@ ) from .abstract import AbstractStorage +from .sqlite_recover import maybe_recover_malformed_sqlite logger = logging.getLogger(__name__) @@ -176,6 +177,7 @@ def __init__(self, testing: bool = True, filepath: Optional[str] = None) -> None + ".db" ) filepath = os.path.join(data_dir, filename) + maybe_recover_malformed_sqlite(filepath) self.db = _db self.db.init( filepath, diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py new file mode 100644 index 0000000..e28c579 --- /dev/null +++ b/aw_datastore/storages/sqlite_recover.py @@ -0,0 +1,357 @@ +"""Recover a malformed SQLite file so aw-server can start instead of restart-looping. + +Ubuntu/Debian sqlite3 is often built without ``SQLITE_ENABLE_DBPAGE_VTAB``, so +the shell ``.recover`` command fails with ``no such table: sqlite_dbpage``. +macOS sqlite3 usually has dbpage and ``.recover`` is the better tool (it +restored 148k events on a real corrupt ``peewee-sqlite.v2.db``). + +This module tries, in order: + +1. ``sqlite3 .recover`` when dbpage is available +2. ``sqlite3 .bail off .dump`` with ``ROLLBACK`` rewritten to ``COMMIT`` +3. Reconstruct any ``eventmodel.bucket_id`` rows missing from ``bucketmodel`` + +The original file is copied aside as ``.corrupt-`` before replacement. +Disable with ``AW_SQLITE_AUTO_RECOVER=0``. +""" + +from __future__ import annotations + +import logging +import os +import shutil +import sqlite3 +import subprocess +from datetime import datetime, timezone + +logger = logging.getLogger(__name__) + +AUTO_RECOVER_ENV = "AW_SQLITE_AUTO_RECOVER" +RECOVER_TIMEOUT_SEC = 300 + + +class SqliteRecoverError(RuntimeError): + """Raised when the database is malformed and recovery did not succeed.""" + + +def auto_recover_enabled() -> bool: + val = os.environ.get(AUTO_RECOVER_ENV, "1").strip().lower() + return val not in {"0", "false", "no", "off"} + + +def is_sqlite_healthy(path: str) -> bool: + """Return True if path is missing (will be created) or PRAGMA quick_check is ok.""" + if not os.path.exists(path) or os.path.getsize(path) == 0: + return True + try: + con = sqlite3.connect(f"file:{path}?mode=ro", uri=True) + try: + row = con.execute("PRAGMA quick_check").fetchone() + return bool(row) and str(row[0]).lower() == "ok" + finally: + con.close() + except sqlite3.Error: + return False + + +def maybe_recover_malformed_sqlite(path: str) -> str | None: + """If ``path`` is malformed, preserve it and replace with a recovered copy. + + Returns the sidecar path when recovery ran, otherwise None. + """ + if not os.path.exists(path) or os.path.getsize(path) == 0: + return None + if is_sqlite_healthy(path): + return None + if not auto_recover_enabled(): + raise SqliteRecoverError( + _manual_instructions(path) + + f"\nAuto-recovery disabled via {AUTO_RECOVER_ENV}=0." + ) + sidecar = _copy_aside(path) + logger.warning( + "SQLite database %s is malformed; original copied to %s. Attempting recovery.", + path, + sidecar, + ) + tmp_dest = path + ".recovered-tmp" + _remove_if_exists(tmp_dest) + try: + recovered = _try_recover(sidecar, tmp_dest) + if ( + not recovered + or not os.path.exists(tmp_dest) + or os.path.getsize(tmp_dest) == 0 + ): + raise SqliteRecoverError("recovery produced no database file") + _reconstruct_missing_buckets(tmp_dest) + if not is_sqlite_healthy(tmp_dest): + raise SqliteRecoverError("recovered file still fails PRAGMA quick_check") + _replace_live_db(path, tmp_dest) + except Exception as exc: + _remove_if_exists(tmp_dest) + _restore_sidecars(path, sidecar) + if isinstance(exc, SqliteRecoverError): + raise SqliteRecoverError( + f"Failed to recover {path} (original preserved at {sidecar}).\n" + + _manual_instructions(sidecar) + ) from exc + raise SqliteRecoverError( + f"Failed to recover {path} (original preserved at {sidecar}): {exc}\n" + + _manual_instructions(sidecar) + ) from exc + event_count, bucket_count = _counts(path) + logger.warning( + "SQLite recovery succeeded for %s (%s events, %s buckets). Original at %s.", + path, + event_count, + bucket_count, + sidecar, + ) + return sidecar + + +def sanitize_dump_sql(sql: str) -> str: + """Turn a ``.dump`` of a corrupt DB into SQL that can be loaded.""" + lines: list[str] = [] + for line in sql.splitlines(): + if "CORRUPTION ERROR" in line: + continue + if line.startswith("ROLLBACK;"): + lines.append("COMMIT;") + continue + lines.append(line) + return "\n".join(lines) + "\n" + + +def _sqlite_bin() -> str | None: + return shutil.which("sqlite3") + + +def _has_dbpage(sqlite_bin: str) -> bool: + try: + proc = subprocess.run( + [ + sqlite_bin, + ":memory:", + "CREATE VIRTUAL TABLE temp.t USING sqlite_dbpage;", + ], + capture_output=True, + text=True, + timeout=10, + ) + except (OSError, subprocess.TimeoutExpired): + return False + return proc.returncode == 0 + + +def _try_recover(src: str, dest: str) -> bool: + sqlite_bin = _sqlite_bin() + if sqlite_bin is None: + raise SqliteRecoverError("sqlite3 CLI not found on PATH; cannot auto-recover.") + if _has_dbpage(sqlite_bin) and _recover_with_dbpage(sqlite_bin, src, dest): + return True + return _recover_with_dump(sqlite_bin, src, dest) + + +def _recover_with_dbpage(sqlite_bin: str, src: str, dest: str) -> bool: + dump = subprocess.Popen( + [sqlite_bin, src, ".recover"], + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + load = subprocess.Popen( + [sqlite_bin, dest], + stdin=dump.stdout, + stdout=subprocess.DEVNULL, + stderr=subprocess.PIPE, + text=True, + ) + if dump.stdout is not None: + dump.stdout.close() + try: + _, load_err = load.communicate(timeout=RECOVER_TIMEOUT_SEC) + _, dump_err = dump.communicate(timeout=30) + except subprocess.TimeoutExpired: + dump.kill() + load.kill() + dump.communicate() + load.communicate() + logger.warning("sqlite3 .recover timed out") + return False + if dump.returncode != 0: + logger.info( + "sqlite3 .recover unavailable or failed (rc=%s): %s", + dump.returncode, + (dump_err or "").strip()[:300], + ) + _remove_if_exists(dest) + return False + if load.returncode != 0: + logger.info("sqlite3 .recover load failed: %s", (load_err or "").strip()[:300]) + _remove_if_exists(dest) + return False + return os.path.exists(dest) and os.path.getsize(dest) > 0 + + +def _recover_with_dump(sqlite_bin: str, src: str, dest: str) -> bool: + try: + dump = subprocess.run( + [sqlite_bin, src, ".bail off", ".dump"], + capture_output=True, + text=True, + timeout=RECOVER_TIMEOUT_SEC, + ) + except subprocess.TimeoutExpired as exc: + raise SqliteRecoverError("sqlite3 .dump timed out") from exc + sql = sanitize_dump_sql(dump.stdout or "") + if "CREATE TABLE" not in sql: + raise SqliteRecoverError( + "sqlite3 .dump produced no schema: " + (dump.stderr or "").strip()[:300] + ) + try: + load = subprocess.run( + [sqlite_bin, dest], + input=sql, + capture_output=True, + text=True, + timeout=RECOVER_TIMEOUT_SEC, + ) + except subprocess.TimeoutExpired as exc: + raise SqliteRecoverError("sqlite3 dump load timed out") from exc + if load.returncode != 0: + raise SqliteRecoverError( + f"sqlite3 dump load failed (rc={load.returncode}): " + + (load.stderr or "").strip()[:300] + ) + if dump.returncode not in (0, 1): + # rc=1 is common on corrupt dumps; rc=0 also happens with .bail off. + logger.info( + "sqlite3 .dump rc=%s stderr=%s", + dump.returncode, + (dump.stderr or "").strip()[:200], + ) + return os.path.exists(dest) and os.path.getsize(dest) > 0 + + +def _reconstruct_missing_buckets(path: str) -> int: + con = sqlite3.connect(path) + try: + tables = { + row[0] + for row in con.execute("SELECT name FROM sqlite_master WHERE type='table'") + } + if "eventmodel" not in tables or "bucketmodel" not in tables: + return 0 + missing = con.execute( + """ + SELECT DISTINCT e.bucket_id + FROM eventmodel e + LEFT JOIN bucketmodel b ON b."key" = e.bucket_id + WHERE b."key" IS NULL AND e.bucket_id IS NOT NULL + """ + ).fetchall() + now = datetime.now(timezone.utc).isoformat() + n = 0 + for (key,) in missing: + con.execute( + """ + INSERT INTO bucketmodel + ("key", id, created, name, type, client, hostname, datastr) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + key, + f"recovered-{key}", + now, + "recovered", + "unknown", + "sqlite-recover", + "unknown", + "{}", + ), + ) + n += 1 + if n: + con.commit() + logger.warning( + "Reconstructed %s missing bucket row(s) as recovered- in %s", + n, + path, + ) + return n + finally: + con.close() + + +def _counts(path: str) -> tuple[str, str]: + try: + con = sqlite3.connect(path) + try: + tables = { + row[0] + for row in con.execute( + "SELECT name FROM sqlite_master WHERE type='table'" + ) + } + events = ( + con.execute("SELECT count(*) FROM eventmodel").fetchone()[0] + if "eventmodel" in tables + else "n/a" + ) + buckets = ( + con.execute("SELECT count(*) FROM bucketmodel").fetchone()[0] + if "bucketmodel" in tables + else "n/a" + ) + return str(events), str(buckets) + finally: + con.close() + except sqlite3.Error: + return "n/a", "n/a" + + +def _copy_aside(path: str) -> str: + ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + sidecar = f"{path}.corrupt-{ts}" + shutil.copy2(path, sidecar) + for suffix in ("-wal", "-shm", "-journal"): + extra = path + suffix + if os.path.exists(extra) and os.path.getsize(extra) > 0: + shutil.copy2(extra, sidecar + suffix) + return sidecar + + +def _replace_live_db(path: str, recovered: str) -> None: + # Drop WAL/SHM first so SQLite cannot apply the old log to the new file. + for suffix in ("-wal", "-shm", "-journal"): + _remove_if_exists(path + suffix) + os.replace(recovered, path) + + +def _restore_sidecars(path: str, sidecar: str) -> None: + """Put WAL/SHM back if replacement failed after they were deleted.""" + for suffix in ("-wal", "-shm", "-journal"): + src = sidecar + suffix + dest = path + suffix + if os.path.exists(src) and not os.path.exists(dest): + shutil.copy2(src, dest) + + +def _remove_if_exists(path: str) -> None: + try: + os.remove(path) + except FileNotFoundError: + return + + +def _manual_instructions(path: str) -> str: + return ( + "SQLite database is malformed. Preserve the file and recover with:\n" + f" sqlite3 {path} '.recover' | sqlite3 recovered.db\n" + "If .recover fails with 'no such table: sqlite_dbpage', use:\n" + f" sqlite3 {path} '.bail off' '.dump' | sed '/CORRUPTION ERROR/d; s/^ROLLBACK;/COMMIT;/'" + " | sqlite3 recovered.db\n" + "Then replace the original database with recovered.db." + ) diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py new file mode 100644 index 0000000..7acd52f --- /dev/null +++ b/tests/test_sqlite_recover.py @@ -0,0 +1,142 @@ +import glob +import os +import sqlite3 +from datetime import datetime, timedelta, timezone + +import pytest +from aw_core.models import Event +from aw_datastore.storages.peewee import PeeweeStorage, _db +from aw_datastore.storages.sqlite_recover import ( + AUTO_RECOVER_ENV, + SqliteRecoverError, + is_sqlite_healthy, + maybe_recover_malformed_sqlite, + sanitize_dump_sql, +) + + +def _checkpoint_and_close(path: str) -> None: + if not _db.is_closed(): + _db.close() + con = sqlite3.connect(path) + con.execute("PRAGMA wal_checkpoint(TRUNCATE)") + con.close() + for suffix in ("-wal", "-shm"): + extra = path + suffix + if os.path.exists(extra) and os.path.getsize(extra) == 0: + os.remove(extra) + + +def _xor_corrupt(path: str, offset: int = 4096, length: int = 200) -> None: + data = bytearray(open(path, "rb").read()) + end = min(offset + length, len(data)) + assert end > offset, "fixture too small to corrupt at offset" + for i in range(offset, end): + data[i] ^= 0xFF + open(path, "wb").write(data) + + +def _seed_peewee_db(path: str, n_events: int = 30) -> None: + if not _db.is_closed(): + _db.close() + store = PeeweeStorage(testing=True, filepath=path) + store.create_bucket( + "aw-watcher-window", + "currentwindow", + "aw-watcher-window", + "host", + datetime.now(timezone.utc).isoformat(), + name="window", + ) + now = datetime.now(timezone.utc) + events = [ + Event( + timestamp=now + timedelta(seconds=i), + duration=timedelta(seconds=1), + data={"app": f"t{i}", "title": "x"}, + ) + for i in range(n_events) + ] + store.insert_many("aw-watcher-window", events) + assert store.get_eventcount("aw-watcher-window") == n_events + _checkpoint_and_close(path) + + +def test_sanitize_dump_sql_rewrites_rollback_and_drops_corruption_markers(): + sql = ( + "PRAGMA foreign_keys=OFF;\n" + "BEGIN TRANSACTION;\n" + "CREATE TABLE t(a);\n" + "/****** CORRUPTION ERROR *******/\n" + "INSERT INTO t VALUES(1);\n" + "ROLLBACK; -- due to errors\n" + ) + out = sanitize_dump_sql(sql) + assert "CORRUPTION ERROR" not in out + assert "ROLLBACK;" not in out + assert "COMMIT;" in out + assert "INSERT INTO t VALUES(1);" in out + + +def test_is_sqlite_healthy_missing_and_ok(tmp_path): + missing = str(tmp_path / "nope.db") + assert is_sqlite_healthy(missing) + path = str(tmp_path / "ok.db") + con = sqlite3.connect(path) + con.execute("CREATE TABLE t(a)") + con.commit() + con.close() + assert is_sqlite_healthy(path) + + +def test_maybe_recover_noop_on_healthy_db(tmp_path): + path = str(tmp_path / "ok.db") + con = sqlite3.connect(path) + con.execute("CREATE TABLE t(a)") + con.commit() + con.close() + assert maybe_recover_malformed_sqlite(path) is None + assert glob.glob(path + ".corrupt-*") == [] + + +def test_peewee_startup_recovers_xor_corrupt_db(tmp_path): + path = str(tmp_path / "peewee-sqlite.v2.db") + _seed_peewee_db(path, n_events=30) + assert is_sqlite_healthy(path) + _xor_corrupt(path) + assert not is_sqlite_healthy(path) + + if not _db.is_closed(): + _db.close() + store = PeeweeStorage(testing=True, filepath=path) + try: + sidecars = glob.glob(path + ".corrupt-*") + assert len(sidecars) == 1 + assert os.path.exists(sidecars[0]) + assert is_sqlite_healthy(path) + buckets = store.buckets() + assert buckets, "recovered DB should expose at least one bucket" + # Original string id survives when the bucket row is intact; otherwise + # sqlite_recover reconstructs recovered-. + bucket_id = ( + "aw-watcher-window" + if "aw-watcher-window" in buckets + else next(iter(buckets)) + ) + assert store.get_eventcount(bucket_id) == 30 + events = store.get_events(bucket_id, limit=100) + apps = {e.data["app"] for e in events} + assert apps == {f"t{i}" for i in range(30)} + finally: + if not _db.is_closed(): + _db.close() + + +def test_auto_recover_disabled_raises(tmp_path, monkeypatch): + path = str(tmp_path / "peewee-sqlite.v2.db") + _seed_peewee_db(path, n_events=5) + _xor_corrupt(path) + monkeypatch.setenv(AUTO_RECOVER_ENV, "0") + with pytest.raises(SqliteRecoverError, match="Auto-recovery disabled"): + maybe_recover_malformed_sqlite(path) + assert glob.glob(path + ".corrupt-*") == [] From c21fc20277d838b5c11da17891987548d84f013e Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:12:14 +0000 Subject: [PATCH 02/11] fix(datastore): Python fallback when sqlite3 CLI is missing Windows CI has no sqlite3.exe, so CLI .recover/.dump never ran and the xor-corrupt fixture failed. Copy schema plus surviving rows through the stdlib sqlite3 module, reopening poisoned connections after DatabaseError. Git-Session-Id: 08b7 --- aw_datastore/storages/sqlite_recover.py | 145 +++++++++++++++++++++++- tests/test_sqlite_recover.py | 28 +++++ 2 files changed, 167 insertions(+), 6 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index e28c579..02e9e9d 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -9,7 +9,8 @@ 1. ``sqlite3 .recover`` when dbpage is available 2. ``sqlite3 .bail off .dump`` with ``ROLLBACK`` rewritten to ``COMMIT`` -3. Reconstruct any ``eventmodel.bucket_id`` rows missing from ``bucketmodel`` +3. Pure-Python schema + row copy (Windows CI has no sqlite3 CLI) +4. Reconstruct any ``eventmodel.bucket_id`` rows missing from ``bucketmodel`` The original file is copied aside as ``.corrupt-`` before replacement. Disable with ``AW_SQLITE_AUTO_RECOVER=0``. @@ -93,7 +94,7 @@ def maybe_recover_malformed_sqlite(path: str) -> str | None: _restore_sidecars(path, sidecar) if isinstance(exc, SqliteRecoverError): raise SqliteRecoverError( - f"Failed to recover {path} (original preserved at {sidecar}).\n" + f"Failed to recover {path} (original preserved at {sidecar}): {exc}\n" + _manual_instructions(sidecar) ) from exc raise SqliteRecoverError( @@ -147,11 +148,21 @@ def _has_dbpage(sqlite_bin: str) -> bool: def _try_recover(src: str, dest: str) -> bool: sqlite_bin = _sqlite_bin() - if sqlite_bin is None: - raise SqliteRecoverError("sqlite3 CLI not found on PATH; cannot auto-recover.") - if _has_dbpage(sqlite_bin) and _recover_with_dbpage(sqlite_bin, src, dest): + if sqlite_bin is not None: + if _has_dbpage(sqlite_bin) and _recover_with_dbpage(sqlite_bin, src, dest): + return True + try: + if _recover_with_dump(sqlite_bin, src, dest): + return True + except SqliteRecoverError as exc: + logger.info("CLI dump recover failed, trying Python fallback: %s", exc) + _remove_if_exists(dest) + if _recover_with_python(src, dest): return True - return _recover_with_dump(sqlite_bin, src, dest) + raise SqliteRecoverError( + "sqlite3 CLI recover/dump failed or is missing, and the Python " + "row-copy fallback produced no database" + ) def _recover_with_dbpage(sqlite_bin: str, src: str, dest: str) -> bool: @@ -235,6 +246,128 @@ def _recover_with_dump(sqlite_bin: str, src: str, dest: str) -> bool: return os.path.exists(dest) and os.path.getsize(dest) > 0 +def _ident(name: str) -> str: + return '"' + str(name).replace('"', '""') + '"' + + +def _open_ro(path: str) -> sqlite3.Connection: + return sqlite3.connect(f"file:{path}?mode=ro", uri=True) + + +def _recover_with_python(src: str, dest: str) -> bool: + """Copy schema + surviving rows without the sqlite3 CLI. + + Used on Windows CI (no sqlite3.exe) and as a last resort when CLI recover + fails. A poisoned connection is reopened after each DatabaseError. + """ + _remove_if_exists(dest) + src_con = _open_ro(src) + dest_con = sqlite3.connect(dest) + try: + try: + tables = src_con.execute( + "SELECT name, sql FROM sqlite_master " + "WHERE type='table' AND name NOT LIKE 'sqlite_%' AND sql IS NOT NULL" + ).fetchall() + indexes = src_con.execute( + "SELECT sql FROM sqlite_master WHERE type='index' AND sql IS NOT NULL" + ).fetchall() + except sqlite3.Error as exc: + logger.info("Python recover cannot read sqlite_master: %s", exc) + dest_con.close() + _remove_if_exists(dest) + return False + if not tables: + dest_con.close() + _remove_if_exists(dest) + return False + for _name, sql in tables: + dest_con.execute(sql) + dest_con.commit() + src_con.close() + src_con = None + copied = 0 + for name, _sql in tables: + copied += _copy_table_rows(src, dest_con, name) + for (sql,) in indexes: + try: + dest_con.execute(sql) + except sqlite3.Error: + continue + dest_con.commit() + logger.info("Python recover copied %s row(s) from %s", copied, src) + return os.path.exists(dest) and os.path.getsize(dest) > 0 + except sqlite3.Error as exc: + logger.info("Python recover failed: %s", exc) + dest_con.close() + _remove_if_exists(dest) + return False + finally: + if src_con is not None: + src_con.close() + try: + dest_con.close() + except sqlite3.Error: + pass + + +def _copy_table_rows(src: str, dest_con: sqlite3.Connection, table: str) -> int: + ident = _ident(table) + src_con = _open_ro(src) + copied = 0 + try: + cols = [row[1] for row in src_con.execute(f"PRAGMA table_info({ident})")] + if not cols: + return 0 + col_list = ", ".join(_ident(c) for c in cols) + placeholders = ", ".join("?" for _ in cols) + insert_sql = ( + f"INSERT OR IGNORE INTO {ident} ({col_list}) VALUES ({placeholders})" + ) + pk_cols = [ + row[1] for row in src_con.execute(f"PRAGMA table_info({ident})") if row[5] + ] + try: + rows = src_con.execute(f"SELECT {col_list} FROM {ident}").fetchall() + dest_con.executemany(insert_sql, rows) + dest_con.commit() + return len(rows) + except sqlite3.Error: + src_con.close() + src_con = _open_ro(src) + if not pk_cols: + return 0 + pk = pk_cols[0] + pk_ident = _ident(pk) + try: + ids = [row[0] for row in src_con.execute(f"SELECT {pk_ident} FROM {ident}")] + except sqlite3.Error: + src_con.close() + src_con = _open_ro(src) + # Integer primary keys are the peewee/eventmodel case. + try: + mx = src_con.execute(f"SELECT max({pk_ident}) FROM {ident}").fetchone() + ids = list(range(1, int(mx[0]) + 1)) if mx and mx[0] else [] + except sqlite3.Error: + return 0 + for key in ids: + try: + row = src_con.execute( + f"SELECT {col_list} FROM {ident} WHERE {pk_ident}=?", + (key,), + ).fetchone() + if row: + dest_con.execute(insert_sql, row) + copied += 1 + except sqlite3.Error: + src_con.close() + src_con = _open_ro(src) + dest_con.commit() + return copied + finally: + src_con.close() + + def _reconstruct_missing_buckets(path: str) -> int: con = sqlite3.connect(path) try: diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index 7acd52f..7107914 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -132,6 +132,34 @@ def test_peewee_startup_recovers_xor_corrupt_db(tmp_path): _db.close() +def test_python_fallback_without_sqlite_cli(tmp_path, monkeypatch): + path = str(tmp_path / "peewee-sqlite.v2.db") + _seed_peewee_db(path, n_events=30) + _xor_corrupt(path) + monkeypatch.setattr( + "aw_datastore.storages.sqlite_recover._sqlite_bin", lambda: None + ) + sidecar = maybe_recover_malformed_sqlite(path) + assert sidecar + assert os.path.exists(sidecar) + assert is_sqlite_healthy(path) + if not _db.is_closed(): + _db.close() + store = PeeweeStorage(testing=True, filepath=path) + try: + buckets = store.buckets() + assert buckets + bucket_id = ( + "aw-watcher-window" + if "aw-watcher-window" in buckets + else next(iter(buckets)) + ) + assert store.get_eventcount(bucket_id) == 30 + finally: + if not _db.is_closed(): + _db.close() + + def test_auto_recover_disabled_raises(tmp_path, monkeypatch): path = str(tmp_path / "peewee-sqlite.v2.db") _seed_peewee_db(path, n_events=5) From bbc6cbd7fa355774089da77040bd6960128fa078 Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:20:24 +0000 Subject: [PATCH 03/11] fix(datastore): satisfy mypy on python recover src_con close Git-Session-Id: 08b7 --- aw_datastore/storages/sqlite_recover.py | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index 02e9e9d..0e54db1 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -263,6 +263,7 @@ def _recover_with_python(src: str, dest: str) -> bool: _remove_if_exists(dest) src_con = _open_ro(src) dest_con = sqlite3.connect(dest) + src_closed = False try: try: tables = src_con.execute( @@ -274,18 +275,16 @@ def _recover_with_python(src: str, dest: str) -> bool: ).fetchall() except sqlite3.Error as exc: logger.info("Python recover cannot read sqlite_master: %s", exc) - dest_con.close() _remove_if_exists(dest) return False if not tables: - dest_con.close() _remove_if_exists(dest) return False for _name, sql in tables: dest_con.execute(sql) dest_con.commit() src_con.close() - src_con = None + src_closed = True copied = 0 for name, _sql in tables: copied += _copy_table_rows(src, dest_con, name) @@ -299,11 +298,10 @@ def _recover_with_python(src: str, dest: str) -> bool: return os.path.exists(dest) and os.path.getsize(dest) > 0 except sqlite3.Error as exc: logger.info("Python recover failed: %s", exc) - dest_con.close() _remove_if_exists(dest) return False finally: - if src_con is not None: + if not src_closed: src_con.close() try: dest_con.close() From e93b30355186c50c908816ec928f99923f9b6e0c Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:23:16 +0000 Subject: [PATCH 04/11] fix(datastore): keep recovered events reachable after incomplete dumps Create bucketmodel when a corrupt dump salvages eventmodel but omits the bucket catalog, then refuse replacement if events would still be unreachable. Preserve the original file mode on os.replace so recovery does not widen local read access under a permissive umask. --- aw_datastore/storages/sqlite_recover.py | 88 +++++++++++++++++++++++-- tests/test_sqlite_recover.py | 74 +++++++++++++++++++++ 2 files changed, 157 insertions(+), 5 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index 0e54db1..10f1407 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -22,6 +22,7 @@ import os import shutil import sqlite3 +import stat import subprocess from datetime import datetime, timezone @@ -88,6 +89,7 @@ def maybe_recover_malformed_sqlite(path: str) -> str | None: _reconstruct_missing_buckets(tmp_dest) if not is_sqlite_healthy(tmp_dest): raise SqliteRecoverError("recovered file still fails PRAGMA quick_check") + _assert_recovered_schema(tmp_dest) _replace_live_db(path, tmp_dest) except Exception as exc: _remove_if_exists(tmp_dest) @@ -366,15 +368,49 @@ def _copy_table_rows(src: str, dest_con: sqlite3.Connection, table: str) -> int: src_con.close() +_BUCKETMODEL_DDL = ( + 'CREATE TABLE "bucketmodel" (' + '"key" INTEGER NOT NULL PRIMARY KEY, ' + '"id" VARCHAR(255) NOT NULL, ' + '"created" DATETIME NOT NULL, ' + '"name" VARCHAR(255), ' + '"type" VARCHAR(255) NOT NULL, ' + '"client" VARCHAR(255) NOT NULL, ' + '"hostname" VARCHAR(255) NOT NULL, ' + '"datastr" VARCHAR(255))' +) +_BUCKETMODEL_ID_INDEX_DDL = ( + 'CREATE UNIQUE INDEX IF NOT EXISTS "bucketmodel_id" ON "bucketmodel" ("id")' +) + + +def _table_names(con: sqlite3.Connection) -> set[str]: + return { + row[0] + for row in con.execute("SELECT name FROM sqlite_master WHERE type='table'") + } + + +def _ensure_bucketmodel(con: sqlite3.Connection) -> None: + """Create the peewee bucket catalog if a dump salvaged events but not buckets.""" + con.execute(_BUCKETMODEL_DDL) + con.execute(_BUCKETMODEL_ID_INDEX_DDL) + + def _reconstruct_missing_buckets(path: str) -> int: con = sqlite3.connect(path) try: - tables = { - row[0] - for row in con.execute("SELECT name FROM sqlite_master WHERE type='table'") - } - if "eventmodel" not in tables or "bucketmodel" not in tables: + tables = _table_names(con) + if "eventmodel" not in tables: return 0 + if "bucketmodel" not in tables: + _ensure_bucketmodel(con) + con.commit() + logger.warning( + "Recovered file %s had eventmodel but no bucketmodel; " + "created an empty bucket catalog before reconstructing keys", + path, + ) missing = con.execute( """ SELECT DISTINCT e.bucket_id @@ -416,6 +452,38 @@ def _reconstruct_missing_buckets(path: str) -> int: con.close() +def _assert_recovered_schema(path: str) -> None: + """Refuse to install a recovered file whose salvaged events would be unreachable. + + ``PRAGMA quick_check`` only validates pages. A dump can keep ``eventmodel`` + and omit ``bucketmodel``; peewee would then create an empty catalog and + hide the recovered events. + """ + con = sqlite3.connect(path) + try: + tables = _table_names(con) + if "eventmodel" not in tables: + return + if "bucketmodel" not in tables: + raise SqliteRecoverError("recovered file has eventmodel but no bucketmodel") + missing = con.execute( + """ + SELECT DISTINCT e.bucket_id + FROM eventmodel e + LEFT JOIN bucketmodel b ON b."key" = e.bucket_id + WHERE b."key" IS NULL AND e.bucket_id IS NOT NULL + """ + ).fetchall() + if missing: + raise SqliteRecoverError( + "recovered file has events whose buckets could not be reconstructed" + ) + except sqlite3.Error as exc: + raise SqliteRecoverError(f"recovered file schema check failed: {exc}") from exc + finally: + con.close() + + def _counts(path: str) -> tuple[str, str]: try: con = sqlite3.connect(path) @@ -454,10 +522,20 @@ def _copy_aside(path: str) -> str: return sidecar +def _preserve_mode(src: str, dest: str) -> None: + """Copy permission bits so replacement does not widen access under umask.""" + try: + mode = stat.S_IMODE(os.stat(src).st_mode) + os.chmod(dest, mode) + except OSError as exc: + logger.info("Could not copy mode from %s to %s: %s", src, dest, exc) + + def _replace_live_db(path: str, recovered: str) -> None: # Drop WAL/SHM first so SQLite cannot apply the old log to the new file. for suffix in ("-wal", "-shm", "-journal"): _remove_if_exists(path + suffix) + _preserve_mode(path, recovered) os.replace(recovered, path) diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index 7107914..af8f04b 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -1,6 +1,7 @@ import glob import os import sqlite3 +import stat from datetime import datetime, timedelta, timezone import pytest @@ -9,6 +10,8 @@ from aw_datastore.storages.sqlite_recover import ( AUTO_RECOVER_ENV, SqliteRecoverError, + _assert_recovered_schema, + _reconstruct_missing_buckets, is_sqlite_healthy, maybe_recover_malformed_sqlite, sanitize_dump_sql, @@ -168,3 +171,74 @@ def test_auto_recover_disabled_raises(tmp_path, monkeypatch): with pytest.raises(SqliteRecoverError, match="Auto-recovery disabled"): maybe_recover_malformed_sqlite(path) assert glob.glob(path + ".corrupt-*") == [] + + +def _eventmodel_only_db(path: str, bucket_key: int = 7, n_events: int = 2) -> None: + con = sqlite3.connect(path) + con.execute( + 'CREATE TABLE "eventmodel" (' + '"id" INTEGER NOT NULL PRIMARY KEY, ' + '"bucket_id" INTEGER NOT NULL, ' + '"timestamp" DATETIME NOT NULL, ' + '"duration" DECIMAL(10, 5) NOT NULL, ' + '"datastr" VARCHAR(255) NOT NULL)' + ) + for i in range(n_events): + con.execute( + "INSERT INTO eventmodel (id, bucket_id, timestamp, duration, datastr) " + "VALUES (?, ?, ?, ?, ?)", + (i + 1, bucket_key, f"2020-01-01 00:00:0{i}+00:00", 1, "{}"), + ) + con.commit() + con.close() + + +def test_reconstruct_creates_bucketmodel_when_missing(tmp_path): + path = str(tmp_path / "partial.db") + _eventmodel_only_db(path, bucket_key=7, n_events=2) + assert _reconstruct_missing_buckets(path) == 1 + con = sqlite3.connect(path) + try: + tables = { + row[0] + for row in con.execute("SELECT name FROM sqlite_master WHERE type='table'") + } + assert "bucketmodel" in tables + row = con.execute('SELECT "key", id, type, client FROM bucketmodel').fetchone() + assert row[0] == 7 + assert row[1] == "recovered-7" + assert row[2] == "unknown" + assert row[3] == "sqlite-recover" + assert con.execute("SELECT count(*) FROM eventmodel").fetchone()[0] == 2 + missing = con.execute( + """ + SELECT count(*) FROM eventmodel e + LEFT JOIN bucketmodel b ON b."key" = e.bucket_id + WHERE b."key" IS NULL + """ + ).fetchone()[0] + assert missing == 0 + finally: + con.close() + _assert_recovered_schema(path) + + +def test_assert_recovered_schema_rejects_events_without_buckets(tmp_path): + path = str(tmp_path / "partial.db") + _eventmodel_only_db(path) + with pytest.raises(SqliteRecoverError, match="no bucketmodel"): + _assert_recovered_schema(path) + + +@pytest.mark.skipif( + os.name == "nt", reason="POSIX permission bits are not preserved on Windows" +) +def test_recovery_preserves_restrictive_mode(tmp_path): + path = str(tmp_path / "peewee-sqlite.v2.db") + _seed_peewee_db(path, n_events=5) + os.chmod(path, 0o600) + _xor_corrupt(path) + sidecar = maybe_recover_malformed_sqlite(path) + assert sidecar + assert is_sqlite_healthy(path) + assert stat.S_IMODE(os.stat(path).st_mode) == 0o600 From a053ccc53ed332798147cca5db9c730eb4cf618c Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:27:36 +0000 Subject: [PATCH 05/11] test(datastore): sidecar glob should ignore WAL/SHM copies Windows leaves non-empty -wal/-shm next to the .corrupt- copy, so glob(path + '.corrupt-*') matched three files and tripped assert len==1. Git-Session-Id: 08b7 --- tests/test_sqlite_recover.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index af8f04b..ae594db 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -30,6 +30,12 @@ def _checkpoint_and_close(path: str) -> None: os.remove(extra) +def _corrupt_sidecars(path: str): + """Sidecar DB copies only — WAL/SHM use the same prefix on Windows.""" + skip = ("-wal", "-shm", "-journal") + return [p for p in glob.glob(path + ".corrupt-*") if not p.endswith(skip)] + + def _xor_corrupt(path: str, offset: int = 4096, length: int = 200) -> None: data = bytearray(open(path, "rb").read()) end = min(offset + length, len(data)) @@ -99,7 +105,7 @@ def test_maybe_recover_noop_on_healthy_db(tmp_path): con.commit() con.close() assert maybe_recover_malformed_sqlite(path) is None - assert glob.glob(path + ".corrupt-*") == [] + assert _corrupt_sidecars(path) == [] def test_peewee_startup_recovers_xor_corrupt_db(tmp_path): @@ -113,7 +119,7 @@ def test_peewee_startup_recovers_xor_corrupt_db(tmp_path): _db.close() store = PeeweeStorage(testing=True, filepath=path) try: - sidecars = glob.glob(path + ".corrupt-*") + sidecars = _corrupt_sidecars(path) assert len(sidecars) == 1 assert os.path.exists(sidecars[0]) assert is_sqlite_healthy(path) @@ -170,7 +176,7 @@ def test_auto_recover_disabled_raises(tmp_path, monkeypatch): monkeypatch.setenv(AUTO_RECOVER_ENV, "0") with pytest.raises(SqliteRecoverError, match="Auto-recovery disabled"): maybe_recover_malformed_sqlite(path) - assert glob.glob(path + ".corrupt-*") == [] + assert _corrupt_sidecars(path) == [] def _eventmodel_only_db(path: str, bucket_key: int = 7, n_events: int = 2) -> None: From ed07f01bc30e12a3fc4dfdb1c9be0584d57e0b50 Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:52:11 +0000 Subject: [PATCH 06/11] fix(datastore): secure recovery file before writing Git-Session-Id: 5f4c1144-c5d9-5e51-a701-52c1c2fa45a9 --- aw_datastore/storages/sqlite_recover.py | 37 +++++++++++++++---------- tests/test_sqlite_recover.py | 36 ++++++++++++++++++++++++ 2 files changed, 59 insertions(+), 14 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index 10f1407..dea090d 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -79,6 +79,7 @@ def maybe_recover_malformed_sqlite(path: str) -> str | None: tmp_dest = path + ".recovered-tmp" _remove_if_exists(tmp_dest) try: + _prepare_recovery_file(path, tmp_dest) recovered = _try_recover(sidecar, tmp_dest) if ( not recovered @@ -158,7 +159,7 @@ def _try_recover(src: str, dest: str) -> bool: return True except SqliteRecoverError as exc: logger.info("CLI dump recover failed, trying Python fallback: %s", exc) - _remove_if_exists(dest) + _truncate(dest) if _recover_with_python(src, dest): return True raise SqliteRecoverError( @@ -199,11 +200,11 @@ def _recover_with_dbpage(sqlite_bin: str, src: str, dest: str) -> bool: dump.returncode, (dump_err or "").strip()[:300], ) - _remove_if_exists(dest) + _truncate(dest) return False if load.returncode != 0: logger.info("sqlite3 .recover load failed: %s", (load_err or "").strip()[:300]) - _remove_if_exists(dest) + _truncate(dest) return False return os.path.exists(dest) and os.path.getsize(dest) > 0 @@ -262,7 +263,7 @@ def _recover_with_python(src: str, dest: str) -> bool: Used on Windows CI (no sqlite3.exe) and as a last resort when CLI recover fails. A poisoned connection is reopened after each DatabaseError. """ - _remove_if_exists(dest) + _truncate(dest) src_con = _open_ro(src) dest_con = sqlite3.connect(dest) src_closed = False @@ -277,10 +278,10 @@ def _recover_with_python(src: str, dest: str) -> bool: ).fetchall() except sqlite3.Error as exc: logger.info("Python recover cannot read sqlite_master: %s", exc) - _remove_if_exists(dest) + _truncate(dest) return False if not tables: - _remove_if_exists(dest) + _truncate(dest) return False for _name, sql in tables: dest_con.execute(sql) @@ -300,7 +301,7 @@ def _recover_with_python(src: str, dest: str) -> bool: return os.path.exists(dest) and os.path.getsize(dest) > 0 except sqlite3.Error as exc: logger.info("Python recover failed: %s", exc) - _remove_if_exists(dest) + _truncate(dest) return False finally: if not src_closed: @@ -522,20 +523,23 @@ def _copy_aside(path: str) -> str: return sidecar -def _preserve_mode(src: str, dest: str) -> None: - """Copy permission bits so replacement does not widen access under umask.""" +def _prepare_recovery_file(src: str, dest: str) -> None: + """Create an empty recovery file with the live database's permissions.""" + mode = stat.S_IMODE(os.stat(src).st_mode) + fd = os.open(dest, os.O_WRONLY | os.O_CREAT | os.O_EXCL, mode) try: - mode = stat.S_IMODE(os.stat(src).st_mode) - os.chmod(dest, mode) - except OSError as exc: - logger.info("Could not copy mode from %s to %s: %s", src, dest, exc) + # The process umask may have removed bits requested above. Applying the + # exact original mode before recovery starts also makes chmod failures + # fail closed, before any activity data is written. + os.fchmod(fd, mode) + finally: + os.close(fd) def _replace_live_db(path: str, recovered: str) -> None: # Drop WAL/SHM first so SQLite cannot apply the old log to the new file. for suffix in ("-wal", "-shm", "-journal"): _remove_if_exists(path + suffix) - _preserve_mode(path, recovered) os.replace(recovered, path) @@ -548,6 +552,11 @@ def _restore_sidecars(path: str, sidecar: str) -> None: shutil.copy2(src, dest) +def _truncate(path: str) -> None: + with open(path, "wb"): + pass + + def _remove_if_exists(path: str) -> None: try: os.remove(path) diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index ae594db..e02ac2a 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -11,6 +11,7 @@ AUTO_RECOVER_ENV, SqliteRecoverError, _assert_recovered_schema, + _prepare_recovery_file, _reconstruct_missing_buckets, is_sqlite_healthy, maybe_recover_malformed_sqlite, @@ -236,6 +237,41 @@ def test_assert_recovered_schema_rejects_events_without_buckets(tmp_path): _assert_recovered_schema(path) +@pytest.mark.skipif( + os.name == "nt", reason="POSIX permission bits are not preserved on Windows" +) +def test_prepare_recovery_file_is_restrictive_before_data_is_written(tmp_path): + source = str(tmp_path / "source.db") + recovered = str(tmp_path / "recovered.db") + open(source, "wb").close() + os.chmod(source, 0o600) + old_umask = os.umask(0) + try: + _prepare_recovery_file(source, recovered) + finally: + os.umask(old_umask) + assert os.path.getsize(recovered) == 0 + assert stat.S_IMODE(os.stat(recovered).st_mode) == 0o600 + + +@pytest.mark.skipif( + os.name == "nt", reason="POSIX permission bits are not preserved on Windows" +) +def test_prepare_recovery_file_fails_closed_on_fchmod_error(tmp_path, monkeypatch): + source = str(tmp_path / "source.db") + recovered = str(tmp_path / "recovered.db") + open(source, "wb").close() + os.chmod(source, 0o600) + + def fail_fchmod(_fd, _mode): + raise PermissionError("denied") + + monkeypatch.setattr(os, "fchmod", fail_fchmod) + with pytest.raises(PermissionError, match="denied"): + _prepare_recovery_file(source, recovered) + assert os.path.getsize(recovered) == 0 + + @pytest.mark.skipif( os.name == "nt", reason="POSIX permission bits are not preserved on Windows" ) From 6c6209b502df99c116a018ac983ccbfedb73715f Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 20:04:39 +0000 Subject: [PATCH 07/11] fix(datastore): support recovery file setup on Windows Git-Session-Id: 5f4c1144-c5d9-5e51-a701-52c1c2fa45a9 --- aw_datastore/storages/sqlite_recover.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index dea090d..7cfb98e 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -528,10 +528,12 @@ def _prepare_recovery_file(src: str, dest: str) -> None: mode = stat.S_IMODE(os.stat(src).st_mode) fd = os.open(dest, os.O_WRONLY | os.O_CREAT | os.O_EXCL, mode) try: - # The process umask may have removed bits requested above. Applying the - # exact original mode before recovery starts also makes chmod failures - # fail closed, before any activity data is written. - os.fchmod(fd, mode) + # POSIX umask may have removed bits requested above. Applying the exact + # original mode before recovery starts also makes permission failures + # fail closed, before any activity data is written. Windows has no + # fchmod and does not preserve POSIX permission bits. + if os.name != "nt": + os.fchmod(fd, mode) finally: os.close(fd) From c544da7d2614da0ddf8e5d7ced091f58f3df4bd3 Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 20:10:42 +0000 Subject: [PATCH 08/11] fix(datastore): typecheck fchmod on Windows Git-Session-Id: 5f4c1144-c5d9-5e51-a701-52c1c2fa45a9 --- aw_datastore/storages/sqlite_recover.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index 7cfb98e..f1dc6e4 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -533,7 +533,7 @@ def _prepare_recovery_file(src: str, dest: str) -> None: # fail closed, before any activity data is written. Windows has no # fchmod and does not preserve POSIX permission bits. if os.name != "nt": - os.fchmod(fd, mode) + os.fchmod(fd, mode) # type: ignore[attr-defined] finally: os.close(fd) From 9f76f5b8eb97ddf267610aa1bc9976119a7963ce Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 16 Sep 2026 10:20:49 +0000 Subject: [PATCH 09/11] fix(datastore): unique corrupt sidecars; don't fail after recovery O_EXCL + microsecond timestamps so a same-second retry cannot overwrite the preserved original. _counts swallows OSError so a post-replace stat/connect hiccup cannot restart-loop a recovered server. Git-Session-Id: cdb864f9-bb65-5b97-b145-103ee098a511 --- aw_datastore/storages/sqlite_recover.py | 24 ++++++++++++++-- tests/test_sqlite_recover.py | 38 +++++++++++++++++++++++++ 2 files changed, 59 insertions(+), 3 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index f1dc6e4..83559ce 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -508,13 +508,31 @@ def _counts(path: str) -> tuple[str, str]: return str(events), str(buckets) finally: con.close() - except sqlite3.Error: + except (OSError, sqlite3.Error): return "n/a", "n/a" +def _unique_path(prefix: str) -> str: + """Return ``prefix`` or ``prefix-N`` that does not already exist. + + Uses ``O_EXCL`` so two recoveries in the same timestamp cannot clobber + each other's sidecar (and the original corrupt file it holds). + """ + n = 0 + while True: + candidate = prefix if n == 0 else f"{prefix}-{n}" + try: + fd = os.open(candidate, os.O_CREAT | os.O_EXCL | os.O_WRONLY) + except FileExistsError: + n += 1 + continue + os.close(fd) + return candidate + + def _copy_aside(path: str) -> str: - ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") - sidecar = f"{path}.corrupt-{ts}" + ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S%fZ") + sidecar = _unique_path(f"{path}.corrupt-{ts}") shutil.copy2(path, sidecar) for suffix in ("-wal", "-shm", "-journal"): extra = path + suffix diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index e02ac2a..6e043b5 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -11,6 +11,8 @@ AUTO_RECOVER_ENV, SqliteRecoverError, _assert_recovered_schema, + _copy_aside, + _counts, _prepare_recovery_file, _reconstruct_missing_buckets, is_sqlite_healthy, @@ -72,6 +74,42 @@ def _seed_peewee_db(path: str, n_events: int = 30) -> None: _checkpoint_and_close(path) +def test_copy_aside_does_not_overwrite_a_previous_sidecar(tmp_path, monkeypatch): + """Same-second recovery must keep the first original, not clobber it.""" + import aw_datastore.storages.sqlite_recover as recover + + class FrozenDateTime(datetime): + @classmethod + def now(cls, tz=None): + return datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) + + monkeypatch.setattr(recover, "datetime", FrozenDateTime) + path = str(tmp_path / "peewee-sqlite.v2.db") + with open(path, "wb") as fh: + fh.write(b"v1") + first = _copy_aside(path) + with open(path, "wb") as fh: + fh.write(b"v2") + second = _copy_aside(path) + assert first != second + with open(first, "rb") as fh: + assert fh.read() == b"v1" + with open(second, "rb") as fh: + assert fh.read() == b"v2" + + +def test_counts_oserror_does_not_raise(tmp_path, monkeypatch): + path = str(tmp_path / "x.db") + with open(path, "wb") as fh: + fh.write(b"x") + + def boom(*_args, **_kwargs): + raise OSError("locked") + + monkeypatch.setattr("aw_datastore.storages.sqlite_recover.sqlite3.connect", boom) + assert _counts(path) == ("n/a", "n/a") + + def test_sanitize_dump_sql_rewrites_rollback_and_drops_corruption_markers(): sql = ( "PRAGMA foreign_keys=OFF;\n" From 1756bf932b62b2c9006fd9b8dc8b382994ace1bb Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 16 Sep 2026 12:44:58 +0000 Subject: [PATCH 10/11] fix(datastore): tolerate partial failures in last-resort python recovery - Skip uncreatable tables in _recover_with_python instead of discarding all previously copied rows when one DDL fails. - Catch an unreadable bucketmodel during missing-bucket reconstruction so a recovered events file is not thrown away. Regression tests for both. Git-Session-Id: 37f48743-f7f9-53f3-b679-f8e58bb505d0 --- aw_datastore/storages/sqlite_recover.py | 36 ++++++++++++++++++------- tests/test_sqlite_recover.py | 36 +++++++++++++++++++++++++ 2 files changed, 62 insertions(+), 10 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index 83559ce..7861fcb 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -283,13 +283,20 @@ def _recover_with_python(src: str, dest: str) -> bool: if not tables: _truncate(dest) return False + created: list[str] = [] for _name, sql in tables: - dest_con.execute(sql) + try: + dest_con.execute(sql) + created.append(_name) + except sqlite3.Error as exc: + logger.info( + "Python recover: skipping uncreatable table %s: %s", _name, exc + ) dest_con.commit() src_con.close() src_closed = True copied = 0 - for name, _sql in tables: + for name in created: copied += _copy_table_rows(src, dest_con, name) for (sql,) in indexes: try: @@ -412,14 +419,23 @@ def _reconstruct_missing_buckets(path: str) -> int: "created an empty bucket catalog before reconstructing keys", path, ) - missing = con.execute( - """ - SELECT DISTINCT e.bucket_id - FROM eventmodel e - LEFT JOIN bucketmodel b ON b."key" = e.bucket_id - WHERE b."key" IS NULL AND e.bucket_id IS NOT NULL - """ - ).fetchall() + try: + missing = con.execute( + """ + SELECT DISTINCT e.bucket_id + FROM eventmodel e + LEFT JOIN bucketmodel b ON b."key" = e.bucket_id + WHERE b."key" IS NULL AND e.bucket_id IS NOT NULL + """ + ).fetchall() + except sqlite3.Error as exc: + logger.warning( + "Could not reconstruct missing buckets in %s " + "(unreadable bucketmodel/eventmodel); keeping recovered events: %s", + path, + exc, + ) + return 0 now = datetime.now(timezone.utc).isoformat() n = 0 for (key,) in missing: diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index 6e043b5..10fd094 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -322,3 +322,39 @@ def test_recovery_preserves_restrictive_mode(tmp_path): assert sidecar assert is_sqlite_healthy(path) assert stat.S_IMODE(os.stat(path).st_mode) == 0o600 + + +def test_recover_with_python_skips_uncreatable_tables(tmp_path): + """A table whose DDL fails on the destination must not abort the recovery.""" + from aw_datastore.storages import sqlite_recover + + src = str(tmp_path / "src.db") + dest = str(tmp_path / "dest.db") + con = sqlite3.connect(src) + con.create_collation("weird", lambda a, b: (a > b) - (a < b)) + con.execute("CREATE TABLE good (id INTEGER PRIMARY KEY, v TEXT)") + con.execute("INSERT INTO good VALUES (1, 'a')") + # DDL text persists into sqlite_master with COLLATE weird, which no + # destination connection can create without the collation registered. + con.execute("CREATE TABLE bad (v TEXT COLLATE weird)") + con.commit() + con.close() + assert sqlite_recover._recover_with_python(src, dest) + dcon = sqlite3.connect(dest) + assert dcon.execute("SELECT count(*) FROM good").fetchone()[0] == 1 + dcon.close() + + +def test_reconstruct_missing_buckets_tolerates_unreadable_bucketmodel(tmp_path): + """An unreadable bucketmodel must not discard the recovered events file.""" + from aw_datastore.storages import sqlite_recover + + db = str(tmp_path / "rec.db") + con = sqlite3.connect(db) + con.execute("CREATE TABLE eventmodel (bucket_id INTEGER)") + # bucketmodel exists but lacks the "key" column the LEFT JOIN needs. + con.execute("CREATE TABLE bucketmodel (id TEXT)") + con.execute("INSERT INTO eventmodel VALUES (1)") + con.commit() + con.close() + assert sqlite_recover._reconstruct_missing_buckets(db) == 0 From 74a59ebae7b923ed8bf7eaf8af4b2eef39e53296 Mon Sep 17 00:00:00 2001 From: Bob Date: Wed, 16 Sep 2026 13:59:31 +0000 Subject: [PATCH 11/11] fix(datastore): bound recover timeout reaps; keep dump data lines sqlite3 .dump comments of the form `/****** CORRUPTION ERROR *******/` were matched as a substring, so an INSERT whose payload contained those words was dropped. Drop comment lines only. After a .recover timeout, communicate() had no timeout and could hang PeeweeStorage.__init__. Bound the post-kill reaps. Git-Session-Id: 9fd15790-8aac-53c3-82f6-d6dd692d66c2 --- aw_datastore/storages/sqlite_recover.py | 18 ++++++++--- tests/test_sqlite_recover.py | 43 ++++++++++++++++++++++++- 2 files changed, 56 insertions(+), 5 deletions(-) diff --git a/aw_datastore/storages/sqlite_recover.py b/aw_datastore/storages/sqlite_recover.py index 7861fcb..65fabe5 100644 --- a/aw_datastore/storages/sqlite_recover.py +++ b/aw_datastore/storages/sqlite_recover.py @@ -119,7 +119,11 @@ def sanitize_dump_sql(sql: str) -> str: """Turn a ``.dump`` of a corrupt DB into SQL that can be loaded.""" lines: list[str] = [] for line in sql.splitlines(): - if "CORRUPTION ERROR" in line: + stripped = line.lstrip() + # sqlite3 .dump inserts `/****** CORRUPTION ERROR *******/` comments + # on bad pages. Drop those comments only — a payload containing the + # same words (window title, JSON) must still load. + if stripped.startswith("/*") and "CORRUPTION ERROR" in line: continue if line.startswith("ROLLBACK;"): lines.append("COMMIT;") @@ -190,8 +194,13 @@ def _recover_with_dbpage(sqlite_bin: str, src: str, dest: str) -> bool: except subprocess.TimeoutExpired: dump.kill() load.kill() - dump.communicate() - load.communicate() + # Bounded reaps only: unbounded communicate() here can hang + # PeeweeStorage.__init__ if a child ignores SIGPIPE with a full pipe. + for proc in (dump, load): + try: + proc.communicate(timeout=10) + except (subprocess.TimeoutExpired, OSError): + pass logger.warning("sqlite3 .recover timed out") return False if dump.returncode != 0: @@ -605,7 +614,8 @@ def _manual_instructions(path: str) -> str: "SQLite database is malformed. Preserve the file and recover with:\n" f" sqlite3 {path} '.recover' | sqlite3 recovered.db\n" "If .recover fails with 'no such table: sqlite_dbpage', use:\n" - f" sqlite3 {path} '.bail off' '.dump' | sed '/CORRUPTION ERROR/d; s/^ROLLBACK;/COMMIT;/'" + f" sqlite3 {path} '.bail off' '.dump'" + r" | sed '/^\/\*.*CORRUPTION ERROR/d; s/^ROLLBACK;/COMMIT;/'" " | sqlite3 recovered.db\n" "Then replace the original database with recovered.db." ) diff --git a/tests/test_sqlite_recover.py b/tests/test_sqlite_recover.py index 10fd094..59a5653 100644 --- a/tests/test_sqlite_recover.py +++ b/tests/test_sqlite_recover.py @@ -2,6 +2,7 @@ import os import sqlite3 import stat +import subprocess from datetime import datetime, timedelta, timezone import pytest @@ -117,13 +118,15 @@ def test_sanitize_dump_sql_rewrites_rollback_and_drops_corruption_markers(): "CREATE TABLE t(a);\n" "/****** CORRUPTION ERROR *******/\n" "INSERT INTO t VALUES(1);\n" + "INSERT INTO t VALUES('title with CORRUPTION ERROR in data');\n" "ROLLBACK; -- due to errors\n" ) out = sanitize_dump_sql(sql) - assert "CORRUPTION ERROR" not in out + assert "/****** CORRUPTION ERROR *******/" not in out assert "ROLLBACK;" not in out assert "COMMIT;" in out assert "INSERT INTO t VALUES(1);" in out + assert "INSERT INTO t VALUES('title with CORRUPTION ERROR in data');" in out def test_is_sqlite_healthy_missing_and_ok(tmp_path): @@ -358,3 +361,41 @@ def test_reconstruct_missing_buckets_tolerates_unreadable_bucketmodel(tmp_path): con.commit() con.close() assert sqlite_recover._reconstruct_missing_buckets(db) == 0 + + +class _TimeoutThenReap: + """First communicate() times out; later ones must be given a timeout.""" + + def __init__(self): + self.stdout = type("S", (), {"close": lambda self: None})() + self.returncode = -9 + self.calls = 0 + + def communicate(self, timeout=None): + self.calls += 1 + if self.calls == 1: + raise subprocess.TimeoutExpired(cmd="sqlite3", timeout=timeout) + if timeout is None: + raise AssertionError("post-kill communicate() must be bounded") + return ("", "") + + def kill(self): + pass + + +def test_recover_with_dbpage_timeout_does_not_hang(monkeypatch): + """A .recover timeout must not call unbounded communicate() after kill().""" + from aw_datastore.storages import sqlite_recover + + dump = _TimeoutThenReap() + load = _TimeoutThenReap() + n = {"i": 0} + + def fake_popen(*_args, **_kwargs): + n["i"] += 1 + return dump if n["i"] == 1 else load + + monkeypatch.setattr(sqlite_recover.subprocess, "Popen", fake_popen) + assert sqlite_recover._recover_with_dbpage("sqlite3", "src", "dest") is False + assert dump.calls >= 1 + assert load.calls >= 1