from __future__ import annotations import subprocess import time import pytest from ops import verify_teleo_v3_epistemic_contract as contract from tests import test_v3_apply_worker_postgres as worker CATALOG_DDL_ADVISORY_LOCK_KEY = 6072343533146485505 TEST_APPLY_BARRIER_LOCK_KEY = 6072343533146485599 def _start_psql( postgres: contract.DisposablePostgres, sql: str, *, role: str = "postgres", password: str | None = None, keep_stdin_open: bool = False, ) -> subprocess.Popen[str]: assert postgres.docker is not None command = [postgres.docker, "exec"] if password is not None: command.extend(["-e", f"PGPASSWORD={password}"]) command.extend( [ "-i", postgres.name, "psql", "--no-psqlrc", "--set=ON_ERROR_STOP=1", "--set=VERBOSITY=verbose", "--username", role, f"--dbname={contract.DATABASE}", ] ) if password is not None: command.extend(["--host=127.0.0.1"]) process = subprocess.Popen( command, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, ) assert process.stdin is not None process.stdin.write(sql) if keep_stdin_open: process.stdin.flush() else: process.stdin.close() process.stdin = None return process def _finish(process: subprocess.Popen[str], *, timeout: float = 20) -> subprocess.CompletedProcess[str]: stdout, stderr = process.communicate(timeout=timeout) return subprocess.CompletedProcess(process.args, process.returncode, stdout, stderr) def _stop(process: subprocess.Popen[str] | None) -> None: if process is None or process.poll() is not None: return process.terminate() try: process.communicate(timeout=3) except subprocess.TimeoutExpired: process.kill() process.communicate(timeout=3) def _wait_for_activity( postgres: contract.DisposablePostgres, application_name: str, *, wait_event_type: str, wait_event: str, timeout: float = 8, ) -> None: deadline = time.monotonic() + timeout while time.monotonic() < deadline: waiting = postgres.psql( f""" select exists ( select 1 from pg_catalog.pg_stat_activity where application_name = '{application_name}' and wait_event_type = '{wait_event_type}' and wait_event = '{wait_event}' ); """ ).stdout.strip() if waiting == "t": return time.sleep(0.05) raise AssertionError(f"{application_name} did not wait on {wait_event_type}/{wait_event}") def _approved_apply_sql( postgres: contract.DisposablePostgres, ids: dict[str, str], ) -> str: payload = worker._payload(ids) packet = worker._insert_pending_proposal(postgres, payload, ids) worker._approve_proposal(postgres, ids["proposal"], packet) proposal = postgres.psql_json( f""" select pg_catalog.jsonb_build_object( 'id', id::text, 'proposal_type', proposal_type, 'status', status, 'payload', payload, 'reviewed_by_handle', reviewed_by_handle, 'reviewed_by_agent_id', reviewed_by_agent_id::text, 'reviewed_at', pg_catalog.to_char( reviewed_at at time zone 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US"Z"' ), 'review_note', review_note, 'applied_by_handle', applied_by_handle, 'applied_by_agent_id', applied_by_agent_id::text, 'applied_at', applied_at )::text from kb_stage.kb_proposals where id = '{ids["proposal"]}'::uuid; """ ) return worker.apply.build_apply_sql(proposal, worker.apply.SERVICE_AGENT_HANDLE) def _instrument_apply_after_catalog_guard(sql: str, application_name: str) -> str: guard_end = "$v3_contract_guard$;" assert sql.count(guard_end) == 1 return sql.replace( "begin;\n", f"begin;\nset local application_name = '{application_name}';\n", 1, ).replace( guard_end, guard_end + f"\nselect pg_catalog.pg_advisory_xact_lock({TEST_APPLY_BARRIER_LOCK_KEY});", 1, ) def _instrument_apply_before_transaction(sql: str, application_name: str) -> str: return f"set application_name = '{application_name}';\n{sql}" def _start_barrier_holder(postgres: contract.DisposablePostgres) -> subprocess.Popen[str]: holder = _start_psql( postgres, f""" set application_name = 'teleo_v3_catalog_barrier_holder'; select pg_catalog.pg_advisory_lock({TEST_APPLY_BARRIER_LOCK_KEY}); """, keep_stdin_open=True, ) _wait_for_activity( postgres, "teleo_v3_catalog_barrier_holder", wait_event_type="Client", wait_event="ClientRead", ) return holder def _release_barrier( postgres: contract.DisposablePostgres, holder: subprocess.Popen[str], ) -> None: released = postgres.psql( """ select pg_catalog.pg_terminate_backend(pid) from pg_catalog.pg_stat_activity where application_name = 'teleo_v3_catalog_barrier_holder'; """ ) assert released.stdout.strip() == "t" if holder.stdin is not None: holder.stdin.close() holder.stdin = None _finish(holder, timeout=5) def _install_and_build_apply( postgres: contract.DisposablePostgres, ids: dict[str, str], ) -> str: environment = postgres.start() assert environment["network_mode"] == "none" postgres.psql(contract.BOOTSTRAP_SQL) postgres.psql(worker._supplemental_bootstrap()) postgres.apply_file(worker.PREREQUISITES) worker._install_v3_contract(postgres) return _approved_apply_sql(postgres, ids) @pytest.mark.skipif(not contract.docker_available(), reason="Docker daemon is required") def test_trigger_ddl_waits_for_apply_transaction_relation_locks() -> None: postgres = contract.DisposablePostgres() holder: subprocess.Popen[str] | None = None apply_process: subprocess.Popen[str] | None = None ddl_process: subprocess.Popen[str] | None = None try: apply_sql = _install_and_build_apply(postgres, worker.RACE_IDS) holder = _start_barrier_holder(postgres) apply_process = _start_psql( postgres, _instrument_apply_after_catalog_guard(apply_sql, "teleo_v3_guarded_apply"), role="kb_apply", password=worker.APPLY_PASSWORD, ) _wait_for_activity( postgres, "teleo_v3_guarded_apply", wait_event_type="Lock", wait_event="advisory", ) ddl_process = _start_psql( postgres, """ set application_name = 'teleo_v3_uncoordinated_trigger_ddl'; create or replace trigger teleo_v3_00_enforce_apply_contract_cutover before insert or update of proposal_type, payload, status on kb_stage.kb_proposals for each row when (false) execute function kb_stage.teleo_v3_enforce_apply_contract_cutover(); """, ) _wait_for_activity( postgres, "teleo_v3_uncoordinated_trigger_ddl", wait_event_type="Lock", wait_event="relation", ) assert ddl_process.poll() is None assert worker._contract_state(postgres) == "V3" _release_barrier(postgres, holder) holder = None apply_result = _finish(apply_process) apply_process = None assert apply_result.returncode == 0, apply_result.stderr ddl_result = _finish(ddl_process) ddl_process = None assert ddl_result.returncode == 0, ddl_result.stderr assert postgres.psql_json( f""" select pg_catalog.jsonb_build_object( 'proposal_status', ( select status from kb_stage.kb_proposals where id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'contract_state', ({worker.apply.build_v3_contract_state_query()}) )::text; """ ) == {"proposal_status": "applied", "contract_state": "INVALID"} finally: if holder is not None and holder.poll() is None: _release_barrier(postgres, holder) _stop(ddl_process) _stop(apply_process) _stop(holder) cleanup = postgres.cleanup() assert cleanup["container_absent"] is True assert cleanup["instance_label_readback_clear"] is True @pytest.mark.skipif(not contract.docker_available(), reason="Docker daemon is required") def test_controlled_function_migration_waits_for_apply_shared_advisory_lock() -> None: postgres = contract.DisposablePostgres() holder: subprocess.Popen[str] | None = None apply_process: subprocess.Popen[str] | None = None migration_process: subprocess.Popen[str] | None = None try: apply_sql = _install_and_build_apply(postgres, worker.RACE_IDS) holder = _start_barrier_holder(postgres) apply_process = _start_psql( postgres, _instrument_apply_after_catalog_guard(apply_sql, "teleo_v3_shared_lock_apply"), role="kb_apply", password=worker.APPLY_PASSWORD, ) _wait_for_activity( postgres, "teleo_v3_shared_lock_apply", wait_event_type="Lock", wait_event="advisory", ) migration_process = _start_psql( postgres, f""" begin; set local application_name = 'teleo_v3_controlled_function_migration'; select pg_catalog.pg_advisory_xact_lock({CATALOG_DDL_ADVISORY_LOCK_KEY}); create or replace function kb_stage.teleo_v3_validate_claim_edge() returns trigger language plpgsql set search_path = pg_catalog, pg_temp as $function$ begin return new; end $function$; commit; """, ) _wait_for_activity( postgres, "teleo_v3_controlled_function_migration", wait_event_type="Lock", wait_event="advisory", ) assert migration_process.poll() is None assert worker._contract_state(postgres) == "V3" _release_barrier(postgres, holder) holder = None apply_result = _finish(apply_process) apply_process = None assert apply_result.returncode == 0, apply_result.stderr migration_result = _finish(migration_process) migration_process = None assert migration_result.returncode == 0, migration_result.stderr assert worker._contract_state(postgres) == "INVALID" finally: if holder is not None and holder.poll() is None: _release_barrier(postgres, holder) _stop(migration_process) _stop(apply_process) _stop(holder) cleanup = postgres.cleanup() assert cleanup["container_absent"] is True assert cleanup["instance_label_readback_clear"] is True @pytest.mark.skipif(not contract.docker_available(), reason="Docker daemon is required") def test_migration_first_apply_resnapshots_and_refuses_post_migration_catalog_drift() -> None: postgres = contract.DisposablePostgres() holder: subprocess.Popen[str] | None = None apply_process: subprocess.Popen[str] | None = None migration_process: subprocess.Popen[str] | None = None try: apply_sql = _install_and_build_apply(postgres, worker.RACE_IDS) holder = _start_barrier_holder(postgres) migration_process = _start_psql( postgres, f""" begin; set local application_name = 'teleo_v3_migration_first_holder'; select pg_catalog.pg_advisory_xact_lock({CATALOG_DDL_ADVISORY_LOCK_KEY}); select pg_catalog.pg_advisory_xact_lock({TEST_APPLY_BARRIER_LOCK_KEY}); create or replace function kb_stage.teleo_v3_validate_claim_edge() returns trigger language plpgsql set search_path = pg_catalog, pg_temp as $function$ begin return new; end $function$; commit; """, ) _wait_for_activity( postgres, "teleo_v3_migration_first_holder", wait_event_type="Lock", wait_event="advisory", ) apply_process = _start_psql( postgres, _instrument_apply_before_transaction(apply_sql, "teleo_v3_migration_first_apply"), role="kb_apply", password=worker.APPLY_PASSWORD, ) _wait_for_activity( postgres, "teleo_v3_migration_first_apply", wait_event_type="Lock", wait_event="advisory", ) _release_barrier(postgres, holder) holder = None migration_result = _finish(migration_process) migration_process = None assert migration_result.returncode == 0, migration_result.stderr apply_result = _finish(apply_process) apply_process = None assert apply_result.returncode != 0 assert "V3 catalog contract is INVALID" in apply_result.stderr assert postgres.psql_json( f""" select pg_catalog.jsonb_build_object( 'proposal_status', ( select status from kb_stage.kb_proposals where id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'applied_at', ( select applied_at from kb_stage.kb_proposals where id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'claims', ( select count(*) from public.claims where accepted_by_proposal_id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'sources', ( select count(*) from public.sources where accepted_by_proposal_id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'evidence', ( select count(*) from public.claim_evidence where accepted_from_proposal_id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'edges', ( select count(*) from public.claim_edges where accepted_by_proposal_id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'assessments', ( select count(*) from public.claim_evidence_assessments where accepted_by_proposal_id = '{worker.RACE_IDS["proposal"]}'::uuid ), 'contract_state', ({worker.apply.build_v3_contract_state_query()}) )::text; """ ) == { "proposal_status": "approved", "applied_at": None, "claims": 0, "sources": 0, "evidence": 0, "edges": 0, "assessments": 0, "contract_state": "INVALID", } lock_cleanup = postgres.psql_json( f""" with acquired as ( select pg_catalog.pg_try_advisory_lock({CATALOG_DDL_ADVISORY_LOCK_KEY}) as value ), released as ( select pg_catalog.pg_advisory_unlock({CATALOG_DDL_ADVISORY_LOCK_KEY}) as value from acquired where value ) select pg_catalog.jsonb_build_object( 'exclusive_lock_acquired', (select value from acquired), 'exclusive_lock_released', coalesce((select value from released), false) )::text; """ ) assert lock_cleanup == { "exclusive_lock_acquired": True, "exclusive_lock_released": True, } finally: if holder is not None and holder.poll() is None: _release_barrier(postgres, holder) _stop(migration_process) _stop(apply_process) _stop(holder) cleanup = postgres.cleanup() assert cleanup["container_absent"] is True assert cleanup["instance_label_readback_clear"] is True