* fix(deriver): eliminate create_documents deadlock and stop silently burning batches on transient errors Two concurrent work units writing the same (workspace, observer, observed) collection deadlocked on times_derived reinforcement UPDATEs issued in batch order (DEV-1975, 682 events in 90 days). The deadlock was swallowed per-document, the loop cascaded PendingRollbackErrors against the dead session, the whole batch was lost, and the queue item was marked processed. - serialize writers per collection with a transaction-scoped advisory lock (pg_advisory_xact_lock + SET LOCAL lock_timeout), skipped for insert-only batches; covers all three row-lock sites in one move - hoist external-vector-store dup-candidate resolution ahead of the first DB statement so the lock's critical section contains no network calls - abort the batch on SQLAlchemyError instead of continuing through an aborted transaction; per-document skip semantics kept for non-DB errors - classify transient errors (new src/utils/retryable_errors.py) and retry them via a bounded in-process counter instead of marking items errored * fix(deriver): replace create_documents advisory lock with id-ordered row locks Advisory locks are database-scoped and would serialize every writer to a collection, including across Groudon tenants that share names. Collect reinforcement and replace ops during the loop, lock target rows with SELECT ... ORDER BY id FOR UPDATE, then apply. populate_existing reloads times_derived so a prefetched identity-map row cannot lose a concurrent increment. * fix(deriver): harden create_documents candidate hoist and test isolation Skip empty embeddings on the external-store path, isolate per-document resolve failures, and keep replacement times_derived in the in-batch ledger. Patch get_external_vector_store in the hoist test and cover in-loop SQLAlchemyError abort. * fix(deriver): address CodeRabbit findings on create_documents deadlock fix - Distinguish external resolve failure ([] skip) from pgvector fallback (None) so _semantic_dup_decision never re-enters external I/O under an open session - Bound external candidate hoist concurrency with a semaphore - Map in-loop IntegrityError to ValidationException for a uniform contract - Persist transient retry attempts on the oldest unprocessed queue item so every deriver instance shares one MAX_RETRYABLE_ATTEMPTS budget - Cover resolve-failure skip and multi-manager reclaim of the retry budget * fix(deriver): harden retry metadata cleanup and stale reinforce fallback - Strip _retry_attempts from payloads in the same transaction as mark_queue_items_as_processed / mark_queue_item_as_errored - Clear shared retry metadata only after a successful terminal mark - On reinforce, if the locked target is gone or soft-deleted, insert the incoming document instead of dropping it - Skip pgvector semantic lookup when embedding is empty so query_documents cannot embed under an open session * fix(deriver): address review on deadlock retry and row-lock apply Strip _retry_attempts before payload validation so non-representation tasks are not burned as extra_forbidden. Re-raise retryable observer save errors after telemetry so the queue actually retries. Skip same-batch reinforce fallbacks after a replace. Revert unordered FOR UPDATE on mark processed/errored and drop post-commit retry cleanup from the success path. * fix: add test and simplify queue query --------- Co-authored-by: Vineeth Voruganti <13438633+VVoruganti@users.noreply.github.com> |
||
|---|---|---|
| .. | ||
| README.md | ||
| __init__.py | ||
| conftest.py | ||
| test_deriver_processing.py | ||
| test_embed_now.py | ||
| test_enqueue_dream.py | ||
| test_prompts.py | ||
| test_queue_operations.py | ||
| test_queue_processing.py | ||
| test_representation_crud.py | ||
| test_scope_backfill.py | ||
| test_vector_reconciliation.py | ||
README.md
Deriver Testing
This directory contains tests for the deriver system, which handles background processing of messages to extract insights and update working representations.
Structure
conftest.py- Shared fixtures for deriver testingtest_queue_operations.py- Tests for basic queue operationstest_deriver_processing.py- Tests for deriver processing logictest_queue_processing.py- Tests for queue manager and work unit processing
Key Fixtures
Database Fixtures
sample_session_with_peers- Creates a session with multiple peers having different observation configurationssample_messages- Creates sample messages for testingsample_queue_items- Creates queue items with various payload types (representation, summary)
Queue Fixtures
create_queue_payload- Helper to create queue payloads for testingadd_queue_items- Helper to add queue items to the databasecreate_active_queue_session- Helper to create active queue sessions for work unit tracking
Mocking Fixtures
mock_critical_analysis_call- Mocks the critical analysis LLM callmock_queue_manager- Mocks the queue manager for testingmock_representation_manager- Mocks the representation manager operations
Testing Patterns
Creating Queue Items
# Create representation payloads
payload = create_queue_payload(
message=message,
task_type="representation",
observer=observer_peer.name,
observed=message.peer_name
)
# Add to queue
queue_items = await add_queue_items([payload], session.id)
Testing Work Units
# Create a work unit
work_unit = WorkUnit(
session_id=session.id,
task_type="representation",
observer=observer,
observed=observed
)
# Test string representation
assert str(work_unit) == f"({session.id}, {observed.name}, {observer.name}, representation)"