diff --git a/migrations/versions/05486ce795d5_make_session_name_required_on_messages.py b/migrations/versions/05486ce795d5_make_session_name_required_on_messages.py index 6f239da9..f8bc4486 100644 --- a/migrations/versions/05486ce795d5_make_session_name_required_on_messages.py +++ b/migrations/versions/05486ce795d5_make_session_name_required_on_messages.py @@ -106,27 +106,63 @@ def upgrade() -> None: f"Created session peer association for peer '{peer_name}' in default session '{default_session_name}'" ) - # Step 4: Assign orphaned messages for this peer to the default session - op.execute( - sa.text(f""" - UPDATE {schema}.messages - SET session_name = '{default_session_name}' - WHERE workspace_name = '{workspace_name}' - AND peer_name = '{peer_name}' - AND session_name IS NULL - """) - ) + # Step 4: Assign orphaned messages for this peer to the default session in batches + batch_size = 5000 + while True: + result = conn.execute( + sa.text(f""" + WITH batch AS ( + SELECT id + FROM {schema}.messages + WHERE workspace_name = :workspace_name + AND peer_name = :peer_name + AND session_name IS NULL + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.messages m + SET session_name = :default_session_name + FROM batch + WHERE m.id = batch.id + """), + { + "workspace_name": workspace_name, + "peer_name": peer_name, + "default_session_name": default_session_name, + "batch_size": batch_size, + }, + ) + if result.rowcount == 0: + break - # Step 4.5: Handle orphaned message embeddings for this peer - op.execute( - sa.text(f""" - UPDATE {schema}.message_embeddings - SET session_name = '{default_session_name}' - WHERE workspace_name = '{workspace_name}' - AND peer_name = '{peer_name}' - AND session_name IS NULL - """) - ) + # Step 4.5: Handle orphaned message embeddings for this peer in batches + batch_size = 5000 + while True: + result = conn.execute( + sa.text(f""" + WITH batch AS ( + SELECT id + FROM {schema}.message_embeddings + WHERE workspace_name = :workspace_name + AND peer_name = :peer_name + AND session_name IS NULL + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.message_embeddings me + SET session_name = :default_session_name + FROM batch + WHERE me.id = batch.id + """), + { + "workspace_name": workspace_name, + "peer_name": peer_name, + "default_session_name": default_session_name, + "batch_size": batch_size, + }, + ) + if result.rowcount == 0: + break # Step 5: Sanity check that no orphaned messages remain remaining_orphaned = conn.execute( diff --git a/migrations/versions/08894082221a_replace_collection_name_with_observer_.py b/migrations/versions/08894082221a_replace_collection_name_with_observer_.py index 2e27a5a8..cadad00f 100644 --- a/migrations/versions/08894082221a_replace_collection_name_with_observer_.py +++ b/migrations/versions/08894082221a_replace_collection_name_with_observer_.py @@ -55,16 +55,30 @@ def upgrade() -> None: ), {"session_id": session_id, "workspace_name": workspace_name}, ) - # Update all documents with NULL session_name - connection.execute( - text( - f""" - UPDATE {schema}.documents - SET session_name = '__global_observations__' - WHERE session_name IS NULL - """ - ), - ) + # Update all documents with NULL session_name in batches + batch_size = 5000 + while True: + result = connection.execute( + text( + f""" + WITH batch AS ( + SELECT id + FROM {schema}.documents + WHERE session_name IS NULL + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.documents d + SET session_name = '__global_observations__' + FROM batch + WHERE d.id = batch.id + AND d.session_name IS NULL + """ + ), + {"batch_size": batch_size}, + ) + if result.rowcount == 0: + break op.alter_column("documents", "session_name", nullable=False, schema=schema) @@ -84,29 +98,98 @@ def upgrade() -> None: schema=schema, ) - # Step 2: Populate collections observer and observed from existing name field + # Step 2a: Identify collections that should be deleted (in memory) + # These are user-created collections that are not used in the system: + # - name = 'global_representation' + # - name starts with peer_name + "_" (pattern: observer_observed) + # - name ends with "_" + peer_name (pattern: observed_observer) + collections_to_delete = connection.execute( + text( + f""" + SELECT id, name, peer_name, workspace_name + FROM {schema}.collections + WHERE name != 'global_representation' + AND name NOT LIKE peer_name || '_%' + AND name NOT LIKE '%_' || peer_name + """ + ) + ).fetchall() + + # Step 2b: Delete documents that reference collections marked for deletion (in batches) + if collections_to_delete: + # Delete documents in batches + batch_size = 5000 + for i in range(0, len(collections_to_delete), batch_size): + batch = collections_to_delete[i : i + batch_size] + collection_ids = [row.id for row in batch] + + connection.execute( + text( + f""" + DELETE FROM {schema}.documents d + USING {schema}.collections c + WHERE d.collection_name = c.name + AND d.peer_name = c.peer_name + AND d.workspace_name = c.workspace_name + AND c.id = ANY(:collection_ids) + """ + ), + {"collection_ids": collection_ids}, + ) + + # Step 2c: Delete the collections identified in step 2a (in batches) + if collections_to_delete: + batch_size = 5000 + for i in range(0, len(collections_to_delete), batch_size): + batch = collections_to_delete[i : i + batch_size] + collection_ids = [row.id for row in batch] + + connection.execute( + text( + f""" + DELETE FROM {schema}.collections + WHERE id = ANY(:collection_ids) + """ + ), + {"collection_ids": collection_ids}, + ) + + # Step 2d: Populate collections observer and observed from existing name field in batches # The logic is: # - observer = peer_name (the exact peer ID) # - If name is "global_representation", observed = peer_name (self-observation) # - If name starts with peer_name + "_", extract the observed part (pattern: observer_observed) # - If name ends with "_" + peer_name, extract the first part (pattern: observed_observer) - # - Otherwise (legacy edge cases), observed = name itself - connection.execute( - text( - f""" - UPDATE {schema}.collections - SET - observer = peer_name, - observed = CASE - WHEN name = 'global_representation' THEN peer_name - WHEN name LIKE peer_name || '_%' THEN substring(name from length(peer_name) + 2) - WHEN name LIKE '%_' || peer_name THEN substring(name from 1 for length(name) - length(peer_name) - 1) - ELSE name - END - WHERE observer IS NULL OR observed IS NULL - """ + # - Any legacy edge cases will have been deleted in step 2a. + batch_size = 5000 + while True: + result = connection.execute( + text( + f""" + WITH batch AS ( + SELECT id + FROM {schema}.collections + WHERE observer IS NULL OR observed IS NULL + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.collections c + SET + observer = c.peer_name, + observed = CASE + WHEN c.name = 'global_representation' THEN c.peer_name + WHEN c.name LIKE c.peer_name || '_%' THEN substring(c.name from length(c.peer_name) + 2) + WHEN c.name LIKE '%_' || c.peer_name THEN substring(c.name from 1 for length(c.name) - length(c.peer_name) - 1) + ELSE c.peer_name + END + FROM batch + WHERE c.id = batch.id + """ + ), + {"batch_size": batch_size}, ) - ) + if result.rowcount == 0: + break # Step 3: Make collections observer and observed NOT NULL op.alter_column("collections", "observer", nullable=False, schema=schema) @@ -151,6 +234,7 @@ def upgrade() -> None: AND d.collection_name = c.name AND d.peer_name = c.peer_name AND d.workspace_name = c.workspace_name + AND (d.observer IS NULL OR d.observed IS NULL) """ ), {"batch_size": batch_size}, @@ -415,19 +499,33 @@ def downgrade() -> None: schema=schema, ) - # Step 5: Populate documents collection_name from observer and observed - connection.execute( - text( - f""" - UPDATE {schema}.documents - SET collection_name = CASE - WHEN observer = observed THEN 'global_representation' - ELSE observer || '_' || observed - END - WHERE collection_name IS NULL - """ + # Step 5: Populate documents collection_name from observer and observed in batches + batch_size = 5000 + while True: + result = connection.execute( + text( + f""" + WITH batch AS ( + SELECT id + FROM {schema}.documents + WHERE collection_name IS NULL + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.documents d + SET collection_name = CASE + WHEN d.observer = d.observed THEN 'global_representation' + ELSE d.observer || '_' || d.observed + END + FROM batch + WHERE d.id = batch.id + AND d.collection_name IS NULL + """ + ), + {"batch_size": batch_size}, ) - ) + if result.rowcount == 0: + break # Step 6: Make documents collection_name NOT NULL op.alter_column("documents", "collection_name", nullable=False, schema=schema) @@ -440,16 +538,30 @@ def downgrade() -> None: schema=schema, ) - # Populate peer_name with observed value - connection.execute( - text( - f""" - UPDATE {schema}.documents - SET peer_name = observed - WHERE peer_name IS NULL - """ + # Populate peer_name with observed value in batches + batch_size = 5000 + while True: + result = connection.execute( + text( + f""" + WITH batch AS ( + SELECT id + FROM {schema}.documents + WHERE peer_name IS NULL + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.documents d + SET peer_name = d.observed + FROM batch + WHERE d.id = batch.id + AND d.peer_name IS NULL + """ + ), + {"batch_size": batch_size}, ) - ) + if result.rowcount == 0: + break # Make peer_name NOT NULL op.alter_column("documents", "peer_name", nullable=False, schema=schema) diff --git a/migrations/versions/564ba40505c5_add_session_name_column_to_documents.py b/migrations/versions/564ba40505c5_add_session_name_column_to_documents.py index d6a5fde4..d7238385 100644 --- a/migrations/versions/564ba40505c5_add_session_name_column_to_documents.py +++ b/migrations/versions/564ba40505c5_add_session_name_column_to_documents.py @@ -34,16 +34,38 @@ def upgrade() -> None: schema=schema, ) - # Step 2: Migrate data from internal_metadata to session_name column - op.execute( - sa.text( - f""" - UPDATE {schema}.documents - SET session_name = internal_metadata->>'session_name' - WHERE internal_metadata ? 'session_name' - """ + # Step 2: Migrate data from internal_metadata to session_name column in batches + # Process in batches to avoid timeout with large datasets + # Only migrate documents that have 'session_name' key with a non-null, non-empty value + bind = op.get_bind() + batch_size = 5000 + + while True: + result = bind.execute( + sa.text( + f""" + WITH batch AS ( + SELECT id + FROM {schema}.documents + WHERE session_name IS NULL + AND internal_metadata ? 'session_name' + AND internal_metadata->>'session_name' IS NOT NULL + AND internal_metadata->>'session_name' != '' + ORDER BY id + LIMIT :batch_size + ) + UPDATE {schema}.documents d + SET session_name = d.internal_metadata->>'session_name' + FROM batch b + WHERE d.id = b.id + AND d.session_name IS NULL + """ + ), + {"batch_size": batch_size}, ) - ) + + if result.rowcount == 0: + break # Step 3: Create index on session_name for efficient querying if not index_exists("documents", "idx_documents_session_name", inspector):