0% found this document useful (0 votes)
4 views7 pages

SQL Code

The document outlines a Python script that sets up a SQLAlchemy ORM with a relational schema for managing articles, entities, and cascade events. It includes definitions for database tables, triggers, views, and stored procedures to handle data integrity and auditing. Additionally, it provides methods for inserting data, checking SQL cache, and retrieving recent cascades from the database.

Uploaded by

tipu
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd
0% found this document useful (0 votes)
4 views7 pages

SQL Code

The document outlines a Python script that sets up a SQLAlchemy ORM with a relational schema for managing articles, entities, and cascade events. It includes definitions for database tables, triggers, views, and stored procedures to handle data integrity and auditing. Additionally, it provides methods for inserting data, checking SQL cache, and retrieving recent cascades from the database.

Uploaded by

tipu
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd

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()

You might also like