From a49da0b9a930da6367f777c73247e3edffc000e6 Mon Sep 17 00:00:00 2001 From: Bob Date: Tue, 1 Sep 2026 19:07:08 +0000 Subject: [PATCH 1/8] 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 2/8] 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 3/8] 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 4/8] 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 5/8] 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 6/8] 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 7/8] 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 8/8] 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)