Spaces:
Paused
Paused
| #!/usr/bin/env python3 | |
| """ | |
| Database Persistence Layer | |
| This module implements production-grade database persistence for collateral packets. | |
| No mocks, no simulations - real database operations. | |
| Architecture: | |
| 1. SQLAlchemy ORM models for all entities | |
| 2. Database connection management | |
| 3. CRUD operations for collateral packets | |
| 4. Transaction tracking and income recording | |
| 5. Risk assessment persistence | |
| """ | |
| from typing import Dict, Any, List, Optional | |
| from datetime import datetime | |
| from enum import Enum | |
| import logging | |
| import os | |
| from sqlalchemy import create_engine, Column, String, Float, Integer, DateTime, Boolean, Text, JSON, ForeignKey | |
| from sqlalchemy.ext.declarative import declarative_base | |
| from sqlalchemy.orm import sessionmaker, relationship | |
| from sqlalchemy.exc import SQLAlchemyError | |
| # Configure logging | |
| logging.basicConfig(level=logging.INFO) | |
| logger = logging.getLogger(__name__) | |
| Base = declarative_base() | |
| class CollateralTypeEnum(str, Enum): | |
| """Collateral type enum for database.""" | |
| CODE_LICENSE = "code_license" | |
| API_ACCESS = "api_access" | |
| WHITE_LABEL = "white_label" | |
| MAINTENANCE_CONTRACT = "maintenance_contract" | |
| SUPPORT_CONTRACT = "support_contract" | |
| TRAINING_CERTIFICATION = "training_certification" | |
| class IncomeSourceEnum(str, Enum): | |
| """Income source enum for database.""" | |
| LICENSING_FEES = "licensing_fees" | |
| API_USAGE_REVENUE = "api_usage_revenue" | |
| MAINTENANCE_REVENUE = "maintenance_revenue" | |
| SUPPORT_REVENUE = "support_revenue" | |
| TRAINING_REVENUE = "training_revenue" | |
| YIELD_FARMING = "yield_farming" | |
| STAKING_REWARDS = "staking_rewards" | |
| class PaymentRailEnum(str, Enum): | |
| """Payment rail enum for database.""" | |
| STRIPE = "stripe" | |
| SOLANA = "solana" | |
| class TransactionStatusEnum(str, Enum): | |
| """Transaction status enum for database.""" | |
| PENDING = "pending" | |
| COMPLETED = "completed" | |
| FAILED = "failed" | |
| REFUNDED = "refunded" | |
| class RiskLevelEnum(str, Enum): | |
| """Risk level enum for database.""" | |
| VERY_LOW = "very_low" | |
| LOW = "low" | |
| MEDIUM = "medium" | |
| HIGH = "high" | |
| VERY_HIGH = "very_high" | |
| class CollateralPacket(Base): | |
| """Collateral packet database model.""" | |
| __tablename__ = 'collateral_packets' | |
| id = Column(String, primary_key=True) | |
| asset_id = Column(String, nullable=False, index=True) | |
| collateral_type = Column(String, nullable=False) | |
| # Valuation | |
| base_value_usd = Column(Float, nullable=False) | |
| collateral_value_usd = Column(Float, nullable=False) | |
| risk_adjusted_value_usd = Column(Float, nullable=False) | |
| collateral_multiplier = Column(Float, nullable=False) | |
| confidence_score = Column(Float, nullable=False) | |
| valuation_method = Column(String, nullable=False) | |
| # Income | |
| total_expected_annual_income_usd = Column(Float, nullable=False) | |
| total_actual_annual_income_usd = Column(Float, default=0.0) | |
| combined_multiplier = Column(Float, nullable=False) | |
| # Metadata | |
| created_at = Column(DateTime, nullable=False, default=datetime.utcnow) | |
| expires_at = Column(DateTime, nullable=True) | |
| status = Column(String, nullable=False, default='active') | |
| # Relationships | |
| income_streams = relationship("IncomeStream", back_populates="collateral_packet", cascade="all, delete-orphan") | |
| transactions = relationship("Transaction", back_populates="collateral_packet", cascade="all, delete-orphan") | |
| risk_assessment = relationship("RiskAssessment", back_populates="collateral_packet", uselist=False, cascade="all, delete-orphan") | |
| class IncomeStream(Base): | |
| """Income stream database model.""" | |
| __tablename__ = 'income_streams' | |
| id = Column(String, primary_key=True) | |
| collateral_packet_id = Column(String, ForeignKey('collateral_packets.id'), nullable=False) | |
| source = Column(String, nullable=False) | |
| expected_annual_income_usd = Column(Float, nullable=False) | |
| actual_annual_income_usd = Column(Float, default=0.0) | |
| multiplier = Column(Float, nullable=False) | |
| start_date = Column(DateTime, nullable=False) | |
| end_date = Column(DateTime, nullable=True) | |
| active = Column(Boolean, default=True) | |
| last_payout_date = Column(DateTime, nullable=True) | |
| payout_frequency = Column(String, default='monthly') | |
| # Relationship | |
| collateral_packet = relationship("CollateralPacket", back_populates="income_streams") | |
| class Transaction(Base): | |
| """Transaction database model.""" | |
| __tablename__ = 'transactions' | |
| id = Column(String, primary_key=True) | |
| collateral_packet_id = Column(String, ForeignKey('collateral_packets.id'), nullable=True) | |
| rail = Column(String, nullable=False) | |
| amount_usd = Column(Float, nullable=False) | |
| currency = Column(String, nullable=False) | |
| status = Column(String, nullable=False) | |
| created_at = Column(DateTime, nullable=False, default=datetime.utcnow) | |
| completed_at = Column(DateTime, nullable=True) | |
| tx_metadata = Column(JSON, nullable=True) | |
| description = Column(Text, nullable=True) | |
| customer_id = Column(String, nullable=True) | |
| invoice_id = Column(String, nullable=True) | |
| # Relationship | |
| collateral_packet = relationship("CollateralPacket", back_populates="transactions") | |
| class RiskAssessment(Base): | |
| """Risk assessment database model.""" | |
| __tablename__ = 'risk_assessments' | |
| id = Column(String, primary_key=True) | |
| collateral_packet_id = Column(String, ForeignKey('collateral_packets.id'), nullable=True) | |
| asset_id = Column(String, nullable=False, index=True) | |
| overall_risk_score = Column(Float, nullable=False) | |
| risk_level = Column(String, nullable=False) | |
| category_scores = Column(JSON, nullable=False) | |
| risk_factors = Column(JSON, nullable=False) | |
| mitigation_recommendations = Column(JSON, nullable=False) | |
| risk_adjusted_return = Column(Float, nullable=False) | |
| confidence_interval_lower = Column(Float, nullable=False) | |
| confidence_interval_upper = Column(Float, nullable=False) | |
| stress_test_results = Column(JSON, nullable=False) | |
| assessment_date = Column(DateTime, nullable=False, default=datetime.utcnow) | |
| # Relationship | |
| collateral_packet = relationship("CollateralPacket", back_populates="risk_assessment") | |
| class DatabaseManager: | |
| """Database manager for all operations.""" | |
| def __init__(self, database_url: Optional[str] = None): | |
| """ | |
| Initialize database manager. | |
| Args: | |
| database_url: Database connection URL. If not provided, uses environment variable. | |
| """ | |
| self.database_url = database_url or os.environ.get( | |
| 'DATABASE_URL', | |
| 'sqlite:///catacomb_underwriter.db' | |
| ) | |
| self.engine = create_engine(self.database_url) | |
| self.SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=self.engine) | |
| # Create tables | |
| self.create_tables() | |
| logger.info(f"Database manager initialized with URL: {self.database_url}") | |
| def create_tables(self): | |
| """Create all database tables.""" | |
| try: | |
| Base.metadata.create_all(bind=self.engine) | |
| logger.info("Database tables created successfully") | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to create database tables: {e}") | |
| raise | |
| def get_session(self): | |
| """Get database session.""" | |
| return self.SessionLocal() | |
| # Collateral Packet Operations | |
| def create_collateral_packet(self, packet_data: Dict[str, Any]) -> CollateralPacket: | |
| """Create a new collateral packet.""" | |
| session = self.get_session() | |
| try: | |
| packet = CollateralPacket( | |
| id=packet_data['packet_id'], | |
| asset_id=packet_data['asset_id'], | |
| collateral_type=packet_data['collateral_type'], | |
| base_value_usd=packet_data['valuation']['base_value_usd'], | |
| collateral_value_usd=packet_data['valuation']['collateral_value_usd'], | |
| risk_adjusted_value_usd=packet_data['valuation']['risk_adjusted_value_usd'], | |
| collateral_multiplier=packet_data['valuation']['collateral_multiplier'], | |
| confidence_score=packet_data['valuation']['confidence_score'], | |
| valuation_method=packet_data['valuation']['valuation_method'], | |
| total_expected_annual_income_usd=packet_data['total_expected_annual_income_usd'], | |
| total_actual_annual_income_usd=packet_data['total_actual_annual_income_usd'], | |
| combined_multiplier=packet_data['combined_multiplier'], | |
| created_at=datetime.fromisoformat(packet_data['created_at']), | |
| expires_at=datetime.fromisoformat(packet_data['expires_at']) if packet_data.get('expires_at') else None, | |
| status=packet_data['status'], | |
| ) | |
| # Add income streams | |
| for stream_data in packet_data.get('income_streams', []): | |
| stream = IncomeStream( | |
| id=stream_data['stream_id'], | |
| collateral_packet_id=packet.id, | |
| source=stream_data['source'], | |
| expected_annual_income_usd=stream_data['expected_annual_income_usd'], | |
| actual_annual_income_usd=stream_data['actual_annual_income_usd'], | |
| multiplier=stream_data['multiplier'], | |
| start_date=datetime.fromisoformat(stream_data['start_date']), | |
| end_date=datetime.fromisoformat(stream_data['end_date']) if stream_data.get('end_date') else None, | |
| active=stream_data['active'], | |
| payout_frequency='monthly', | |
| ) | |
| packet.income_streams.append(stream) | |
| session.add(packet) | |
| session.commit() | |
| session.refresh(packet) | |
| logger.info(f"Created collateral packet: {packet.id}") | |
| return packet | |
| except SQLAlchemyError as e: | |
| session.rollback() | |
| logger.error(f"Failed to create collateral packet: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def get_collateral_packet(self, packet_id: str) -> Optional[CollateralPacket]: | |
| """Get collateral packet by ID.""" | |
| session = self.get_session() | |
| try: | |
| packet = session.query(CollateralPacket).filter( | |
| CollateralPacket.id == packet_id | |
| ).first() | |
| return packet | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to get collateral packet: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def list_collateral_packets( | |
| self, | |
| asset_id: Optional[str] = None, | |
| status: Optional[str] = None, | |
| limit: int = 100 | |
| ) -> List[CollateralPacket]: | |
| """List collateral packets with optional filters.""" | |
| session = self.get_session() | |
| try: | |
| query = session.query(CollateralPacket) | |
| if asset_id: | |
| query = query.filter(CollateralPacket.asset_id == asset_id) | |
| if status: | |
| query = query.filter(CollateralPacket.status == status) | |
| packets = query.limit(limit).all() | |
| return packets | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to list collateral packets: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def update_collateral_packet(self, packet_id: str, updates: Dict[str, Any]) -> Optional[CollateralPacket]: | |
| """Update collateral packet.""" | |
| session = self.get_session() | |
| try: | |
| packet = session.query(CollateralPacket).filter( | |
| CollateralPacket.id == packet_id | |
| ).first() | |
| if not packet: | |
| return None | |
| for key, value in updates.items(): | |
| if hasattr(packet, key): | |
| setattr(packet, key, value) | |
| session.commit() | |
| session.refresh(packet) | |
| logger.info(f"Updated collateral packet: {packet_id}") | |
| return packet | |
| except SQLAlchemyError as e: | |
| session.rollback() | |
| logger.error(f"Failed to update collateral packet: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def delete_collateral_packet(self, packet_id: str) -> bool: | |
| """Delete collateral packet.""" | |
| session = self.get_session() | |
| try: | |
| packet = session.query(CollateralPacket).filter( | |
| CollateralPacket.id == packet_id | |
| ).first() | |
| if not packet: | |
| return False | |
| session.delete(packet) | |
| session.commit() | |
| logger.info(f"Deleted collateral packet: {packet_id}") | |
| return True | |
| except SQLAlchemyError as e: | |
| session.rollback() | |
| logger.error(f"Failed to delete collateral packet: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| # Transaction Operations | |
| def create_transaction(self, transaction_data: Dict[str, Any]) -> Transaction: | |
| """Create a new transaction.""" | |
| session = self.get_session() | |
| try: | |
| transaction = Transaction( | |
| id=transaction_data['transaction_id'], | |
| collateral_packet_id=transaction_data.get('collateral_packet_id'), | |
| rail=transaction_data['rail'], | |
| amount_usd=transaction_data['amount_usd'], | |
| currency=transaction_data['currency'], | |
| status=transaction_data['status'], | |
| created_at=datetime.fromisoformat(transaction_data['created_at']), | |
| completed_at=datetime.fromisoformat(transaction_data['completed_at']) if transaction_data.get('completed_at') else None, | |
| metadata=transaction_data.get('metadata'), | |
| description=transaction_data.get('description'), | |
| customer_id=transaction_data.get('customer_id'), | |
| invoice_id=transaction_data.get('invoice_id'), | |
| ) | |
| session.add(transaction) | |
| session.commit() | |
| session.refresh(transaction) | |
| logger.info(f"Created transaction: {transaction.id}") | |
| return transaction | |
| except SQLAlchemyError as e: | |
| session.rollback() | |
| logger.error(f"Failed to create transaction: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def get_transaction(self, transaction_id: str) -> Optional[Transaction]: | |
| """Get transaction by ID.""" | |
| session = self.get_session() | |
| try: | |
| transaction = session.query(Transaction).filter( | |
| Transaction.id == transaction_id | |
| ).first() | |
| return transaction | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to get transaction: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def list_transactions( | |
| self, | |
| collateral_packet_id: Optional[str] = None, | |
| rail: Optional[str] = None, | |
| status: Optional[str] = None, | |
| limit: int = 100 | |
| ) -> List[Transaction]: | |
| """List transactions with optional filters.""" | |
| session = self.get_session() | |
| try: | |
| query = session.query(Transaction) | |
| if collateral_packet_id: | |
| query = query.filter(Transaction.collateral_packet_id == collateral_packet_id) | |
| if rail: | |
| query = query.filter(Transaction.rail == rail) | |
| if status: | |
| query = query.filter(Transaction.status == status) | |
| transactions = query.limit(limit).all() | |
| return transactions | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to list transactions: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| # Risk Assessment Operations | |
| def create_risk_assessment(self, assessment_data: Dict[str, Any]) -> RiskAssessment: | |
| """Create a new risk assessment.""" | |
| session = self.get_session() | |
| try: | |
| assessment = RiskAssessment( | |
| id=f"risk_{assessment_data['asset_id']}_{datetime.utcnow().timestamp()}", | |
| collateral_packet_id=assessment_data.get('collateral_packet_id'), | |
| asset_id=assessment_data['asset_id'], | |
| overall_risk_score=assessment_data['overall_risk_score'], | |
| risk_level=assessment_data['risk_level'], | |
| category_scores=assessment_data['category_scores'], | |
| risk_factors=assessment_data['risk_factors'], | |
| mitigation_recommendations=assessment_data['mitigation_recommendations'], | |
| risk_adjusted_return=assessment_data['risk_adjusted_return'], | |
| confidence_interval_lower=assessment_data['confidence_interval']['lower'], | |
| confidence_interval_upper=assessment_data['confidence_interval']['upper'], | |
| stress_test_results=assessment_data['stress_test_results'], | |
| assessment_date=datetime.fromisoformat(assessment_data['assessment_date']), | |
| ) | |
| session.add(assessment) | |
| session.commit() | |
| session.refresh(assessment) | |
| logger.info(f"Created risk assessment for asset: {assessment_data['asset_id']}") | |
| return assessment | |
| except SQLAlchemyError as e: | |
| session.rollback() | |
| logger.error(f"Failed to create risk assessment: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def get_risk_assessment(self, asset_id: str) -> Optional[RiskAssessment]: | |
| """Get risk assessment by asset ID.""" | |
| session = self.get_session() | |
| try: | |
| assessment = session.query(RiskAssessment).filter( | |
| RiskAssessment.asset_id == asset_id | |
| ).order_by(RiskAssessment.assessment_date.desc()).first() | |
| return assessment | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to get risk assessment: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| # Income Recording Operations | |
| def record_income( | |
| self, | |
| packet_id: str, | |
| stream_id: str, | |
| amount_usd: float, | |
| timestamp: Optional[datetime] = None | |
| ) -> bool: | |
| """Record income for a stream.""" | |
| session = self.get_session() | |
| try: | |
| # Get the income stream | |
| stream = session.query(IncomeStream).filter( | |
| IncomeStream.id == stream_id, | |
| IncomeStream.collateral_packet_id == packet_id | |
| ).first() | |
| if not stream: | |
| logger.error(f"Income stream not found: {stream_id}") | |
| return False | |
| # Update actual income | |
| stream.actual_annual_income_usd += amount_usd | |
| stream.last_payout_date = timestamp or datetime.utcnow() | |
| # Update packet total | |
| packet = session.query(CollateralPacket).filter( | |
| CollateralPacket.id == packet_id | |
| ).first() | |
| if packet: | |
| packet.total_actual_annual_income_usd += amount_usd | |
| session.commit() | |
| logger.info(f"Recorded ${amount_usd:,.2f} income for stream {stream_id}") | |
| return True | |
| except SQLAlchemyError as e: | |
| session.rollback() | |
| logger.error(f"Failed to record income: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| def get_income_summary(self, packet_id: str) -> Dict[str, Any]: | |
| """Get income summary for a packet.""" | |
| session = self.get_session() | |
| try: | |
| packet = session.query(CollateralPacket).filter( | |
| CollateralPacket.id == packet_id | |
| ).first() | |
| if not packet: | |
| return { | |
| 'total_income_usd': 0.0, | |
| 'transaction_count': 0, | |
| 'last_income_date': None, | |
| } | |
| # Get transactions for this packet | |
| transactions = session.query(Transaction).filter( | |
| Transaction.collateral_packet_id == packet_id, | |
| Transaction.status == 'completed' | |
| ).all() | |
| total_income = sum(t.amount_usd for t in transactions) | |
| last_date = max((t.completed_at for t in transactions if t.completed_at), default=None) | |
| return { | |
| 'total_income_usd': round(total_income, 2), | |
| 'transaction_count': len(transactions), | |
| 'last_income_date': last_date.isoformat() if last_date else None, | |
| } | |
| except SQLAlchemyError as e: | |
| logger.error(f"Failed to get income summary: {e}") | |
| raise | |
| finally: | |
| session.close() | |
| # Global database manager instance | |
| db_manager: Optional[DatabaseManager] = None | |
| def get_database_manager() -> DatabaseManager: | |
| """Get or create global database manager.""" | |
| global db_manager | |
| if db_manager is None: | |
| db_manager = DatabaseManager() | |
| return db_manager | |