feat: SessionStore persistence to JSON file + startup auto recovery
This commit is contained in:
@@ -5,9 +5,12 @@ Singleton pattern with RLock guard. Stores:
|
||||
- Cluster results: result_id -> (labels, n_clusters, params, parent_dataset_id)
|
||||
- Feature results: result_id -> (features, parent_cluster_result_id)
|
||||
"""
|
||||
import json
|
||||
import os
|
||||
import threading
|
||||
import gc
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
import polars as pl
|
||||
@@ -43,11 +46,26 @@ def _generate_id(prefix: str = 'ds') -> str:
|
||||
return f"{prefix}_{uuid.uuid4().hex[:12]}"
|
||||
|
||||
|
||||
def _get_store_path() -> Path:
|
||||
"""Return the persistent JSON file path for dataset metadata.
|
||||
|
||||
Uses ``%APPDATA%/TianXuan/.session_store.json`` (Windows) or
|
||||
``~/.tianxuan_session_store.json`` (other platforms).
|
||||
"""
|
||||
appdata = os.environ.get('APPDATA', os.path.expanduser('~'))
|
||||
store_dir = Path(appdata) / 'TianXuan'
|
||||
store_dir.mkdir(parents=True, exist_ok=True)
|
||||
return store_dir / '.session_store.json'
|
||||
|
||||
|
||||
class SessionStore:
|
||||
"""Singleton in-memory store for datasets and analysis results.
|
||||
|
||||
Thread-safe via RLock. Datasets store Polars LazyFrames (query plans,
|
||||
not materialised data), so clone/shallow-copy is effectively free.
|
||||
|
||||
Dataset metadata is persisted to ``_get_store_path()`` so that after a
|
||||
server restart datasets whose CSV files still exist can be re-loaded.
|
||||
"""
|
||||
|
||||
_instance: Optional['SessionStore'] = None
|
||||
@@ -62,6 +80,65 @@ class SessionStore:
|
||||
cls._instance._datalock = threading.RLock()
|
||||
return cls._instance
|
||||
|
||||
# ── persistence ────────────────────────────────────────────────────
|
||||
|
||||
def _save_to_disk(self) -> None:
|
||||
"""Persist current dataset metadata to JSON file.
|
||||
|
||||
Only metadata (schema, csv_glob, row_count, …) is saved — never the
|
||||
LazyFrame itself, which is an in-memory query plan.
|
||||
"""
|
||||
with self._datalock:
|
||||
serializable = {}
|
||||
for ds_id, entry in self._stores.items():
|
||||
if entry.get('type') != 'dataset':
|
||||
continue
|
||||
meta = entry.get('metadata', {})
|
||||
serializable[ds_id] = {
|
||||
'schema': entry.get('schema'),
|
||||
'metadata': meta,
|
||||
'csv_glob': meta.get('csv_glob'),
|
||||
}
|
||||
path = _get_store_path()
|
||||
path.write_text(
|
||||
json.dumps(serializable, ensure_ascii=False, indent=2),
|
||||
encoding='utf-8',
|
||||
)
|
||||
|
||||
def restore_from_disk(self) -> int:
|
||||
"""Re-register datasets whose CSV files still exist on disk.
|
||||
|
||||
Uses ``csv_glob`` saved in metadata to re-load each dataset via
|
||||
:func:`analysis.data_loader.load_csv_directory`.
|
||||
|
||||
Returns the number of successfully restored datasets.
|
||||
"""
|
||||
path = _get_store_path()
|
||||
if not path.exists():
|
||||
return 0
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding='utf-8'))
|
||||
except (json.JSONDecodeError, OSError):
|
||||
return 0
|
||||
|
||||
# Lazy import to avoid circular dependency at module level
|
||||
from analysis.data_loader import load_csv_directory # fmt: skip
|
||||
|
||||
restored = 0
|
||||
for ds_id, info in data.items():
|
||||
csv_glob = info.get('csv_glob')
|
||||
if not csv_glob:
|
||||
continue
|
||||
try:
|
||||
lf, schema, row_count, file_count, memory_mb = load_csv_directory(csv_glob)
|
||||
meta = dict(info.get('metadata', {}))
|
||||
self.store_dataset(ds_id, lf, schema=schema, metadata=meta)
|
||||
restored += 1
|
||||
except Exception:
|
||||
# File(s) no longer exist or schema changed — skip gracefully
|
||||
continue
|
||||
return restored
|
||||
|
||||
# ── dataset operations ──────────────────────────────────────────────
|
||||
|
||||
def store_dataset(
|
||||
@@ -90,6 +167,7 @@ class SessionStore:
|
||||
'metadata': metadata or {},
|
||||
'parent_id': parent_id,
|
||||
}
|
||||
self._save_to_disk()
|
||||
|
||||
def get_dataset(self, dataset_id: str) -> Optional[dict]:
|
||||
"""Retrieve a stored dataset entry (or *None*)."""
|
||||
@@ -110,6 +188,7 @@ class SessionStore:
|
||||
del self._stores[dataset_id]
|
||||
if existed:
|
||||
gc.collect()
|
||||
self._save_to_disk()
|
||||
return existed
|
||||
|
||||
def clone_dataset(self, dataset_id: str) -> Optional[str]:
|
||||
|
||||
Reference in New Issue
Block a user