From 8312046f58b3ac0cfe5e9a93cab1cd9f68b22a41 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9D=8E=E4=BC=9F?= Date: Sun, 27 Sep 2026 17:27:39 +0800 Subject: [PATCH 1/2] =?UTF-8?q?fix(ai):=20=E9=98=BB=E6=AD=A2=E6=9C=AC?= =?UTF-8?q?=E5=9C=B0=E6=A3=80=E7=B4=A2=E8=BF=9B=E5=BA=A6=E4=BA=8B=E4=BB=B6?= =?UTF-8?q?=E6=97=A0=E9=99=90=E5=A2=9E=E9=95=BF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 本地检索每次进度变化都会向 events 表追加一份约 216 KB 的完整任务快照, 且没有 unique_key、没有保留策略,任务运行期间以约 3 行/秒持续写入, 使 output/local_search/ai.sqlite3 增长到 144 GB。 - 进度事件改用 replace=True 原地替换,每个任务只保留最新一行;自增 id 仍 递增,依赖 Last-Event-ID 重连的 SSE 依然能收到最新进度。 - 提醒类通知保留 INSERT OR IGNORE 去重语义,重复 key 不会重置 delivered, 避免已确认的提醒被重新投递。 - 进度事件剔除可从 records 重建的大字段(config/coverage/segments/ read_starts),完整快照仍以 records 表为准;前端改为合并增量字段, 不再整体替换任务对象。 - 新增启动维护:折叠旧版遗留的重复快照、按 24 小时 TTL 回收事件、空闲页 超过 64 MB 时执行 VACUUM,并在后台线程对 summary 与 search 两个 AI 库 执行,避免大型遗留库拖慢启动;初始化不再同步建 events 索引,杜绝阻塞启动。 - 新增 tests/test_ai_storage_retention.py 覆盖原地替换、非 replace 去重、 遗留快照折叠、TTL 与 VACUUM;tests/test_ai_message_pages.py 同步更新断言。 --- frontend/components/LocalSearchSettings.vue | 3 +- src/wechat_decrypt_tool/ai/lifecycle.py | 13 ++ src/wechat_decrypt_tool/ai/storage.py | 140 ++++++++++++++++-- .../local_search/downloads.py | 2 +- src/wechat_decrypt_tool/local_search/gpu.py | 2 +- .../local_search/service.py | 9 +- .../local_search/totals.py | 2 +- tests/test_ai_message_pages.py | 13 +- tests/test_ai_storage_retention.py | 85 +++++++++++ 9 files changed, 246 insertions(+), 23 deletions(-) create mode 100644 tests/test_ai_storage_retention.py diff --git a/frontend/components/LocalSearchSettings.vue b/frontend/components/LocalSearchSettings.vue index 07bc9144..f2dce5f8 100644 --- a/frontend/components/LocalSearchSettings.vue +++ b/frontend/components/LocalSearchSettings.vue @@ -378,7 +378,8 @@ function subscribe(){ if(event?.kind==='local_search_index' && event.account===current){ const incoming=event.body, existing=state.value.jobs?.find(j=>j.id===incoming?.id) if(existing && incoming.updated>=existing.updated){ - state.value.jobs=state.value.jobs.map(j=>j.id===incoming.id ? incoming : j) + // 进度事件只带高频字段;合并保留首次加载时拿到的 config/segments 等快照。 + state.value.jobs=state.value.jobs.map(j=>j.id===incoming.id ? {...j, ...incoming} : j) now.value=Date.now()/1000 if(incoming.status==='done')refresh() }else if(!existing){refresh()}else{api.diagnostic?.('response.stale',{task_id:incoming.id,component:'search'})} diff --git a/src/wechat_decrypt_tool/ai/lifecycle.py b/src/wechat_decrypt_tool/ai/lifecycle.py index cfb06d94..31c63e00 100644 --- a/src/wechat_decrypt_tool/ai/lifecycle.py +++ b/src/wechat_decrypt_tool/ai/lifecycle.py @@ -1,10 +1,21 @@ """AI 服务统一启停诊断,独立于其他后台服务的生命周期。""" import logging +import threading from importlib.metadata import PackageNotFoundError, version from .diagnostics import observed, event +def _maintain_store(store, name): + """在后台回收过期事件并压缩数据库,避免大型遗留库拖慢启动。""" + try: + deduplicated, removed, freed = store.maintain() + event('storage.maintenance.finished', component=name, events_deduplicated=deduplicated, + events_removed=removed, bytes_freed=freed) + except Exception as error: + event('storage.maintenance.failed', level=logging.WARNING, component=name, error=error) + + @observed('lifecycle.start') async def start_services(): for package in ('deepagents', 'langchain-openai', 'langchain-anthropic', 'langgraph', 'onnxruntime', 'sqlite-vec', 'tokenizers'): @@ -19,6 +30,8 @@ async def start_services(): get_ai_service().start() await get_agent_service().start() await get_local_search().start() + for name, store in (('summary', get_ai_service().store), ('search', get_local_search().store)): + threading.Thread(target=_maintain_store, args=(store, name), name=f'ai-store-maintenance-{name}', daemon=True).start() @observed('lifecycle.stop') diff --git a/src/wechat_decrypt_tool/ai/storage.py b/src/wechat_decrypt_tool/ai/storage.py index 2998d930..f0d55660 100644 --- a/src/wechat_decrypt_tool/ai/storage.py +++ b/src/wechat_decrypt_tool/ai/storage.py @@ -13,6 +13,25 @@ from ..app_paths import get_output_dir +# 事件默认保留窗口;超过该窗口且无需重放的事件会被回收。 +EVENT_RETENTION_SECONDS = 24 * 3600 +# 事件表空闲页超过该阈值才执行 VACUUM,避免频繁全库重写。 +COMPACT_MINIMUM_BYTES = 64 * 1024 * 1024 + +SCHEMA_SQL = """ + CREATE TABLE IF NOT EXISTS records ( + kind TEXT NOT NULL, id TEXT NOT NULL, account TEXT NOT NULL DEFAULT '', + body TEXT NOT NULL, updated REAL NOT NULL, PRIMARY KEY(kind,id)); + CREATE INDEX IF NOT EXISTS records_account ON records(kind,account,updated); + CREATE INDEX IF NOT EXISTS records_status ON records(kind,json_extract(body,'$.status'),updated); + CREATE INDEX IF NOT EXISTS records_task_usage ON records(kind,account,json_extract(body,'$.task_id')); + CREATE TABLE IF NOT EXISTS events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, account TEXT NOT NULL, + kind TEXT NOT NULL, body TEXT NOT NULL, unique_key TEXT UNIQUE, + delivered INTEGER NOT NULL DEFAULT 0, created REAL NOT NULL); +""" + + class AIStore: """短事务业务存储;与工作流检查点分开,避免模型请求持有数据库锁。""" @@ -27,18 +46,7 @@ def __init__(self, root: Path | None = None): self._event_condition = threading.Condition() self._event_revisions = {} with self.connection() as db: - db.executescript(""" - CREATE TABLE IF NOT EXISTS records ( - kind TEXT NOT NULL, id TEXT NOT NULL, account TEXT NOT NULL DEFAULT '', - body TEXT NOT NULL, updated REAL NOT NULL, PRIMARY KEY(kind,id)); - CREATE INDEX IF NOT EXISTS records_account ON records(kind,account,updated); - CREATE INDEX IF NOT EXISTS records_status ON records(kind,json_extract(body,'$.status'),updated); - CREATE INDEX IF NOT EXISTS records_task_usage ON records(kind,account,json_extract(body,'$.task_id')); - CREATE TABLE IF NOT EXISTS events ( - id INTEGER PRIMARY KEY AUTOINCREMENT, account TEXT NOT NULL, - kind TEXT NOT NULL, body TEXT NOT NULL, unique_key TEXT UNIQUE, - delivered INTEGER NOT NULL DEFAULT 0, created REAL NOT NULL); - """) + db.executescript(SCHEMA_SQL) @contextmanager def connection(self): @@ -108,13 +116,30 @@ def delete(self, kind, id): with self.connection() as db: db.execute("DELETE FROM records WHERE kind=? AND id=?", (kind, id)) - def event(self, account, kind, body, unique_key=None): - inserted = False + def event(self, account, kind, body, unique_key=None, replace=False): + """写入事件供 SSE 重放。 + + 默认行为保持不变:提供 `unique_key` 时按去重语义写入(同 key 已存在则忽略), + 用于提醒等只应投递一次的事件。`replace=True` 时改为用最新快照替换旧行, + 让高频进度事件每个逻辑任务只保留一行,同时因 INSERT OR REPLACE 会删除旧行、 + 新行仍获得递增的自增 id,断线重连的 EventSource 依然能收到最新状态。 + """ with self.connection() as db: if account in self.revoked_accounts: return - cursor = db.execute("INSERT OR IGNORE INTO events(account,kind,body,unique_key,created) VALUES(?,?,?,?,?)", - (account, kind, json.dumps(body, ensure_ascii=False), unique_key, time.time())) + payload = json.dumps(body, ensure_ascii=False) + if unique_key is None: + cursor = db.execute( + "INSERT INTO events(account,kind,body,created) VALUES(?,?,?,?)", + (account, kind, payload, time.time())) + elif replace: + cursor = db.execute( + "INSERT OR REPLACE INTO events(account,kind,body,unique_key,created) VALUES(?,?,?,?,?)", + (account, kind, payload, unique_key, time.time())) + else: + cursor = db.execute( + "INSERT OR IGNORE INTO events(account,kind,body,unique_key,created) VALUES(?,?,?,?,?)", + (account, kind, payload, unique_key, time.time())) inserted = cursor.rowcount > 0 if inserted: with self._event_condition: @@ -151,6 +176,89 @@ def acknowledge(self, id): with self.connection() as db: db.execute("UPDATE events SET delivered=1 WHERE id=?", (id,)) + @observed('storage.prune_duplicates') + def prune_duplicate_events(self, batch=2000): + """一次性折叠旧版追加式进度事件:每个逻辑任务只保留最新快照。 + + 新写入的进度事件已带 unique_key、本身只保留一行;这里主要清理升级前 + 历史遗留的、同一任务多次追加的整份快照,避免巨型库只能等 TTL 慢慢过期。 + """ + total = 0 + for kind, field in (('local_search_index', '$.id'), ('local_search_download', '$.id'), + ('local_search_total', '$.job_id')): + with self.connection() as db: + db.execute("CREATE TEMP TABLE IF NOT EXISTS keep_event_ids(id INTEGER PRIMARY KEY)") + db.execute("DELETE FROM keep_event_ids") + db.execute( + f"INSERT INTO keep_event_ids SELECT max(id) FROM events " + f"WHERE kind=? AND unique_key IS NULL GROUP BY json_extract(body,'{field}')", (kind,)) + while True: + removed = db.execute( + "DELETE FROM events WHERE id IN (SELECT id FROM events " + "WHERE kind=? AND unique_key IS NULL AND id NOT IN (SELECT id FROM keep_event_ids) LIMIT ?)", + (kind, batch)).rowcount + total += removed + db.commit() + if removed < batch: + break + db.execute("DELETE FROM keep_event_ids") + for kind in ('local_search_device', 'local_search_gpu'): + with self.connection() as db: + total += db.execute( + "DELETE FROM events WHERE kind=? AND unique_key IS NULL AND id < (SELECT max(id) FROM events WHERE kind=?)", + (kind, kind)).rowcount + if total: + diagnostic_event('storage.events.deduplicated', count=total) + return total + + @observed('storage.prune_events') + def prune_events(self, max_age=EVENT_RETENTION_SECONDS, batch=2000): + """按 TTL 批量回收事件:非通知事件直接过期;已投递通知同样回收,未投递通知保留。""" + cutoff = time.time() - max_age + total = 0 + while True: + with self.connection() as db: + removed = db.execute( + "DELETE FROM events WHERE id IN (SELECT id FROM events " + "WHERE created e['body']['processed'] for e in events) + and e['body'].get('read_count', 0) > e['body']['processed'] for e in emitted) + # 同一任务在事件表只保留最新一行,且不携带可重建的大字段。 + stored = [e for e in events if e['kind'] == 'local_search_index' and e['body'].get('id') == job['id']] + assert len(stored) == 1 + assert 'config' not in stored[0]['body'] and 'segments' not in stored[0]['body'] and 'coverage' not in stored[0]['body'] await service.stop() asyncio.run(run()) diff --git a/tests/test_ai_storage_retention.py b/tests/test_ai_storage_retention.py new file mode 100644 index 00000000..3590746b --- /dev/null +++ b/tests/test_ai_storage_retention.py @@ -0,0 +1,85 @@ +"""回归:进度事件必须原地替换、按 TTL 回收并可压缩,避免 ai.sqlite3 无界膨胀。""" +import json +import sys +import time +from pathlib import Path +from types import SimpleNamespace + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) +from wechat_decrypt_tool.ai.storage import AIStore +from wechat_decrypt_tool.local_search.service import LocalSearch + + +def test_local_search_update_emits_compact_deduplicated_event(tmp_path): + service = LocalSearch(tmp_path, engine=SimpleNamespace(status={}, gpu_failed=False)) + job = {'id': 'job1', 'account': 'a', 'config': {'usernames': ['chat']}, 'coverage': {'0': {}}, + 'segments': [{}], 'read_starts': {'chat': 0}, 'processed': 0} + for processed in range(10): + service.update(job, processed=processed) + events = [event for event in service.store.events() if event['kind'] == 'local_search_index'] + assert len(events) == 1 + body = events[0]['body'] + assert body['processed'] == 9 and body['id'] == 'job1' + for heavy in ('config', 'coverage', 'segments', 'read_starts'): + assert heavy not in body + + +def test_prune_duplicate_events_collapses_legacy_snapshots(tmp_path): + store = AIStore(tmp_path) + now = time.time() + with store.connection() as db: + for processed in range(5): + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','local_search_index',?,NULL,0,?)", + (json.dumps({'id': 'job-a', 'processed': processed}), now)) + for processed in (1, 2): + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','local_search_index',?,NULL,0,?)", + (json.dumps({'id': 'job-b', 'processed': processed}), now)) + assert store.prune_duplicate_events(batch=2) == 5 + kept = sorted(row['body']['processed'] for row in store.events()) + assert kept == [2, 4] + + +def test_unique_key_event_keeps_only_latest_snapshot(tmp_path): + store = AIStore(tmp_path) + for processed in range(50): + store.event('', 'local_search_index', {'id': 'job', 'processed': processed}, + unique_key='index_job:job', replace=True) + rows = store.events() + assert len(rows) == 1 + assert rows[0]['body']['processed'] == 49 + # 替换后仍产生新的自增 id,保证 SSE 客户端能收到更新。 + assert rows[0]['id'] > 0 + + +def test_unique_key_without_replace_stays_idempotent(tmp_path): + store = AIStore(tmp_path) + store.event('', 'notification', {'n': 1}, unique_key='summary:1') + store.event('', 'notification', {'n': 2}, unique_key='summary:1') + rows = store.events() + assert len(rows) == 1 and rows[0]['body']['n'] == 1 + + +def test_prune_events_keeps_undelivered_notifications(tmp_path): + store = AIStore(tmp_path) + old = time.time() - 10 * 24 * 3600 + with store.connection() as db: + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','local_search_index','{}','old-progress',0,?)", (old,)) + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','notification','{}',NULL,0,?)", (old,)) + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','notification','{}',NULL,1,?)", (old,)) + assert store.prune_events(max_age=24 * 3600) == 2 + remaining = store.events() + assert len(remaining) == 1 + assert remaining[0]['kind'] == 'notification' and remaining[0]['delivered'] == 0 + + +def test_compact_reclaims_free_pages_after_prune(tmp_path): + store = AIStore(tmp_path) + for index in range(2000): + store.event('', 'task', {'payload': 'x' * 2000, 'index': index}) + with store.connection() as db: + db.execute("DELETE FROM events") + before = store.path.stat().st_size + assert before > 0 + freed = store.compact(minimum_bytes=1) + assert freed > 0 + assert store.path.stat().st_size < before From a52aecaadcc42135007e9e823a7d8fe0a0b8fd2d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9D=8E=E4=BC=9F?= Date: Sun, 27 Sep 2026 17:27:57 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix(ai):=20=E8=87=AA=E5=8A=A8=E9=87=8D?= =?UTF-8?q?=E5=BB=BA=E5=B7=B2=E8=86=A8=E8=83=80=E7=9A=84=20AI=20=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E5=BA=93?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 升级前已膨胀的库只靠 TTL 回收会残留很久,直接删除重建又会丢失本地检索配置 (records 中的 active 索引指针),导致需要重新整理/重新向量化。 - 库超过 512 MB 时改为重建:把 records 与未投递提醒复制到新库, 丢弃可再生的 events;耗时与库体积无关。 - 先原子替换主库、成功后再清理旧 WAL/SHM;替换失败自动重试, 仍失败则保留原库并告警,应用照常运行,下次启动再试。 - 保留 sqlite_sequence,重建后事件 id 继续递增,SSE 不断档、前端无感。 - maintain() 在超大库上走重建分支,不再原地删除/VACUUM。 - tests/test_ai_storage_retention.py 增加超大库重建用例。 --- src/wechat_decrypt_tool/ai/storage.py | 96 ++++++++++++++++++++++++++- tests/test_ai_storage_retention.py | 23 +++++++ 2 files changed, 117 insertions(+), 2 deletions(-) diff --git a/src/wechat_decrypt_tool/ai/storage.py b/src/wechat_decrypt_tool/ai/storage.py index f0d55660..7c49fca8 100644 --- a/src/wechat_decrypt_tool/ai/storage.py +++ b/src/wechat_decrypt_tool/ai/storage.py @@ -3,6 +3,7 @@ import logging import json +import os import sqlite3 import threading import time @@ -17,6 +18,9 @@ EVENT_RETENTION_SECONDS = 24 * 3600 # 事件表空闲页超过该阈值才执行 VACUUM,避免频繁全库重写。 COMPACT_MINIMUM_BYTES = 64 * 1024 * 1024 +# 超过该体积的库不做原地删除/VACUUM(对遗留巨型库会放大 WAL),改为重建: +# 仅保留 records 与未投递提醒,丢弃可再生的 events。 +MAINTENANCE_MAX_DATABASE_BYTES = 512 * 1024 * 1024 SCHEMA_SQL = """ CREATE TABLE IF NOT EXISTS records ( @@ -252,13 +256,101 @@ def compact(self, minimum_bytes=COMPACT_MINIMUM_BYTES): diagnostic_event('storage.compacted', freed_bytes=free) return free - def maintain(self, max_age=EVENT_RETENTION_SECONDS, minimum_bytes=COMPACT_MINIMUM_BYTES): - """启动维护:折叠遗留重复事件、回收过期事件,再按需压缩数据库文件。""" + @observed('storage.repair') + def repair_oversized(self, max_database_bytes=MAINTENANCE_MAX_DATABASE_BYTES): + """重建过大的库:保留 records 与未投递提醒,丢弃可再生的 events。 + + 对遗留巨型库,原地 DELETE + VACUUM 会把 WAL 放大到库体积且长时间占锁; + 这里改为把少量存活数据复制到新库再原子替换,耗时与库体积无关。 + 返回重建前的库大小(字节),未触发或失败返回 0。 + """ + with self.lock: + try: + probe = sqlite3.connect(self.path, timeout=30) + try: + database_bytes = self._database_bytes(probe) + finally: + probe.close() + if database_bytes <= max_database_bytes: + return 0 + + temporary = self.path.with_name(self.path.name + '.repair') + for suffix in ('', '-wal', '-shm'): + try: + os.remove(str(temporary) + suffix) + except FileNotFoundError: + pass + + source = sqlite3.connect(self.path, timeout=30) + target = sqlite3.connect(temporary, timeout=30) + try: + source.execute('PRAGMA wal_checkpoint(TRUNCATE)') + target.executescript(SCHEMA_SQL) + target.executemany( + 'INSERT INTO records(kind,id,account,body,updated) VALUES(?,?,?,?,?)', + source.execute('SELECT kind,id,account,body,updated FROM records')) + # 未投递提醒不可再生,随 records 一起保留;其余 events 只是进度快照。 + target.executemany( + 'INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES(?,?,?,?,?,?)', + source.execute("SELECT account,kind,body,unique_key,delivered,created FROM events " + "WHERE kind='notification' AND delivered=0")) + sequence = source.execute("SELECT seq FROM sqlite_sequence WHERE name='events'").fetchone() + if sequence: + target.execute("INSERT INTO sqlite_sequence(name,seq) VALUES('events',?)", (sequence[0],)) + target.commit() + finally: + source.close() + target.close() + + # 先原子替换主库,再清理旧的 WAL/SHM:即使替换失败,原库仍完整。 + for attempt in range(4): + try: + os.replace(temporary, self.path) + break + except OSError as error: + if attempt == 3: + raise + diagnostic_event('storage.repair.retry', level=logging.WARNING, error=error) + time.sleep(0.3 * (attempt + 1)) + for suffix in ('-wal', '-shm'): + try: + os.remove(str(self.path) + suffix) + except FileNotFoundError: + pass + diagnostic_event('storage.repaired', bytes_before=database_bytes) + return database_bytes + except Exception as error: + for suffix in ('', '-wal', '-shm'): + try: + os.remove(str(self.path.with_name(self.path.name + '.repair')) + suffix) + except FileNotFoundError: + pass + diagnostic_event('storage.repair.failed', level=logging.WARNING, error=error) + return 0 + + def maintain(self, max_age=EVENT_RETENTION_SECONDS, minimum_bytes=COMPACT_MINIMUM_BYTES, + max_database_bytes=MAINTENANCE_MAX_DATABASE_BYTES): + """启动维护:折叠遗留重复事件、回收过期事件,再按需压缩数据库文件。 + + 超过 `max_database_bytes` 的遗留巨型库改为重建(保留 records 与未投递提醒), + 避免原地删除/VACUUM 长时间占锁并放大 WAL。 + """ + with self.connection() as db: + database_bytes = self._database_bytes(db) + if database_bytes > max_database_bytes: + repaired = self.repair_oversized(max_database_bytes) + return 0, 0, repaired deduplicated = self.prune_duplicate_events() removed = self.prune_events(max_age) freed = self.compact(minimum_bytes) return deduplicated, removed, freed + @staticmethod + def _database_bytes(db): + page_size = db.execute('PRAGMA page_size').fetchone()[0] + page_count = db.execute('PRAGMA page_count').fetchone()[0] + return page_size * page_count + @observed('storage.purge_account') def purge_account(self, account): with self.connection() as db: diff --git a/tests/test_ai_storage_retention.py b/tests/test_ai_storage_retention.py index 3590746b..f6296553 100644 --- a/tests/test_ai_storage_retention.py +++ b/tests/test_ai_storage_retention.py @@ -72,6 +72,29 @@ def test_prune_events_keeps_undelivered_notifications(tmp_path): assert remaining[0]['kind'] == 'notification' and remaining[0]['delivered'] == 0 +def test_maintain_repairs_oversized_database_preserving_records(tmp_path): + store = AIStore(tmp_path) + store.put('config', {'enabled': True}, id='acct', account='acct') + now = time.time() + with store.connection() as db: + for index in range(300): + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','local_search_index',?,NULL,0,?)", + (json.dumps({'id': 'job', 'index': index, 'pad': 'x' * 2000}), now)) + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','notification','{\"n\":1}','k1',0,?)", (now,)) + db.execute("INSERT INTO events(account,kind,body,unique_key,delivered,created) VALUES('','notification','{\"n\":2}','k2',1,?)", (now,)) + before = store._database_bytes(db) + # 用极小阈值模拟遗留巨型库:应重建而不是长时间原地删除。 + deduplicated, pruned, repaired = store.maintain(max_database_bytes=1) + assert repaired > 0 + with store.connection() as db: + assert store._database_bytes(db) < before + # records 保留;未投递提醒保留,可再生的进度事件与已投递提醒丢弃。 + assert store.get('config', 'acct')['enabled'] is True + events = store.events() + assert [event['kind'] for event in events] == ['notification'] + assert events[0]['body']['n'] == 1 + + def test_compact_reclaims_free_pages_after_prune(tmp_path): store = AIStore(tmp_path) for index in range(2000):