Files
mosaic/storage/migrations.py
T

192 lines
8.5 KiB
Python

"""Schema versioning and migration runner for Mosaic database."""
import sqlite3
import logging
logger = logging.getLogger(__name__)
MIGRATIONS = [
# Version 1: Initial schema
{
"version": 1,
"description": "Initial schema",
"sql": [
"""CREATE TABLE IF NOT EXISTS documents (
id TEXT PRIMARY KEY,
source TEXT NOT NULL DEFAULT '',
source_url TEXT NOT NULL DEFAULT '',
doc_type TEXT NOT NULL DEFAULT '',
title TEXT NOT NULL DEFAULT '',
date TEXT,
classification TEXT,
raw_path TEXT NOT NULL DEFAULT '',
text TEXT NOT NULL DEFAULT '',
metadata TEXT NOT NULL DEFAULT '{}',
content_hash TEXT NOT NULL DEFAULT '',
fetch_date TEXT NOT NULL DEFAULT '',
etag TEXT,
last_modified TEXT,
char_count INTEGER NOT NULL DEFAULT 0,
status TEXT NOT NULL DEFAULT 'PENDING',
parse_status TEXT NOT NULL DEFAULT 'PENDING',
last_extracted_at TEXT
)""",
"""CREATE INDEX IF NOT EXISTS idx_documents_source ON documents(source)""",
"""CREATE INDEX IF NOT EXISTS idx_documents_status ON documents(status)""",
"""CREATE INDEX IF NOT EXISTS idx_documents_parse_status ON documents(parse_status)""",
"""CREATE INDEX IF NOT EXISTS idx_documents_date ON documents(date)""",
"""CREATE INDEX IF NOT EXISTS idx_documents_source_url ON documents(source_url)""",
"""CREATE VIRTUAL TABLE IF NOT EXISTS documents_fts USING fts5(
title, text, source,
content=documents,
content_rowid=rowid
)""",
# Triggers to keep FTS in sync
"""CREATE TRIGGER IF NOT EXISTS documents_ai AFTER INSERT ON documents BEGIN
INSERT INTO documents_fts(rowid, title, text, source)
VALUES (new.rowid, new.title, new.text, new.source);
END""",
"""CREATE TRIGGER IF NOT EXISTS documents_ad AFTER DELETE ON documents BEGIN
INSERT INTO documents_fts(documents_fts, rowid, title, text, source)
VALUES ('delete', old.rowid, old.title, old.text, old.source);
END""",
"""CREATE TRIGGER IF NOT EXISTS documents_au AFTER UPDATE ON documents BEGIN
INSERT INTO documents_fts(documents_fts, rowid, title, text, source)
VALUES ('delete', old.rowid, old.title, old.text, old.source);
INSERT INTO documents_fts(rowid, title, text, source)
VALUES (new.rowid, new.title, new.text, new.source);
END""",
"""CREATE TABLE IF NOT EXISTS extracted_tools (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL DEFAULT '',
aliases TEXT NOT NULL DEFAULT '[]',
description TEXT NOT NULL DEFAULT '',
capability TEXT NOT NULL DEFAULT '',
target_platforms TEXT NOT NULL DEFAULT '[]',
target_software TEXT NOT NULL DEFAULT '[]',
cves TEXT NOT NULL DEFAULT '[]',
source_doc_ids TEXT NOT NULL DEFAULT '[]',
source_context TEXT NOT NULL DEFAULT '',
mitre_techniques TEXT NOT NULL DEFAULT '[]',
dedup_key TEXT NOT NULL DEFAULT ''
)""",
"""CREATE UNIQUE INDEX IF NOT EXISTS idx_tools_dedup ON extracted_tools(dedup_key)""",
"""CREATE TABLE IF NOT EXISTS ttps (
id INTEGER PRIMARY KEY AUTOINCREMENT,
technique TEXT NOT NULL DEFAULT '',
category TEXT NOT NULL DEFAULT '',
description TEXT NOT NULL DEFAULT '',
actor TEXT,
mitre_id TEXT,
source_doc_ids TEXT NOT NULL DEFAULT '[]',
source_context TEXT NOT NULL DEFAULT '',
dedup_key TEXT NOT NULL DEFAULT ''
)""",
"""CREATE UNIQUE INDEX IF NOT EXISTS idx_ttps_dedup ON ttps(dedup_key)""",
"""CREATE TABLE IF NOT EXISTS tradecraft (
id INTEGER PRIMARY KEY AUTOINCREMENT,
method TEXT NOT NULL DEFAULT '',
domain TEXT NOT NULL DEFAULT '',
description TEXT NOT NULL DEFAULT '',
operational_notes TEXT NOT NULL DEFAULT '',
countermeasures TEXT NOT NULL DEFAULT '',
source_doc_ids TEXT NOT NULL DEFAULT '[]',
source_context TEXT NOT NULL DEFAULT ''
)""",
"""CREATE TABLE IF NOT EXISTS entities (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL DEFAULT '',
type TEXT NOT NULL DEFAULT '',
description TEXT NOT NULL DEFAULT '',
relationships TEXT NOT NULL DEFAULT '[]',
source_doc_ids TEXT NOT NULL DEFAULT '[]',
source_context TEXT NOT NULL DEFAULT ''
)""",
"""CREATE TABLE IF NOT EXISTS infrastructure (
id INTEGER PRIMARY KEY AUTOINCREMENT,
indicator TEXT NOT NULL DEFAULT '',
indicator_type TEXT NOT NULL DEFAULT '',
context TEXT NOT NULL DEFAULT '',
source_doc_ids TEXT NOT NULL DEFAULT '[]',
source_context TEXT NOT NULL DEFAULT '',
collection TEXT NOT NULL DEFAULT '',
doc_title TEXT NOT NULL DEFAULT '',
doc_date TEXT
)""",
"""CREATE INDEX IF NOT EXISTS idx_infra_indicator ON infrastructure(indicator)""",
"""CREATE INDEX IF NOT EXISTS idx_infra_type ON infrastructure(indicator_type)""",
"""CREATE TABLE IF NOT EXISTS extraction_status (
doc_id TEXT NOT NULL DEFAULT '',
extractor TEXT NOT NULL DEFAULT '',
status TEXT NOT NULL DEFAULT 'PENDING',
error TEXT,
last_run TEXT,
token_count INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (doc_id, extractor)
)""",
"""CREATE TABLE IF NOT EXISTS token_usage (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp TEXT NOT NULL DEFAULT '',
source TEXT NOT NULL DEFAULT '',
doc_id TEXT NOT NULL DEFAULT '',
extractor TEXT NOT NULL DEFAULT '',
input_tokens INTEGER NOT NULL DEFAULT 0,
output_tokens INTEGER NOT NULL DEFAULT 0,
model TEXT NOT NULL DEFAULT ''
)""",
]
},
]
class MigrationRunner:
"""Run schema migrations on a SQLite database."""
def __init__(self, conn: sqlite3.Connection):
self.conn = conn
self._ensure_version_table()
def _ensure_version_table(self):
self.conn.execute("""
CREATE TABLE IF NOT EXISTS schema_version (
version INTEGER PRIMARY KEY,
description TEXT NOT NULL DEFAULT '',
applied_at TEXT NOT NULL DEFAULT (datetime('now'))
)
""")
self.conn.commit()
def get_current_version(self) -> int:
cursor = self.conn.execute(
"SELECT MAX(version) FROM schema_version"
)
row = cursor.fetchone()
return row[0] if row[0] is not None else 0
def run(self):
current = self.get_current_version()
pending = [m for m in MIGRATIONS if m['version'] > current]
if not pending:
logger.debug("Database schema is up to date (version %d)", current)
return
for migration in sorted(pending, key=lambda m: m['version']):
logger.info(
"Applying migration %d: %s",
migration['version'],
migration['description']
)
try:
for sql in migration['sql']:
self.conn.execute(sql)
self.conn.execute(
"INSERT INTO schema_version (version, description) VALUES (?, ?)",
(migration['version'], migration['description'])
)
self.conn.commit()
logger.info("Migration %d applied successfully", migration['version'])
except Exception as e:
self.conn.rollback()
logger.error("Migration %d failed: %s", migration['version'], e)
raise