From 5846f7166a7361525d298be2fe904a30007b89ef Mon Sep 17 00:00:00 2001 From: Rahul Johny Date: Wed, 30 Sep 2026 18:10:41 +0530 Subject: [PATCH] [MISC] Raise the rig Postgres lock limit and mock dispatch in pg_barrier enqueue tests The integration tier fails intermittently for two unrelated reasons. Backend: every xdist worker migrates its own test database in one transaction, holding a lock per table and constraint. The rig's Postgres runs with the default max_locks_per_transaction=64, so the shared lock table overflows at random ("out of shared memory") and errors hundreds of backend tests at setup. Start the container with 256. Workers: five TestPgBarrierEnqueue tests never mocked queue_backend.dispatch.dispatch. Since UN-4078 made the PG queue the only transport, enqueue really dispatches the headers, opening a connection from DB_* env (default host unstract-db) instead of the TEST_DB_* test database. Wrap them in the same patch their neighbours use. Also drop a duplicated _barrier_pg_decrement import. Co-Authored-By: Claude Opus 5.5 --- tests/rig/runtime.py | 8 +++- workers/tests/test_pg_barrier.py | 76 +++++++++++++++++--------------- 2 files changed, 47 insertions(+), 37 deletions(-) diff --git a/tests/rig/runtime.py b/tests/rig/runtime.py index 87941e11a1..52a2840c22 100644 --- a/tests/rig/runtime.py +++ b/tests/rig/runtime.py @@ -333,7 +333,13 @@ def up(self) -> PlatformEndpoints: # __post_init__ invariant (host without port) is self-cleaning; # otherwise a partial spec leaks the four containers we just started. try: - pg = PostgresContainer("pgvector/pgvector:pg15") + # Every xdist worker migrates its own test database in one + # transaction, each holding a lock per table and constraint. The + # default of 64 overflows the shared lock table intermittently + # ("out of shared memory") and errors the whole backend group. + pg = PostgresContainer("pgvector/pgvector:pg15").with_command( + "postgres -c max_locks_per_transaction=256" + ) pg.start() self._stack.append(pg) redis = RedisContainer("redis:7.2.3").start() diff --git a/workers/tests/test_pg_barrier.py b/workers/tests/test_pg_barrier.py index dcae190e94..6d22d604e5 100644 --- a/workers/tests/test_pg_barrier.py +++ b/workers/tests/test_pg_barrier.py @@ -27,7 +27,6 @@ _barrier_pg_decrement, _fire_barrier_callback, barrier_pg_abort, - _barrier_pg_decrement, claim_batch, run_batch_with_barrier, try_claim_orchestration, @@ -850,13 +849,14 @@ def test_enqueue_sets_expires_cap_and_fresh_progress(self, barrier_db, monkeypat # expired nor stale. (UN-3661) monkeypatch.setenv("WORKER_BARRIER_KEY_TTL_SECONDS", "600") task, _ = _mock_header_task() - PgBarrier().enqueue( - [task], - callback_task_name="cb", - callback_kwargs={"execution_id": "exec-SD"}, - callback_queue="general", - app_instance=None, - ) + with patch("queue_backend.dispatch.dispatch"): + PgBarrier().enqueue( + [task], + callback_task_name="cb", + callback_kwargs={"execution_id": "exec-SD"}, + callback_queue="general", + app_instance=None, + ) assert 590 <= _expires_in_seconds(barrier_db, "exec-SD") <= 600 # ~ttl cap assert _last_progress_age_seconds(barrier_db, "exec-SD") < 5 # fresh @@ -920,34 +920,37 @@ def test_upsert_overwrites_stale_state(self, barrier_db): " now() + interval '1h', now())" ) task, _ = _mock_header_task() - PgBarrier().enqueue( - [task, _mock_header_task()[0]], - callback_task_name="cb", - callback_kwargs={"execution_id": "exec-R"}, - callback_queue="general", - app_instance=None, - ) + with patch("queue_backend.dispatch.dispatch"): + PgBarrier().enqueue( + [task, _mock_header_task()[0]], + callback_task_name="cb", + callback_kwargs={"execution_id": "exec-R"}, + callback_queue="general", + app_instance=None, + ) assert _row(barrier_db, "exec-R") == (2, []) def test_enqueue_stamps_organization_id(self, barrier_db): # The whole reason the org column + migration exist (reaper recovery). - PgBarrier().enqueue( - [_mock_header_task()[0]], - callback_task_name="cb", - callback_kwargs={"execution_id": "exec-ORG", "organization_id": "org-42"}, - callback_queue="general", - app_instance=None, - ) + with patch("queue_backend.dispatch.dispatch"): + PgBarrier().enqueue( + [_mock_header_task()[0]], + callback_task_name="cb", + callback_kwargs={"execution_id": "exec-ORG", "organization_id": "org-42"}, + callback_queue="general", + app_instance=None, + ) assert _org(barrier_db, "exec-ORG") == "org-42" def test_enqueue_defaults_org_to_empty_when_absent(self, barrier_db): - PgBarrier().enqueue( - [_mock_header_task()[0]], - callback_task_name="cb", - callback_kwargs={"execution_id": "exec-NOORG"}, # no organization_id - callback_queue="general", - app_instance=None, - ) + with patch("queue_backend.dispatch.dispatch"): + PgBarrier().enqueue( + [_mock_header_task()[0]], + callback_task_name="cb", + callback_kwargs={"execution_id": "exec-NOORG"}, # no organization_id + callback_queue="general", + app_instance=None, + ) assert _org(barrier_db, "exec-NOORG") == "" def test_upsert_refreshes_org_on_reenqueue(self, barrier_db): @@ -960,13 +963,14 @@ def test_upsert_refreshes_org_on_reenqueue(self, barrier_db): "VALUES ('exec-REORG', 'old-org', 1, '[]'::jsonb, now(), " " now() + interval '1h', now())" ) - PgBarrier().enqueue( - [_mock_header_task()[0]], - callback_task_name="cb", - callback_kwargs={"execution_id": "exec-REORG", "organization_id": "new-org"}, - callback_queue="general", - app_instance=None, - ) + with patch("queue_backend.dispatch.dispatch"): + PgBarrier().enqueue( + [_mock_header_task()[0]], + callback_task_name="cb", + callback_kwargs={"execution_id": "exec-REORG", "organization_id": "new-org"}, + callback_queue="general", + app_instance=None, + ) assert _org(barrier_db, "exec-REORG") == "new-org" def test_mid_loop_dispatch_failure_deletes_row(self, barrier_db):