import psycopg2 import psycopg2.extras import uuid import json import re from datetime import datetime from typing import Optional, List class ViewMigrator: def __init__(self, source_conn_params: dict, target_conn_params: dict): self.source_conn = psycopg2.connect(**source_conn_params) self.target_conn = psycopg2.connect(**target_conn_params) def migrate_views(self): """Main migration function""" try: # Fetch data from source old_views = self._fetch_old_views() print(f"Found {len(old_views)} views to migrate") # Migrate each view for old_view in old_views: self._migrate_single_view(old_view) print("Migration completed successfully!") except Exception as e: print(f"Migration failed: {e}") raise finally: self.source_conn.close() self.target_conn.close() def _fetch_old_views(self) -> List[dict]: """Fetch all views from the old table""" cursor = self.source_conn.cursor() query = ''' SELECT "Id", "Title", "Description", "Definition", "SysInserted", "SysUpdated", "SysInsertedUser", "SysUpdatedUser", "Parameter" FROM public."View" ORDER BY "Id" ''' cursor.execute(query) columns = [desc[0] for desc in cursor.description] rows = cursor.fetchall() return [dict(zip(columns, row)) for row in rows] def _migrate_single_view(self, old_view: dict): """Migrate a single view record""" # Generate new UUID for the view view_id = uuid.uuid4() title = self._get_title(old_view['Title']) alias = self._get_alias(title) parameters = self._parse_parameters(old_view.get('Parameter')) # Insert into new View table self._insert_view( view_id=view_id, title=title, description=old_view.get('Description'), alias=alias, created_at=datetime.now(), updated_at=datetime.now(), created_by=old_view.get('SysInsertedUser', 'migration_script'), updated_by=old_view.get('SysUpdatedUser', 'migration_script') ) # Insert into ViewVersion table self._insert_view_version( view_id=view_id, version=1, state="draft", published_at=None, created_at=datetime.now(), updated_at=datetime.now(), created_by=old_view.get('SysInsertedUser', 'migration_script'), updated_by=old_view.get('SysUpdatedUser', 'migration_script') ) # Insert into ViewData table self._insert_view_data( view_id=view_id, version=1, query=old_view.get('Definition', ''), parameters=parameters ) print(f"Migrated view: {old_view['Title']} (ID: {old_view['Id']} -> {view_id})") def _get_values_of_column_from_db(self, column: str) -> list: cursor = self.target_conn.cursor() query = f''' SELECT "{column}" FROM public."View" ORDER BY "Id" ''' cursor.execute(query) rows = cursor.fetchall() return [item[0] for item in rows] def _get_title(self, title: str) -> str: titles_from_db = self._get_values_of_column_from_db("Title") return self._check_title(title, titles_from_db, 1) def _check_title(self, title: str, titles_from_db: list, counter: int) -> str: while title in titles_from_db: title = f"{title} ({counter})" counter+=1 self._check_title(title, titles_from_db, counter) return title def _get_alias(self, title: str) -> str: aliases_from_db = self._get_values_of_column_from_db("Alias") alias = self._create_alias_from_title(title) return self._check_alias(alias, aliases_from_db, 1) def _check_alias(self, alias: str, aliases_from_db: list, counter: int) -> str: while alias in aliases_from_db: alias = f"{alias}-{counter}" counter+=1 self._check_alias(alias, aliases_from_db, counter) return alias def _create_alias_from_title(self, title: str) -> str: if not title: return "" alias = title.lower() alias = re.sub(r'[^a-z0-9-]', '-', alias) alias = re.sub(r'-+', '-', alias) alias = alias.strip('-') return alias def _parse_parameters(self, parameter_text: Optional[str]) -> List[str]: """Parse the old Parameter field into a string array""" if not parameter_text: return [] try: # Try to parse as JSON array first if parameter_text.strip().startswith('['): return json.loads(parameter_text) except (json.JSONDecodeError, AttributeError): pass # If not JSON, try to split by common delimiters if ',' in parameter_text: return [p.strip() for p in parameter_text.split(',')] elif ';' in parameter_text: return [p.strip() for p in parameter_text.split(';')] elif '|' in parameter_text: return [p.strip() for p in parameter_text.split('|')] else: # Single parameter or unknown format return [parameter_text.strip()] def _insert_view(self, view_id: uuid.UUID, title: str, description: str, alias: str, created_at: datetime, updated_at: datetime, created_by: str, updated_by: str): """Insert into the new View table""" cursor = self.target_conn.cursor() query = ''' INSERT INTO public."View" ("Id", "Title", "Description", "Alias", "CreatedAt", "UpdatedAt", "CreatedBy", "UpdatedBy") VALUES (%s, %s, %s, %s, %s, %s, %s, %s) ''' cursor.execute(query, (view_id, title, description, alias, created_at, updated_at, created_by, updated_by)) self.target_conn.commit() def _insert_view_version(self, view_id: uuid.UUID, version: int, state: str, published_at: Optional[datetime], created_at: datetime, updated_at: datetime, created_by: str, updated_by: str): """Insert into the ViewVersion table""" cursor = self.target_conn.cursor() query = ''' INSERT INTO public."ViewVersion" ("ViewId", "Version", "State", "PublishedAt", "CreatedAt", "UpdatedAt", "CreatedBy", "UpdatedBy") VALUES (%s, %s, %s, %s, %s, %s, %s, %s) ''' cursor.execute(query, (view_id, version, state, published_at, created_at, updated_at, created_by, updated_by)) self.target_conn.commit() def _insert_view_data(self, view_id: uuid.UUID, version: int, query: str, parameters: List[str]): """Insert into the ViewData table""" cursor = self.target_conn.cursor() insert_query = ''' INSERT INTO public."ViewData" ("ViewId", "Version", "Query", "Parameters") VALUES (%s, %s, %s, %s) ''' cursor.execute(insert_query, (view_id, version, query, parameters)) self.target_conn.commit() # Usage example if __name__ == "__main__": # Source database connection parameters (PostgreSQL with old View table) source_params = { 'host': '', 'database': 'dch_repo', 'user': 'postgres', 'password': '', 'port': 5432 } # Target database connection parameters (new database structure) target_params = { 'host': '', 'database': 'linked_data_api', 'user': 'postgres', 'password': '', 'port': 5432 } psycopg2.extras.register_uuid() # Run migration migrator = ViewMigrator(source_params, target_params) migrator.migrate_views()