Research-Stack/5-Applications/scripts/manifold_cache.py

463 lines
20 KiB
Python

#!/usr/bin/env python3
"""
MANIFOLD CACHE SYSTEM
Disk-backed persistent cache for incremental English manifold building.
Prevents memory blow-up and enables resume across runs.
Features:
- SQLite-backed sentence→fingerprint cache (dedup)
- SQLite-backed form→count accumulator (persistent Counter)
- Streaming checkpoint/resume
- Memory-capped with LRU spill to disk
- Batch insert for speed (10K+ sentences/sec)
- 1:1 restorable from cache alone
Schema:
sentences (hash PRIMARY KEY, text, fingerprint, source_file)
forms (fingerprint PRIMARY KEY, count, examples_json)
meta (key PRIMARY KEY, value)
"""
import os
import re
import sys
import json
import math
import hashlib
import sqlite3
import tempfile
from pathlib import Path
from collections import Counter, defaultdict
from datetime import datetime
from typing import Iterator, Optional, Dict, List, Tuple
BASE = Path("/home/allaun/Documents/Research Stack")
DEFAULT_CACHE_DIR = BASE / "3-Mathematical-Models/redpajama_english_manifold/cache"
# ── POS Tagger (reused) ───────────────────────────────────────────────────────
CLOSED = {'the','a','an','and','or','but','in','on','at','to','for','of','with','by','from','as',
'is','was','are','were','be','been','being','have','has','had','do','does','did','will',
'would','could','should','may','might','can','shall','this','that','these','those','it',
'its','they','them','their','he','she','his','her','him','we','us','our','you','your',
'my','mine','i','me','who','which','what','when','where','why','how','all','each',
'every','both','either','neither','some','any','no','none','more','most','many','much',
'few','little','other','another','such','only','own','same','so','than','too','very',
'just','now','then','here','there','up','down','out','off','over','under','again',
'further','once','not','also','always','never','often','sometimes','usually','still'}
PREP = {'in','on','at','by','for','with','about','against','between','into','through','during',
'before','after','above','below','to','from','up','down','of','off','over','under',
'again','further','then','once','around','behind','beyond','despite','except','inside',
'near','past','since','toward','upon','within','without','across','along','among',
'beside','besides','concerning','considering','following','including','like','minus',
'plus','regarding','round','save','till','until','via','worth'}
CONJ = {'and','or','but','nor','yet','so','for','although','because','before','if','since',
'though','unless','until','when','while','whereas','whether','either','neither',
'both','not','only','than','rather','however','moreover','furthermore','nevertheless',
'otherwise','therefore','thus','hence','consequently','meanwhile'}
AUX = {'be','am','is','are','was','were','being','been','have','has','had','do','does','did',
'will','would','shall','should','may','might','can','could','must','ought','need',
'dare','used','get','gets','got','getting','become','becomes','became','seem','seems',
'seemed','appear','appears','appeared'}
PRON = {'i','me','my','mine','myself','you','your','yours','yourself','he','him','his',
'himself','she','her','hers','herself','it','its','itself','we','us','our','ours',
'ourselves','they','them','their','theirs','themselves','this','that','these','those',
'who','whom','whose','which','what','whatever','whoever','whomever','anyone','someone',
'everyone','nobody','nothing','something','anything','everything'}
DET = {'the','a','an','this','that','these','those','my','your','his','her','its','our',
'their','some','any','no','each','every','either','neither','both','all','half',
'enough','several','many','much','few','little','other','another','such','what',
'which','whose','one','two','three','first','last','next','various','certain'}
def _tag(word: str) -> str:
w = word.lower().strip("'\"")
if w in DET: return "DET"
if w in PRON: return "PRON"
if w in PREP: return "PREP"
if w in CONJ: return "CONJ"
if w in AUX: return "AUX"
if w in CLOSED: return "FUNC"
if w.endswith("ing"): return "VBG"
if w.endswith("ed"): return "VBN"
if w.endswith(("ly","ily","ally")): return "ADV"
if w.endswith(("tion","sion","ment","ness","ity","ance","ence","hood","ship")): return "NOUN"
if w.endswith(("able","ible","ful","ous","ive","less","ish","al")): return "ADJ"
if w.endswith(("ize","ise","ify","ate")): return "VERB"
if len(w) <= 3: return "SHORT"
return "LEX"
def fingerprint(sentence: str) -> str:
words = re.findall(r"[a-zA-Z']+", sentence)
if len(words) < 3 or len(words) > 40:
return ""
tags = [_tag(w) for w in words]
collapsed = []
prev = None
for t in tags:
if t != prev:
collapsed.append(t)
prev = t
elif t == "LEX" and (len(collapsed) < 2 or collapsed[-2] != "LEX+"):
collapsed[-1] = "LEX+"
return " ".join(collapsed)
def _hash(text: str) -> str:
return hashlib.sha256(text.encode('utf-8')).hexdigest()[:32]
# ── Cache Engine ──────────────────────────────────────────────────────────────
class ManifoldCache:
"""SQLite-backed persistent cache for manifold construction."""
def __init__(self, cache_dir: Optional[Path] = None, max_ram_cache: int = 1_000_000):
self.cache_dir = Path(cache_dir) if cache_dir else DEFAULT_CACHE_DIR
self.cache_dir.mkdir(parents=True, exist_ok=True)
self.db_path = self.cache_dir / "manifold_cache.db"
self.max_ram_cache = max_ram_cache
# RAM buffers (spill to SQLite when full)
self._sentence_buffer: Dict[str, Tuple[str, str]] = {} # hash -> (sentence, fingerprint)
self._form_buffer: Dict[str, int] = Counter() # fingerprint -> count
self._examples_buffer: Dict[str, List[str]] = defaultdict(list)
self.sentences_cached = 0
self.sentences_hit = 0
self.forms_cached = 0
self._init_db()
self._load_resume_state()
def _init_db(self):
"""Initialize SQLite schema."""
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
conn.execute("""
CREATE TABLE IF NOT EXISTS sentences (
hash TEXT PRIMARY KEY,
text TEXT NOT NULL,
fingerprint TEXT NOT NULL,
indexed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS forms (
fingerprint TEXT PRIMARY KEY,
count INTEGER NOT NULL DEFAULT 0,
examples TEXT DEFAULT '[]'
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)
""")
conn.execute("""
CREATE INDEX IF NOT EXISTS idx_sent_fp ON sentences(fingerprint)
""")
conn.commit()
def _load_resume_state(self):
"""Load previous run counts."""
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
cursor = conn.execute("SELECT key, value FROM meta")
self.meta = {k: v for k, v in cursor.fetchall()}
# Load form counts
cursor = conn.execute("SELECT fingerprint, count, examples FROM forms")
for fp, cnt, ex_json in cursor.fetchall():
self._form_buffer[fp] = cnt
self._examples_buffer[fp] = json.loads(ex_json)
self.sentences_cached = int(self.meta.get('sentences_cached', '0'))
print(f" Cache resumed: {self.sentences_cached:,} sentences, {len(self._form_buffer):,} forms")
def get_fingerprint(self, sentence: str) -> Optional[str]:
"""Check cache for sentence fingerprint. Returns None if not found."""
h = _hash(sentence)
# Check RAM
if h in self._sentence_buffer:
self.sentences_hit += 1
return self._sentence_buffer[h][1]
# Check SQLite
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
cursor = conn.execute("SELECT fingerprint FROM sentences WHERE hash=?", (h,))
row = cursor.fetchone()
if row:
self.sentences_hit += 1
return row[0]
return None
def add_sentence(self, sentence: str, fingerprint: str):
"""Add sentence→fingerprint to cache."""
h = _hash(sentence)
self._sentence_buffer[h] = (sentence, fingerprint)
if len(self._sentence_buffer) >= self.max_ram_cache:
self._flush_sentences()
def increment_form(self, fingerprint: str, example: str):
"""Increment form count and store example."""
self._form_buffer[fingerprint] += 1
if len(self._examples_buffer[fingerprint]) < 3:
self._examples_buffer[fingerprint].append(example[:200])
if len(self._form_buffer) >= self.max_ram_cache:
self._flush_forms()
def _flush_sentences(self):
"""Flush sentence buffer to SQLite."""
if not self._sentence_buffer:
return
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
data = [(h, text, fp) for h, (text, fp) in self._sentence_buffer.items()]
conn.executemany(
"INSERT OR IGNORE INTO sentences (hash, text, fingerprint) VALUES (?, ?, ?)",
data
)
conn.commit()
self.sentences_cached += len(self._sentence_buffer)
self._sentence_buffer.clear()
print(f" Flushed {len(data):,} sentences to cache")
def _flush_forms(self):
"""Flush form counts to SQLite."""
if not self._form_buffer:
return
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
for fp, cnt in self._form_buffer.items():
ex_json = json.dumps(self._examples_buffer.get(fp, [])[:3])
conn.execute(
"""INSERT INTO forms (fingerprint, count, examples)
VALUES (?, ?, ?)
ON CONFLICT(fingerprint) DO UPDATE SET
count = count + excluded.count,
examples = excluded.examples""",
(fp, cnt, ex_json)
)
conn.commit()
self.forms_cached += len(self._form_buffer)
self._form_buffer.clear()
print(f" Flushed {self.forms_cached:,} forms to cache")
def save_meta(self, key: str, value: str):
"""Save metadata key-value."""
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)",
(key, value)
)
conn.commit()
def get_stats(self) -> dict:
"""Return cache statistics."""
return {
"db_path": str(self.db_path),
"db_size_mb": round(self.db_path.stat().st_size / (1024*1024), 2),
"sentences_cached": self.sentences_cached,
"sentences_hit": self.sentences_hit,
"forms_cached": self.forms_cached,
"ram_sentence_buffer": len(self._sentence_buffer),
"ram_form_buffer": len(self._form_buffer),
}
def load_all_forms(self) -> Counter:
"""Load all form counts from database + RAM."""
self._flush_forms() # ensure DB is up to date
forms = Counter()
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
cursor = conn.execute("SELECT fingerprint, count FROM forms")
for fp, cnt in cursor.fetchall():
forms[fp] += cnt
# Add RAM buffer
forms.update(self._form_buffer)
return forms
def load_examples(self) -> Dict[str, List[str]]:
"""Load all examples from database + RAM."""
self._flush_forms()
examples = defaultdict(list)
with sqlite3.connect(str(self.db_path), timeout=30) as conn:
cursor = conn.execute("SELECT fingerprint, examples FROM forms")
for fp, ex_json in cursor.fetchall():
examples[fp] = json.loads(ex_json)
# Merge RAM buffer
for fp, exs in self._examples_buffer.items():
examples[fp].extend(exs)
return dict(examples)
def close(self):
"""Flush all buffers and close."""
self._flush_sentences()
self._flush_forms()
self.save_meta('sentences_cached', str(self.sentences_cached))
self.save_meta('last_run', datetime.now().strftime("%Y%m%d_%H%M%S"))
print(f" Cache closed: {self.db_path}")
print(f" DB size: {self.db_path.stat().st_size / (1024*1024):.1f} MB")
# ── Cached Manifold Builder ──────────────────────────────────────────────────
class CachedManifoldBuilder:
"""Build English manifold with full disk-backed caching."""
def __init__(self, cache_dir: Optional[Path] = None, max_ram: int = 1_000_000):
self.cache = ManifoldCache(cache_dir, max_ram)
self.sentences_processed = 0
self.bytes_processed = 0
def process_sentence(self, sentence: str):
"""Process a single sentence with cache-first logic."""
sentence = sentence.strip()
if not (10 < len(sentence) < 300 and sentence[0].isupper()):
return
# Cache lookup
fp = self.cache.get_fingerprint(sentence)
if fp is None:
fp = fingerprint(sentence)
if fp and len(fp.split()) >= 3:
self.cache.add_sentence(sentence, fp)
if fp:
self.cache.increment_form(fp, sentence)
self.sentences_processed += 1
def process_text(self, text: str):
"""Process all sentences in a text block."""
self.bytes_processed += len(text.encode('utf-8'))
for sent in re.split(r'(?<=[.!?])\s+', text):
self.process_sentence(sent)
def build_report(self, top_n: int = 1000) -> dict:
"""Generate manifold report from cache."""
self.cache.close()
forms = self.cache.load_all_forms()
examples = self.cache.load_examples()
# Taxonomy
cats = Counter()
for fp, count in forms.items():
tags = fp.split()
if "DET" in tags and "NOUN" in tags and "VERB" in tags:
cats["SVO" if tags.index("VERB") > tags.index("NOUN") else "VSO"] += count
elif "DET" in tags and "NOUN" in tags and "PREP" in tags:
cats["NP_PP"] += count
elif "AUX" in tags and "VERB" in tags:
cats["AUX_V"] += count
elif "CONJ" in tags:
cats["COMPOUND"] += count
elif "PRON" in tags and "VERB" in tags:
cats["PRON_V"] += count
elif tags.count("PREP") >= 2:
cats["PP_CHAIN"] += count
elif "LEX+" in tags:
cats["DENSE_NP"] += count
else:
cats["OTHER"] += count
# Shannon entropy
total = sum(forms.values())
entropy = 0.0
for count in forms.values():
p = count / total
if p > 0:
entropy -= p * math.log2(p)
return {
"timestamp": datetime.now().strftime("%Y%m%d_%H%M%S"),
"sentences_processed": self.sentences_processed,
"bytes_processed": self.bytes_processed,
"unique_forms": len(forms),
"shannon_entropy_bits": round(entropy, 4),
"taxonomy": dict(cats),
"top_forms": [{"fingerprint": fp, "count": c, "example": examples.get(fp, [""])[0][:120]} for fp, c in forms.most_common(top_n)],
"cache_stats": self.cache.get_stats(),
}
def save_report(self, path: Path):
"""Save report to JSON."""
report = self.build_report()
with open(path, "w") as f:
json.dump(report, f, indent=2)
print(f" Report saved: {path}")
return report
# ── CLI ──────────────────────────────────────────────────────────────────────
def main():
import argparse
parser = argparse.ArgumentParser(description="Build cached English manifold")
parser.add_argument("--input", type=Path, nargs="+", help="Input file(s)")
parser.add_argument("--cache-dir", type=Path, default=DEFAULT_CACHE_DIR, help="Cache directory")
parser.add_argument("--max-ram", type=int, default=1_000_000, help="Max RAM buffer size")
parser.add_argument("--limit", type=int, default=1_000_000, help="Max sentences")
args = parser.parse_args()
print("=" * 70)
print(" CACHED ENGLISH MANIFOLD BUILDER")
print(" Persistent SQLite cache — resume anytime, no memory blow-up")
print("=" * 70)
builder = CachedManifoldBuilder(cache_dir=args.cache_dir, max_ram=args.max_ram)
# Default to enwik9 if no input
if not args.input:
args.input = [BASE / "shared-data/data/hutter_archive/enwik9_purified.bin"]
for path in args.input:
print(f"\nProcessing: {path}")
raw = path.read_bytes()
text_blocks = re.findall(rb'<text[^>]*>(.*?)</text>', raw, re.DOTALL)
for block in text_blocks:
text = block.decode('utf-8', errors='ignore')
text = re.sub(r'\{\{.*?\}\}', ' ', text, flags=re.DOTALL)
text = re.sub(r'\[\[.*?\|', ' ', text)
text = re.sub(r'\[\[|\]\]', ' ', text)
text = re.sub(r"'{2,}", ' ', text)
text = re.sub(r'<.*?>', ' ', text, flags=re.DOTALL)
text = re.sub(r'&\w+;', ' ', text)
text = re.sub(r'https?://\S+', ' ', text)
text = re.sub(r'[#*|=\{\}\[\]\|]', ' ', text)
builder.process_text(text)
if builder.sentences_processed >= args.limit:
print(f" Reached limit: {args.limit:,}")
break
if builder.sentences_processed >= args.limit:
break
# Save report
print(f"\n{'='*70}")
print(" BUILDING FINAL REPORT")
print(f"{'='*70}")
out_path = args.cache_dir / f"cached_manifold_report_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json"
report = builder.save_report(out_path)
print(f"\n{'='*70}")
print(" CACHED MANIFOLD COMPLETE")
print(f"{'='*70}")
print(f" Sentences: {report['sentences_processed']:,}")
print(f" Unique forms: {report['unique_forms']:,}")
print(f" Entropy: {report['shannon_entropy_bits']:.2f} bits/form")
print(f" Cache DB: {report['cache_stats']['db_size_mb']:.1f} MB")
print(f" Cache hits: {report['cache_stats']['sentences_hit']:,}")
print(f"{'='*70}")
if __name__ == "__main__":
main()