import os import logging from dotenv import load_dotenv from sqlalchemy import create_engine, Column, Integer, String, Float,
DateTime, text, ForeignKey from [Link]
import declarative_base, sessionmaker, relationship from [Link] import func
Load environment variables
load_dotenv()
[Link](level=[Link], format="%(asctime)s [%(levelname)s] %(message)s") logger = [Link](name)
Initialize the declarative base for
SQLAlchemy ORM
Base = declarative_base()
---------------------------------------------------------
1. Define the Relational Schema (Normalized
to 3NF)
---------------------------------------------------------
class Article(Base): tablename = 'articles' id = Column(Integer, primary_key=True, autoincrement=True) title = Column(String(500), nullable=False, unique=True)
class Entity(Base): tablename = 'entities' id = Column(Integer, primary_key=True, autoincrement=True) name = Column(String(255), nullable=False, unique=True, index=True)
class CascadeEvent(Base): tablename = 'cascade_events' id = Column(Integer, primary_key=True, autoincrement=True)
# Foreign Keys for Referential Integrity
source_id = Column(Integer, ForeignKey('[Link]'), nullable=False)
target_id = Column(Integer, ForeignKey('[Link]'), nullable=False)
article_id = Column(Integer, ForeignKey('[Link]'), nullable=False)
relationship_type = Column(String(255), nullable=False)
probability_score = Column(Float, default=0.5)
time_lag_days = Column(Integer, default=0)
# Timestamps for our trigger to manage
created_at = Column(DateTime, server_default=[Link]())
last_updated = Column(DateTime, server_default=[Link]())
class AuditLog(Base): tablename = 'audit_logs' id = Column(Integer, primary_key=True, autoincrement=True) table_name = Column(String(50)) record_id = Column(Integer) action =
Column(String(50)) old_value = Column(String(255)) new_value = Column(String(255)) changed_at = Column(DateTime, server_default=[Link]())
---------------------------------------------------------
2. The Database Manager
---------------------------------------------------------
class SQLBuilder: def init(self): db_url = [Link]("DATABASE_URL") if not db_url: raise ValueError("DATABASE_URL not found in .env file.")
# Ensure pymysql is used if using MySQL
if db_url.startswith("mysql://"):
db_url = db_url.replace("mysql://", "mysql+pymysql://", 1)
try:
[Link] = create_engine(db_url, echo=False)
[Link] = sessionmaker(autocommit=False, autoflush=False, bind=[Link])
[Link]("SQLBuilder successfully connected to MySQL.")
self.initialize_schema()
except Exception as e:
[Link](f"Failed to connect to MySQL: {e}")
raise
def initialize_schema(self):
"""
Generates the 3NF tables, triggers, and stored procedures.
"""
try:
[Link].create_all(bind=[Link])
[Link]("MySQL base tables validated/created successfully.")
self._create_database_logic()
except Exception as e:
[Link](f"Failed to create schema: {e}")
def _create_database_logic(self):
"""
Executes raw MySQL to inject triggers, procedures, and cursors.
"""
# Trigger: Update timestamp before update
trigger_drop_sql = "DROP TRIGGER IF EXISTS set_timestamp;"
trigger_sql = """
CREATE TRIGGER set_timestamp
BEFORE UPDATE ON cascade_events
FOR EACH ROW
BEGIN
SET NEW.last_updated = CURRENT_TIMESTAMP;
END;
"""
# Trigger 1: AFTER UPDATE
trigger_audit_update_drop = "DROP TRIGGER IF EXISTS audit_cascade_update;"
trigger_audit_update_sql = """
CREATE TRIGGER audit_cascade_update
AFTER UPDATE ON cascade_events
FOR EACH ROW
BEGIN
IF OLD.probability_score != NEW.probability_score THEN
INSERT INTO audit_logs (table_name, record_id, action, old_value, new_value)
VALUES ('cascade_events', [Link], 'UPDATE_PROB', CAST(OLD.probability_score AS CHAR), CAST(NEW.probability_score AS CHAR)
END IF;
END;
"""
# Trigger 2: AFTER DELETE
trigger_audit_delete_drop = "DROP TRIGGER IF EXISTS audit_cascade_delete;"
trigger_audit_delete_sql = """
CREATE TRIGGER audit_cascade_delete
AFTER DELETE ON cascade_events
FOR EACH ROW
BEGIN
INSERT INTO audit_logs (table_name, record_id, action, old_value, new_value)
VALUES ('cascade_events', [Link], 'DELETE', CONCAT('Src:', CAST(OLD.source_id AS CHAR), ' Tgt:', CAST(OLD.target_id AS CHAR))
END;
"""
# View: High Risk Cascades
view_drop = "DROP VIEW IF EXISTS high_risk_cascades_view;"
view_sql = """
CREATE VIEW high_risk_cascades_view AS
SELECT [Link], [Link] AS source, [Link] AS target, c.relationship_type, c.probability_score, [Link]
FROM cascade_events c
JOIN entities e1 ON c.source_id = [Link]
JOIN entities e2 ON c.target_id = [Link]
JOIN articles a ON c.article_id = [Link]
WHERE c.probability_score > 0.8;
"""
# Stored Procedure 1: Bulk Update (Uses Joins)
procedure_1_drop = "DROP PROCEDURE IF EXISTS bulk_update_probabilities;"
procedure_1_sql = """
CREATE PROCEDURE bulk_update_probabilities(IN t_name VARCHAR(255), IN new_prob FLOAT)
BEGIN
UPDATE cascade_events c
JOIN entities e ON c.target_id = [Link]
SET c.probability_score = new_prob
WHERE [Link] = t_name;
END;
"""
# Stored Procedure 2: Cursor-based Impact Report using Recursive CTE for Downstream Causal Traversal
procedure_2_drop = "DROP PROCEDURE IF EXISTS generate_impact_report;"
procedure_2_sql = """
CREATE PROCEDURE generate_impact_report(IN t_name VARCHAR(255))
BEGIN
DECLARE done INT DEFAULT FALSE;
DECLARE src_name VARCHAR(255);
DECLARE tgt_name VARCHAR(255);
DECLARE rel_type VARCHAR(255);
DECLARE prob FLOAT;
DECLARE d_val INT;
DECLARE path_str TEXT;
DECLARE total_risk FLOAT DEFAULT 0;
DECLARE event_count INT DEFAULT 0;
DECLARE cascade_cursor CURSOR FOR
WITH RECURSIVE de_duped_events AS (
SELECT source_id, target_id, relationship_type, MAX(probability_score) AS probability_score
FROM cascade_events
GROUP BY source_id, target_id, relationship_type
),
cascade_path AS (
SELECT
c.source_id,
c.target_id,
[Link] AS source_name,
[Link] AS target_name,
c.relationship_type,
c.probability_score AS path_prob,
1 AS depth,
CONCAT([Link], ' -> ', [Link]) AS path_string
FROM de_duped_events c
JOIN entities e1 ON c.source_id = [Link]
JOIN entities e2 ON c.target_id = [Link]
WHERE [Link] = t_name
UNION ALL
SELECT
c.source_id,
c.target_id,
cp.target_name AS source_name,
[Link] AS target_name,
c.relationship_type,
ROUND(cp.path_prob * c.probability_score, 2) AS path_prob,
[Link] + 1 AS depth,
CONCAT(cp.path_string, ' -> ', [Link]) AS path_string
FROM de_duped_events c
JOIN entities e2 ON c.target_id = [Link]
JOIN cascade_path cp ON c.source_id = cp.target_id
WHERE [Link] < 7 AND INSTR(cp.path_string, [Link]) = 0
)
SELECT source_name, target_name, relationship_type, path_prob, depth, path_string
FROM (
SELECT
source_name,
target_name,
relationship_type,
path_prob,
depth,
path_string,
ROW_NUMBER() OVER(PARTITION BY target_name ORDER BY path_prob DESC, depth ASC) as rn
FROM cascade_path
) t
WHERE rn = 1
ORDER BY depth ASC, path_prob DESC;
DECLARE CONTINUE HANDLER FOR NOT FOUND SET done = TRUE;
CREATE TEMPORARY TABLE IF NOT EXISTS impact_report_logs (
log_id INT AUTO_INCREMENT PRIMARY KEY,
message TEXT
);
TRUNCATE TABLE impact_report_logs;
INSERT INTO impact_report_logs (message) VALUES (CONCAT('--- Generating Impact Report for ', t_name, ' ---'));
OPEN cascade_cursor;
read_loop: LOOP
FETCH cascade_cursor INTO src_name, tgt_name, rel_type, prob, d_val, path_str;
IF done THEN
LEAVE read_loop;
END IF;
SET event_count = event_count + 1;
SET total_risk = total_risk + prob;
INSERT INTO impact_report_logs (message) VALUES (CONCAT('Impact ', event_count, ': ', path_str, ' via ', rel_type, ' (Con
END LOOP;
CLOSE cascade_cursor;
IF event_count > 0 THEN
INSERT INTO impact_report_logs (message) VALUES (CONCAT('Total Events: ', event_count, '. Average Path Confidence: ', ROU
ELSE
INSERT INTO impact_report_logs (message) VALUES (CONCAT('No downstream ripple effects found starting from ', t_name, '.')
END IF;
SELECT message FROM impact_report_logs ORDER BY log_id;
DROP TEMPORARY TABLE IF EXISTS impact_report_logs;
END;
"""
# Stored Procedure 3: Transactional Merge Entities
procedure_3_drop = "DROP PROCEDURE IF EXISTS merge_entities;"
procedure_3_sql = """
CREATE PROCEDURE merge_entities(IN old_ent VARCHAR(255), IN new_ent VARCHAR(255))
BEGIN
DECLARE old_id INT;
DECLARE new_id INT;
DECLARE EXIT HANDLER FOR SQLEXCEPTION
BEGIN
ROLLBACK;
SELECT 'Error occurred. Transaction rolled back.' AS message;
END;
START TRANSACTION;
SELECT id INTO old_id FROM entities WHERE name = old_ent LIMIT 1;
SELECT id INTO new_id FROM entities WHERE name = new_ent LIMIT 1;
IF old_id IS NULL OR new_id IS NULL THEN
ROLLBACK;
SELECT CONCAT('One or both entities do not exist: ', old_ent, ' / ', new_ent) AS message;
ELSE
UPDATE cascade_events SET source_id = new_id WHERE source_id = old_id;
UPDATE cascade_events SET target_id = new_id WHERE target_id = old_id;
DELETE FROM entities WHERE id = old_id;
COMMIT;
SELECT CONCAT('Successfully merged ', old_ent, ' into ', new_ent, '.') AS message;
END IF;
END;
"""
with [Link]() as conn:
[Link](text(trigger_drop_sql))
[Link](text(trigger_sql))
[Link](text(trigger_audit_update_drop))
[Link](text(trigger_audit_update_sql))
[Link](text(trigger_audit_delete_drop))
[Link](text(trigger_audit_delete_sql))
[Link](text(view_drop))
[Link](text(view_sql))
[Link](text(procedure_1_drop))
[Link](text(procedure_1_sql))
[Link](text(procedure_2_drop))
[Link](text(procedure_2_sql))
[Link](text(procedure_3_drop))
[Link](text(procedure_3_sql))
[Link]("MySQL Triggers, Views, and Procedures successfully installed.")
def insert_cascades(self, article_title: str, cascades: list):
"""
Inserts dynamic extracted data respecting the 3NF relationships.
"""
if not cascades:
return
session = [Link]()
try:
def get_or_create_entity(name):
ent = [Link](Entity).filter_by(name=name).first()
if not ent:
ent = Entity(name=name)
[Link](ent)
[Link]()
return ent
article = [Link](Article).filter_by(title=article_title).first()
if not article:
article = Article(title=article_title)
[Link](article)
[Link]()
for cascade in cascades:
source_ent = get_or_create_entity([Link]("source_node", "UNKNOWN").upper())
target_ent = get_or_create_entity([Link]("target_node", "UNKNOWN").upper())
new_cascade = CascadeEvent(
source_id=source_ent.id,
target_id=target_ent.id,
article_id=[Link],
relationship_type=[Link]("relationship", "AFFECTS"),
probability_score=[Link]("probability_score", 0.5),
time_lag_days=[Link]("time_lag_days", 0)
)
[Link](new_cascade)
[Link]()
[Link](f"Successfully inserted {len(cascades)} cascades into {CascadeEvent.__tablename__}.")
except Exception as e:
[Link]()
[Link](f"Error inserting row: {e}")
finally:
[Link]()
def check_sql_cache(self, headline):
session = [Link]()
try:
query = text("""
SELECT [Link] AS source_node, [Link] AS target_node, c.relationship_type, c.probability_score
FROM cascade_events c
JOIN articles a ON c.article_id = [Link]
JOIN entities e1 ON c.source_id = [Link]
JOIN entities e2 ON c.target_id = [Link]
WHERE [Link] LIKE :headline
""")
results = [Link](query, {"headline": f"%{headline}%"}).fetchall()
if results:
[Link](f" CACHE HIT: Found {len(results)} existing cascades for '{headline}'.")
return [dict(row._mapping) for row in results]
[Link](" CACHE MISS: Headline not found in SQL memory.")
return None
finally:
[Link]()
def get_recent_cascades(self, limit: int = 20):
session = [Link]()
try:
query = text("""
SELECT [Link], [Link] AS article_title, [Link] AS source_node, [Link] AS target_node,
c.relationship_type, c.probability_score, c.time_lag_days, c.created_at, c.last_updated
FROM cascade_events c
JOIN articles a ON c.article_id = [Link]
JOIN entities e1 ON c.source_id = [Link]
JOIN entities e2 ON c.target_id = [Link]
ORDER BY c.created_at DESC
LIMIT :limit
""")
results = [Link](query, {"limit": limit}).fetchall()
return [dict(row._mapping) for row in results]
except Exception as e:
[Link](f"Failed to fetch recent cascades: {e}")
return []
finally:
[Link]()
if name == "main": [Link]("Starting SQL Builder skeleton...") builder = SQLBuilder() builder.initialize_schema()