teleo-infrastructure/tests/test_v3_apply_catalog_serialization_postgres.py

467 lines
15 KiB
Python

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