From 24375e66ffe066575be70c6393dd21ffd0082c71 Mon Sep 17 00:00:00 2001 From: Cameron Pfiffer Date: Fri, 30 Jan 2026 08:17:12 -0800 Subject: [PATCH] Add SQLite infrastructure for local state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Phase 1 of infrastructure upgrade: - tools/db.py: SQLite database layer with helpers for: - Published messages (dedup tracking) - Concepts, social graph, likes, consent - tools/migrate_to_sqlite.py: Migration script from JSON files - Update hooks to use SQLite instead of text files - Add data/central.db to .gitignore (runtime state) Migration results: - 6313 → 4503 published messages (1810 duplicates removed!) - 24 concepts, 73 social nodes, 6 likes migrated Benefits: - Indexed lookups (O(log n) vs O(n) text scan) - SQL queries for social graph analysis - Automatic deduplication - No more scattered JSON files 🐙 Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta --- .gitignore | 1 + hooks/livestream.py | 21 +-- hooks/publish-response.py | 26 ++- tools/db.py | 340 +++++++++++++++++++++++++++++++++++++ tools/migrate_to_sqlite.py | 163 ++++++++++++++++++ 5 files changed, 521 insertions(+), 30 deletions(-) create mode 100644 tools/db.py create mode 100644 tools/migrate_to_sqlite.py diff --git a/.gitignore b/.gitignore index 88d3361..5d4b9de 100644 --- a/.gitignore +++ b/.gitignore @@ -23,3 +23,4 @@ recon/ # Runtime state (not source code) data/published_messages.txt data/chroma/ +data/central.db diff --git a/hooks/livestream.py b/hooks/livestream.py index 454c96e..806f138 100755 --- a/hooks/livestream.py +++ b/hooks/livestream.py @@ -19,6 +19,10 @@ from pathlib import Path import httpx from dotenv import load_dotenv +# Add parent to path for tools import +sys.path.insert(0, str(Path(__file__).parent.parent)) +from tools.db import is_message_published, mark_message_published, init_db + # Redaction patterns for secrets (backup safety) REDACT_PATTERNS = [ (r'[A-Za-z_]*API_KEY[=:]\s*\S+', '[REDACTED_KEY]'), @@ -136,22 +140,16 @@ def get_recent_messages(limit: int = 20) -> list: def publish_messages(): """Publish any new assistant/reasoning messages.""" - published_file = Path(__file__).parent.parent / "data" / "published_messages.txt" - - try: - published_ids = set(published_file.read_text().splitlines()) - except: - published_ids = set() + init_db() messages = get_recent_messages(limit=30) session = get_session() if not session: return - new_ids = [] for msg in messages: msg_id = msg.get("id", "") - if msg_id in published_ids: + if is_message_published(msg_id): continue content = msg.get("content", "") @@ -176,13 +174,8 @@ def publish_messages(): timeout=10 ) if resp.status_code == 200: - new_ids.append(msg_id) + mark_message_published(msg_id) print(f"Published to {collection}", file=sys.stderr) - - if new_ids: - with open(published_file, "a") as f: - for mid in new_ids: - f.write(mid + "\n") def main(): diff --git a/hooks/publish-response.py b/hooks/publish-response.py index bda788d..1ade898 100755 --- a/hooks/publish-response.py +++ b/hooks/publish-response.py @@ -16,6 +16,10 @@ from pathlib import Path import httpx from dotenv import load_dotenv +# Add parent to path for tools import +sys.path.insert(0, str(Path(__file__).parent.parent)) +from tools.db import is_message_published, mark_message_published, init_db + # Use script directory for relative paths SCRIPT_DIR = Path(__file__).parent.parent load_dotenv(SCRIPT_DIR / ".env") @@ -185,21 +189,15 @@ def main(): if event_type not in ("PreToolUse", "Stop"): sys.exit(0) - # Track which messages we've already published - published_file = SCRIPT_DIR / "data/published_messages.txt" - try: - with open(published_file) as f: - published_ids = set(f.read().splitlines()) - except: - published_ids = set() + # Ensure database exists + init_db() # Fetch recent assistant + reasoning messages from Letta API messages = get_recent_messages(limit=30) - new_ids = [] for msg in messages: msg_id = msg.get("id", "") - if msg_id in published_ids: + if is_message_published(msg_id): continue content = msg.get("content", "") @@ -214,13 +212,9 @@ def main(): post_to_collection(content, "network.comind.prompt", "network.comind.prompt") else: post_to_collection(content, "network.comind.response", "network.comind.response") - new_ids.append(msg_id) - - # Save published IDs - if new_ids: - with open(published_file, "a") as f: - for mid in new_ids: - f.write(mid + "\n") + + # Mark as published in SQLite + mark_message_published(msg_id) sys.exit(0) diff --git a/tools/db.py b/tools/db.py new file mode 100644 index 0000000..14d5357 --- /dev/null +++ b/tools/db.py @@ -0,0 +1,340 @@ +""" +SQLite database layer for central's operational state. + +Consolidates: +- published_messages.txt -> published_messages table +- concepts.json -> concepts table +- social_graph.json -> social_nodes table +- liked.json -> likes table +- metrics.jsonl -> metrics table +- consent.json -> consent table +""" + +import sqlite3 +import json +from pathlib import Path +from contextlib import contextmanager +from typing import Optional, List, Dict, Any +from datetime import datetime + +# Database location +DB_PATH = Path(__file__).parent.parent / "data" / "central.db" + +# Schema definition +SCHEMA = """ +-- Dedup tracking (replaces published_messages.txt) +CREATE TABLE IF NOT EXISTS published_messages ( + id TEXT PRIMARY KEY, + timestamp TEXT DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_messages_timestamp ON published_messages(timestamp); + +-- Concepts (replaces concepts.json) +CREATE TABLE IF NOT EXISTS concepts ( + slug TEXT PRIMARY KEY, + confidence INTEGER, + tags TEXT, + summary TEXT, + updated TEXT +); + +-- Social graph nodes (replaces social_graph.json) +CREATE TABLE IF NOT EXISTS social_nodes ( + handle TEXT PRIMARY KEY, + did TEXT, + display_name TEXT, + first_seen TEXT, + relationship TEXT, + interactions INTEGER DEFAULT 0 +); + +-- Likes (replaces liked.json and liked_posts.txt) +CREATE TABLE IF NOT EXISTS likes ( + uri TEXT PRIMARY KEY, + timestamp TEXT DEFAULT CURRENT_TIMESTAMP +); + +-- Metrics (replaces metrics.jsonl) +CREATE TABLE IF NOT EXISTS metrics ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + timestamp TEXT DEFAULT CURRENT_TIMESTAMP, + data TEXT +); + +-- Consent (replaces consent.json) +CREATE TABLE IF NOT EXISTS consent ( + handle TEXT PRIMARY KEY, + opted_in INTEGER DEFAULT 0, + timestamp TEXT +); +""" + + +def init_db(): + """Initialize database with schema.""" + DB_PATH.parent.mkdir(parents=True, exist_ok=True) + with get_connection() as conn: + conn.executescript(SCHEMA) + print(f"Database initialized at {DB_PATH}") + + +@contextmanager +def get_connection(): + """Get a database connection with context management.""" + conn = sqlite3.connect(DB_PATH) + conn.row_factory = sqlite3.Row + try: + yield conn + conn.commit() + finally: + conn.close() + + +# --- Published Messages (dedup) --- + +def is_message_published(message_id: str) -> bool: + """Check if a message has been published.""" + with get_connection() as conn: + result = conn.execute( + "SELECT 1 FROM published_messages WHERE id = ?", + (message_id,) + ).fetchone() + return result is not None + + +def mark_message_published(message_id: str) -> None: + """Mark a message as published.""" + with get_connection() as conn: + conn.execute( + "INSERT OR IGNORE INTO published_messages (id) VALUES (?)", + (message_id,) + ) + + +def get_published_count() -> int: + """Get count of published messages.""" + with get_connection() as conn: + result = conn.execute("SELECT COUNT(*) FROM published_messages").fetchone() + return result[0] + + +# --- Concepts --- + +def get_concept(slug: str) -> Optional[Dict[str, Any]]: + """Get a concept by slug.""" + with get_connection() as conn: + row = conn.execute( + "SELECT * FROM concepts WHERE slug = ?", + (slug,) + ).fetchone() + if row: + return { + "slug": row["slug"], + "confidence": row["confidence"], + "tags": json.loads(row["tags"]) if row["tags"] else [], + "summary": row["summary"], + "updated": row["updated"] + } + return None + + +def upsert_concept(slug: str, confidence: int, tags: List[str], summary: str) -> None: + """Insert or update a concept.""" + with get_connection() as conn: + conn.execute(""" + INSERT INTO concepts (slug, confidence, tags, summary, updated) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT(slug) DO UPDATE SET + confidence = excluded.confidence, + tags = excluded.tags, + summary = excluded.summary, + updated = excluded.updated + """, (slug, confidence, json.dumps(tags), summary, datetime.utcnow().isoformat())) + + +def list_concepts() -> List[Dict[str, Any]]: + """List all concepts.""" + with get_connection() as conn: + rows = conn.execute("SELECT * FROM concepts ORDER BY slug").fetchall() + return [ + { + "slug": row["slug"], + "confidence": row["confidence"], + "tags": json.loads(row["tags"]) if row["tags"] else [], + "summary": row["summary"], + "updated": row["updated"] + } + for row in rows + ] + + +# --- Social Graph --- + +def get_social_node(handle: str) -> Optional[Dict[str, Any]]: + """Get a social graph node by handle.""" + with get_connection() as conn: + row = conn.execute( + "SELECT * FROM social_nodes WHERE handle = ?", + (handle,) + ).fetchone() + if row: + return { + "handle": row["handle"], + "did": row["did"], + "display_name": row["display_name"], + "first_seen": row["first_seen"], + "relationship": json.loads(row["relationship"]) if row["relationship"] else [], + "interactions": row["interactions"] + } + return None + + +def upsert_social_node( + handle: str, + did: str, + display_name: str, + relationship: List[str], + interactions: int = 0 +) -> None: + """Insert or update a social graph node.""" + with get_connection() as conn: + conn.execute(""" + INSERT INTO social_nodes (handle, did, display_name, first_seen, relationship, interactions) + VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT(handle) DO UPDATE SET + did = excluded.did, + display_name = excluded.display_name, + relationship = excluded.relationship, + interactions = excluded.interactions + """, (handle, did, display_name, datetime.utcnow().isoformat(), json.dumps(relationship), interactions)) + + +def increment_interactions(handle: str) -> None: + """Increment interaction count for a handle.""" + with get_connection() as conn: + conn.execute( + "UPDATE social_nodes SET interactions = interactions + 1 WHERE handle = ?", + (handle,) + ) + + +def list_social_nodes() -> List[Dict[str, Any]]: + """List all social graph nodes.""" + with get_connection() as conn: + rows = conn.execute("SELECT * FROM social_nodes ORDER BY interactions DESC").fetchall() + return [ + { + "handle": row["handle"], + "did": row["did"], + "display_name": row["display_name"], + "first_seen": row["first_seen"], + "relationship": json.loads(row["relationship"]) if row["relationship"] else [], + "interactions": row["interactions"] + } + for row in rows + ] + + +# --- Likes --- + +def is_liked(uri: str) -> bool: + """Check if a URI has been liked.""" + with get_connection() as conn: + result = conn.execute( + "SELECT 1 FROM likes WHERE uri = ?", + (uri,) + ).fetchone() + return result is not None + + +def add_like(uri: str) -> None: + """Record a like.""" + with get_connection() as conn: + conn.execute( + "INSERT OR IGNORE INTO likes (uri) VALUES (?)", + (uri,) + ) + + +def list_likes() -> List[str]: + """List all liked URIs.""" + with get_connection() as conn: + rows = conn.execute("SELECT uri FROM likes").fetchall() + return [row["uri"] for row in rows] + + +# --- Metrics --- + +def record_metric(data: Dict[str, Any]) -> None: + """Record a metric.""" + with get_connection() as conn: + conn.execute( + "INSERT INTO metrics (data) VALUES (?)", + (json.dumps(data),) + ) + + +def get_recent_metrics(limit: int = 100) -> List[Dict[str, Any]]: + """Get recent metrics.""" + with get_connection() as conn: + rows = conn.execute( + "SELECT * FROM metrics ORDER BY timestamp DESC LIMIT ?", + (limit,) + ).fetchall() + return [ + { + "id": row["id"], + "timestamp": row["timestamp"], + "data": json.loads(row["data"]) if row["data"] else {} + } + for row in rows + ] + + +# --- Consent --- + +def is_opted_in(handle: str) -> bool: + """Check if a handle has opted in.""" + with get_connection() as conn: + result = conn.execute( + "SELECT opted_in FROM consent WHERE handle = ?", + (handle,) + ).fetchone() + return result is not None and result["opted_in"] == 1 + + +def set_consent(handle: str, opted_in: bool) -> None: + """Set consent for a handle.""" + with get_connection() as conn: + conn.execute(""" + INSERT INTO consent (handle, opted_in, timestamp) + VALUES (?, ?, ?) + ON CONFLICT(handle) DO UPDATE SET + opted_in = excluded.opted_in, + timestamp = excluded.timestamp + """, (handle, 1 if opted_in else 0, datetime.utcnow().isoformat())) + + +# --- CLI --- + +if __name__ == "__main__": + import sys + + if len(sys.argv) < 2: + print("Usage: python -m tools.db ") + print("Commands: init, stats") + sys.exit(1) + + cmd = sys.argv[1] + + if cmd == "init": + init_db() + elif cmd == "stats": + init_db() # Ensure DB exists + print(f"Published messages: {get_published_count()}") + print(f"Concepts: {len(list_concepts())}") + print(f"Social nodes: {len(list_social_nodes())}") + print(f"Likes: {len(list_likes())}") + else: + print(f"Unknown command: {cmd}") + sys.exit(1) diff --git a/tools/migrate_to_sqlite.py b/tools/migrate_to_sqlite.py new file mode 100644 index 0000000..0e02a25 --- /dev/null +++ b/tools/migrate_to_sqlite.py @@ -0,0 +1,163 @@ +""" +Migration script to move JSON files to SQLite. + +Migrates: +- data/published_messages.txt -> published_messages table +- data/concepts.json -> concepts table +- data/social_graph.json -> social_nodes table +- data/liked.json -> likes table +- data/consent.json -> consent table +""" + +import json +from pathlib import Path +from tools.db import ( + init_db, + get_connection, + mark_message_published, + upsert_concept, + upsert_social_node, + add_like, + set_consent, + get_published_count, + list_concepts, + list_social_nodes, + list_likes, +) + +DATA_DIR = Path(__file__).parent.parent / "data" + + +def migrate_published_messages(): + """Migrate published_messages.txt to SQLite.""" + txt_path = DATA_DIR / "published_messages.txt" + if not txt_path.exists(): + print("No published_messages.txt found, skipping") + return 0 + + count = 0 + with open(txt_path) as f: + for line in f: + message_id = line.strip() + if message_id: + mark_message_published(message_id) + count += 1 + + print(f"Migrated {count} published messages") + return count + + +def migrate_concepts(): + """Migrate concepts.json to SQLite.""" + json_path = DATA_DIR / "concepts.json" + if not json_path.exists(): + print("No concepts.json found, skipping") + return 0 + + with open(json_path) as f: + data = json.load(f) + + count = 0 + for slug, concept in data.items(): + upsert_concept( + slug=slug, + confidence=concept.get("confidence", 0), + tags=concept.get("tags", []), + summary=concept.get("summary", "") + ) + count += 1 + + print(f"Migrated {count} concepts") + return count + + +def migrate_social_graph(): + """Migrate social_graph.json to SQLite.""" + json_path = DATA_DIR / "social_graph.json" + if not json_path.exists(): + print("No social_graph.json found, skipping") + return 0 + + with open(json_path) as f: + data = json.load(f) + + nodes = data.get("nodes", {}) + count = 0 + for handle, node in nodes.items(): + upsert_social_node( + handle=handle, + did=node.get("did", ""), + display_name=node.get("display_name", ""), + relationship=node.get("relationship", []), + interactions=node.get("interactions", 0) + ) + count += 1 + + print(f"Migrated {count} social nodes") + return count + + +def migrate_likes(): + """Migrate liked.json to SQLite.""" + json_path = DATA_DIR / "liked.json" + if not json_path.exists(): + print("No liked.json found, skipping") + return 0 + + with open(json_path) as f: + uris = json.load(f) + + count = 0 + for uri in uris: + add_like(uri) + count += 1 + + print(f"Migrated {count} likes") + return count + + +def migrate_consent(): + """Migrate consent.json to SQLite.""" + json_path = DATA_DIR / "consent.json" + if not json_path.exists(): + print("No consent.json found, skipping") + return 0 + + with open(json_path) as f: + data = json.load(f) + + count = 0 + for handle, info in data.items(): + opted_in = info.get("opted_in", False) if isinstance(info, dict) else bool(info) + set_consent(handle, opted_in) + count += 1 + + print(f"Migrated {count} consent records") + return count + + +def run_migration(): + """Run all migrations.""" + print("Initializing database...") + init_db() + + print("\n--- Starting migration ---\n") + + migrate_published_messages() + migrate_concepts() + migrate_social_graph() + migrate_likes() + migrate_consent() + + print("\n--- Migration complete ---\n") + + # Print stats + print(f"Final counts:") + print(f" Published messages: {get_published_count()}") + print(f" Concepts: {len(list_concepts())}") + print(f" Social nodes: {len(list_social_nodes())}") + print(f" Likes: {len(list_likes())}") + + +if __name__ == "__main__": + run_migration() -- 2.51.2