Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion tests/rig/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
76 changes: 40 additions & 36 deletions workers/tests/test_pg_barrier.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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):
Expand All @@ -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):
Expand Down
Loading