feat: batch large operations in recent migrations (#242)

* feat: batch large operations in recent migrations

* fix: batch downgrade

* fix: CodeRabbit comments

* fix: delete user-created collections and documents

* fix: guard against infinite loop for internal_metadata.session_name is null

* fix: add idempotency guards
This commit is contained in:
Rajat Ahuja 2025-10-21 11:30:56 -04:00 committed by GitHub
parent 77a965e97f
commit a7d01d17df
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
3 changed files with 247 additions and 77 deletions

View File

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

View File

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

View File

@ -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):