From f8e7167bdfddf3aea0f92cce284b607a029434b1 Mon Sep 17 00:00:00 2001 From: Vineeth Voruganti <13438633+VVoruganti@users.noreply.github.com> Date: Fri, 16 Jan 2026 13:28:00 -0500 Subject: [PATCH] fix: Migration naming and long held connection --- ...19a52b73c60_support_external_embeddings.py | 10 +- ...b1f_add_unique_index_for_pending_dreams.py | 4 +- src/crud/document.py | 4 + src/models.py | 26 +++ src/reconciler/sync_vectors.py | 197 ++++++++++++------ ...19a52b73c60_support_external_embeddings.py | 2 +- ...b1f_add_unique_index_for_pending_dreams.py | 2 +- 7 files changed, 171 insertions(+), 74 deletions(-) diff --git a/migrations/versions/119a52b73c60_support_external_embeddings.py b/migrations/versions/119a52b73c60_support_external_embeddings.py index 23a688a6..ea8aed97 100644 --- a/migrations/versions/119a52b73c60_support_external_embeddings.py +++ b/migrations/versions/119a52b73c60_support_external_embeddings.py @@ -187,9 +187,11 @@ def upgrade() -> None: # Add partial unique index on queue table for reconciler task deduplication # This ensures only one pending reconciler task exists per work_unit_key - if not index_exists("queue", "uq_queue_work_unit_key", inspector): + if not index_exists( + "queue", "uq_queue_reconciler_pending_work_unit_key", inspector + ): op.create_index( - "uq_queue_work_unit_key", + "uq_queue_reconciler_pending_work_unit_key", "queue", ["work_unit_key"], unique=True, @@ -292,9 +294,9 @@ def downgrade() -> None: ) # Drop reconciler queue index if it exists - if index_exists("queue", "uq_queue_work_unit_key", inspector): + if index_exists("queue", "uq_queue_reconciler_pending_work_unit_key", inspector): op.drop_index( - "uq_queue_work_unit_key", + "uq_queue_reconciler_pending_work_unit_key", table_name="queue", schema=schema, ) diff --git a/migrations/versions/7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py b/migrations/versions/7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py index 5d669652..fcec958f 100644 --- a/migrations/versions/7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py +++ b/migrations/versions/7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py @@ -28,7 +28,7 @@ schema = get_schema() def upgrade() -> None: """Add a partial unique index to prevent duplicate pending dream queue items.""" op.create_index( - "ux_queue_dream_pending_work_unit_key", + "uq_queue_dream_pending_work_unit_key", "queue", ["work_unit_key"], unique=True, @@ -40,5 +40,5 @@ def upgrade() -> None: def downgrade() -> None: """Drop the partial unique index for pending dream queue items.""" op.drop_index( - "ux_queue_dream_pending_work_unit_key", table_name="queue", schema=schema + "uq_queue_dream_pending_work_unit_key", table_name="queue", schema=schema ) diff --git a/src/crud/document.py b/src/crud/document.py index aa8697e1..0f59d4b2 100644 --- a/src/crud/document.py +++ b/src/crud/document.py @@ -140,6 +140,7 @@ async def query_documents_recent( models.Document.workspace_name == workspace_name, models.Document.observer == observer, models.Document.observed == observed, + models.Document.deleted_at.is_(None), ) if session_name is not None: @@ -178,6 +179,7 @@ async def query_documents_most_derived( models.Document.workspace_name == workspace_name, models.Document.observer == observer, models.Document.observed == observed, + models.Document.deleted_at.is_(None), ) .order_by(models.Document.times_derived.desc()) .limit(limit) @@ -956,6 +958,7 @@ async def get_documents_by_ids( stmt = select(models.Document).where( models.Document.workspace_name == workspace_name, models.Document.id.in_(document_ids), + models.Document.deleted_at.is_(None), ) result = await db.execute(stmt) return result.scalars().all() @@ -989,6 +992,7 @@ async def get_child_observations( stmt = select(models.Document).where( models.Document.workspace_name == workspace_name, models.Document.source_ids.contains([parent_id]), + models.Document.deleted_at.is_(None), ) if observer: stmt = stmt.where(models.Document.observer == observer) diff --git a/src/models.py b/src/models.py index 040ee517..12d342a3 100644 --- a/src/models.py +++ b/src/models.py @@ -319,6 +319,12 @@ class MessageEmbedding(Base): postgresql_with={"m": 16, "ef_construction": 64}, postgresql_ops={"embedding": "vector_cosine_ops"}, ), + # Composite index for efficient reconciliation queries + Index( + "ix_message_embeddings_sync_state_last_sync_at", + "sync_state", + "last_sync_at", + ), ) @@ -455,6 +461,12 @@ class Document(Base): "source_ids", postgresql_using="gin", ), + # Composite index for efficient reconciliation queries + Index( + "ix_documents_sync_state_last_sync_at", + "sync_state", + "last_sync_at", + ), ) @@ -497,6 +509,20 @@ class QueueItem(Base): "processed", "id", ), + # Partial unique index for reconciler task deduplication + Index( + "uq_queue_reconciler_pending_work_unit_key", + "work_unit_key", + unique=True, + postgresql_where=text("task_type = 'reconciler' AND processed = false"), + ), + # Partial unique index for dream task deduplication + Index( + "uq_queue_dream_pending_work_unit_key", + "work_unit_key", + unique=True, + postgresql_where=text("task_type = 'dream' AND processed = false"), + ), ) def __repr__(self) -> str: diff --git a/src/reconciler/sync_vectors.py b/src/reconciler/sync_vectors.py index 2a3c514d..89a3008d 100644 --- a/src/reconciler/sync_vectors.py +++ b/src/reconciler/sync_vectors.py @@ -459,12 +459,113 @@ async def _cleanup_soft_deleted_documents_pgvector( return len(doc_ids) +async def _reconcile_documents_batch( + external_vector_store: VectorStore, + metrics: ReconciliationMetrics, +) -> bool: + """ + Reconcile a single batch of documents. + + Returns True if work was done, False otherwise. + """ + async with tracked_db("reconciliation_docs") as db: + docs = await _get_documents_needing_sync(db) + if not docs: + return False + + synced, failed = await _sync_documents(db, docs, external_vector_store) + metrics.documents_synced += synced + metrics.documents_failed += failed + await db.commit() + return True + + +async def _reconcile_message_embeddings_batch( + external_vector_store: VectorStore, + metrics: ReconciliationMetrics, +) -> bool: + """ + Reconcile a single batch of message embeddings. + + Returns True if work was done, False otherwise. + """ + async with tracked_db("reconciliation_embs") as db: + embs = await _get_message_embeddings_needing_sync(db) + if not embs: + return False + + try: + synced, failed = await _sync_message_embeddings( + db, embs, external_vector_store + ) + except Exception: + logger.exception( + "Message embedding reconciliation failed for %s embeddings", + len(embs), + ) + await _bump_message_embedding_sync_attempts(db, embs) + synced = 0 + failed = len(embs) + + metrics.message_embeddings_synced += synced + metrics.message_embeddings_failed += failed + await db.commit() + return True + + +async def _cleanup_documents_batch( + external_vector_store: VectorStore, + metrics: ReconciliationMetrics, +) -> bool: + """ + Clean up a single batch of soft-deleted documents. + + Returns True if work was done, False otherwise. + """ + from src.crud.document import cleanup_soft_deleted_documents + + async with tracked_db("reconciliation_cleanup") as db: + cleaned = await cleanup_soft_deleted_documents( + db, + external_vector_store, + batch_size=RECONCILIATION_BATCH_SIZE, + ) + if not cleaned: + return False + + metrics.documents_cleaned += cleaned + await db.commit() + return True + + +async def _cleanup_pgvector_batch( + metrics: ReconciliationMetrics, +) -> bool: + """ + Clean up a single batch of soft-deleted documents in pgvector-only mode. + + Returns True if work was done, False otherwise. + """ + async with tracked_db("reconciliation_pgvector_cleanup") as db: + cleaned = await _cleanup_soft_deleted_documents_pgvector( + db, batch_size=RECONCILIATION_BATCH_SIZE + ) + if not cleaned: + return False + + metrics.documents_cleaned += cleaned + await db.commit() + return True + + async def run_vector_reconciliation_cycle() -> ReconciliationMetrics: """ Run a complete reconciliation cycle. Runs a rolling sweep to reconcile missing vectors and clean up soft deletes. Uses batching and FOR UPDATE SKIP LOCKED for safe concurrent operation. + Each batch operation uses its own database session to avoid holding + connections open for the entire cycle duration. Returns metrics about what was synced. """ @@ -472,74 +573,38 @@ async def run_vector_reconciliation_cycle() -> ReconciliationMetrics: external_vector_store = get_external_vector_store() deadline = time.monotonic() + RECONCILIATION_TIME_BUDGET_SECONDS - from src.crud.document import cleanup_soft_deleted_documents - - async with tracked_db("reconciliation") as db: - # If no external vector store (pgvector mode), only clean up soft-deleted documents - if external_vector_store is None: - while time.monotonic() < deadline: - cleaned = await _cleanup_soft_deleted_documents_pgvector( - db, batch_size=RECONCILIATION_BATCH_SIZE - ) - if cleaned: - metrics.documents_cleaned += cleaned - await db.commit() - else: - break - return metrics - + # If no external vector store (pgvector mode), only clean up soft-deleted documents + if external_vector_store is None: while time.monotonic() < deadline: - did_work = False - - # Reconcile documents - docs = await _get_documents_needing_sync(db) - if docs: - synced, failed = await _sync_documents(db, docs, external_vector_store) - metrics.documents_synced += synced - metrics.documents_failed += failed - await db.commit() - did_work = True - - if time.monotonic() >= deadline: - break - - # Reconcile message embeddings - embs = await _get_message_embeddings_needing_sync(db) - if embs: - try: - synced, failed = await _sync_message_embeddings( - db, embs, external_vector_store - ) - except Exception: - logger.exception( - "Message embedding reconciliation failed for %s embeddings", - len(embs), - ) - await _bump_message_embedding_sync_attempts(db, embs) - synced = 0 - failed = len(embs) - metrics.message_embeddings_synced += synced - metrics.message_embeddings_failed += failed - await db.commit() - did_work = True - - if time.monotonic() >= deadline: - break - - # Clean up soft-deleted documents - cleaned = await cleanup_soft_deleted_documents( - db, - external_vector_store, - batch_size=RECONCILIATION_BATCH_SIZE, - ) - if cleaned: - metrics.documents_cleaned += cleaned - await db.commit() - did_work = True - + did_work = await _cleanup_pgvector_batch(metrics) if not did_work: - logger.debug("No work done, breaking reconciliation loop") break - logger.info("Vector reconciliation cycle completed") + logger.info("Vector reconciliation cycle completed (pgvector mode)") + return metrics + # External vector store mode - reconcile documents, embeddings, and cleanup + while time.monotonic() < deadline: + # Reconcile documents + docs_work = await _reconcile_documents_batch(external_vector_store, metrics) + + if time.monotonic() >= deadline: + break + + # Reconcile message embeddings + embs_work = await _reconcile_message_embeddings_batch( + external_vector_store, metrics + ) + + if time.monotonic() >= deadline: + break + + # Clean up soft-deleted documents + cleanup_work = await _cleanup_documents_batch(external_vector_store, metrics) + + # Continue only if any operation did work + if not (docs_work or embs_work or cleanup_work): + logger.debug("No work done, breaking reconciliation loop") + break + + logger.info("Vector reconciliation cycle completed") return metrics diff --git a/tests/alembic/revisions/test_119a52b73c60_support_external_embeddings.py b/tests/alembic/revisions/test_119a52b73c60_support_external_embeddings.py index d9e9bfe5..65ffb115 100644 --- a/tests/alembic/revisions/test_119a52b73c60_support_external_embeddings.py +++ b/tests/alembic/revisions/test_119a52b73c60_support_external_embeddings.py @@ -12,7 +12,7 @@ INDEXES = ( ("documents", "ix_documents_sync_state_last_sync_at"), ("message_embeddings", "ix_message_embeddings_sync_state"), ("message_embeddings", "ix_message_embeddings_sync_state_last_sync_at"), - ("queue", "uq_queue_work_unit_key"), + ("queue", "uq_queue_reconciler_pending_work_unit_key"), ) diff --git a/tests/alembic/revisions/test_7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py b/tests/alembic/revisions/test_7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py index 306f40e4..4028f242 100644 --- a/tests/alembic/revisions/test_7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py +++ b/tests/alembic/revisions/test_7c0d9a4e3b1f_add_unique_index_for_pending_dreams.py @@ -6,7 +6,7 @@ from tests.alembic.registry import register_after_upgrade, register_before_upgra from tests.alembic.verifier import MigrationVerifier # The partial unique index created by this migration -INDEX = ("queue", "ux_queue_dream_pending_work_unit_key") +INDEX = ("queue", "uq_queue_dream_pending_work_unit_key") @register_before_upgrade("7c0d9a4e3b1f")