hub-core/tests/test_postgres_integration.py

313 lines
13 KiB
Python
Raw Normal View History

"""Opt-in, real PostgreSQL tests. Always creates its own disposable container.
Run make postgres-test. No caller-provided database URL is accepted.
"""
import asyncio
import json
import os
from pathlib import Path
import secrets
import subprocess
import sys
import time
from uuid import uuid4
import pytest
import sqlalchemy as sa
from sqlalchemy.engine import URL
if os.getenv('HUB_CORE_TEST_POSTGRES') != '1':
pytest.skip('run make postgres-test for disposable PostgreSQL', allow_module_level=True)
from alembic import command
from alembic.config import Config
from hub_core.runtime.models import MessageCommand
from hub_core.runtime.postgres_store import PostgresPortStore
from hub_core.runtime.tables import runtime_messages, runtime_audit_ledger, runtime_outcome_outbox
from hub_core.security.context import current_authorization
from test_access_boundary import Owners
IMAGE = 'postgres@sha256:57c72fd2a128e416c7fcc499958864df5301e940bca0a56f58fddf30ffc07777'
HEAD = '0006_outcome_outbox'
def docker(*args):
return subprocess.run(['docker',*args], check=True, capture_output=True, text=True, timeout=120).stdout.strip()
@pytest.fixture(scope='module')
def postgres():
password = secrets.token_hex(16)
env = {**os.environ,'POSTGRES_PASSWORD':password}
process = subprocess.run(['docker','run','--detach','--rm',
'--label','hub-core.test=postgres-integration','--tmpfs','/var/lib/postgresql/data:rw',
'--publish','127.0.0.1::5432','--env','POSTGRES_PASSWORD',IMAGE],
env=env,check=True,capture_output=True,text=True,timeout=120)
container = process.stdout.strip()
try:
port = int(docker('port',container,'5432/tcp').rsplit(':',1)[1])
deadline = time.monotonic()+40
while True:
result = subprocess.run(['docker','exec',container,'pg_isready','-h','127.0.0.1','-U','postgres'],
capture_output=True,timeout=5)
if result.returncode == 0:
break
if time.monotonic() >= deadline:
pytest.fail('disposable PostgreSQL did not become ready')
time.sleep(.2)
yield URL.create('postgresql+psycopg2',username='postgres',password=password,
host='127.0.0.1',port=port,database='postgres')
finally:
docker('rm','--force',container)
@pytest.fixture
def database(postgres,monkeypatch):
name = 'hub_test_' + uuid4().hex
admin = sa.create_engine(postgres,isolation_level='AUTOCOMMIT')
with admin.connect() as connection:
connection.exec_driver_sql('CREATE DATABASE ' + name)
url = postgres.set(database=name)
# Alembic's environment can override config. Bind it explicitly to this owned DB.
monkeypatch.setenv('DATABASE_URL',url.render_as_string(hide_password=False))
monkeypatch.delenv('HUB_CORE_MIGRATION_ROLE',raising=False)
monkeypatch.delenv('HUB_CORE_MIGRATION_SCHEMA',raising=False)
config = Config()
config.set_main_option('script_location',str(Path('hub_core/migrations').resolve()))
config.set_main_option('sqlalchemy.url',url.render_as_string(hide_password=False))
try:
command.upgrade(config,HEAD)
yield url, config
finally:
with admin.connect() as connection:
connection.exec_driver_sql('DROP DATABASE ' + name + ' WITH (FORCE)')
admin.dispose()
def store_for(url):
from sqlalchemy.ext.asyncio import create_async_engine
return PostgresPortStore(create_async_engine(url.set(drivername='postgresql+asyncpg')))
async def commit_message(store):
owners = Owners()
context = await owners.controller().authorize('verified-root','hub.message.write',
'/ports/messaging/messages',str(uuid4()),'a'*64)
binding = current_authorization.set(context)
try:
return await store.send_message(MessageCommand(schema_version='0.1.0',correlation_id=uuid4(),
from_address='agent:root',to_addresses=['agent:reader'],body='disposable fixture'))
finally:
current_authorization.reset(binding)
async def pending(store):
async with store.sessions() as session:
return [dict(r) for r in (await session.execute(sa.select(runtime_outcome_outbox))).mappings()]
def test_full_migration_chain_and_guarded_downgrade(database):
url,config = database
engine = sa.create_engine(url)
try:
with engine.connect() as connection:
assert connection.exec_driver_sql('SELECT version_num FROM alembic_version').scalar_one() == HEAD
assert sa.inspect(connection).get_foreign_keys('runtime_outcome_outbox')[0]['referred_table'] == 'runtime_audit_ledger'
async def write():
store = store_for(url)
try:
assert await store.readiness_checks() == {'database':'ok'}
await commit_message(store)
finally:
await store.aclose()
asyncio.run(write())
with pytest.raises(RuntimeError,match='drained outbox'):
command.downgrade(config,'0005_message_identity_aliases')
with engine.begin() as connection:
# Synthetic custody receipt permits testing the schema rollback path.
connection.execute(runtime_outcome_outbox.update().values(delivered_at=time.time()))
command.downgrade(config,'0005_message_identity_aliases')
command.upgrade(config,HEAD)
with engine.connect() as connection:
assert connection.execute(sa.select(sa.func.count()).select_from(runtime_messages)).scalar_one() == 1
assert connection.execute(sa.select(sa.func.count()).select_from(runtime_audit_ledger)).scalar_one() == 1
finally:
engine.dispose()
def test_two_workers_skip_locked_row_without_duplicate_claim(database):
url,_ = database
async def run():
first,second = store_for(url),store_for(url)
held,release = asyncio.Event(),asyncio.Event()
seen = []
class HoldingSink:
async def append_outcome(self,event):
seen.append(event['id'])
held.set()
await release.wait()
class OtherSink:
async def append_outcome(self,event):
seen.append(event['id'])
task = None
try:
await commit_message(first)
await commit_message(first)
task = asyncio.create_task(first.deliver_outcomes(HoldingSink(),limit=1))
await asyncio.wait_for(held.wait(),2)
# Must finish while first worker still holds its transaction's row lock.
assert await asyncio.wait_for(second.deliver_outcomes(OtherSink(),limit=1),2) == 1
assert not task.done()
release.set()
assert await task == 1
assert len(seen) == len(set(seen)) == 2
assert all(row['delivered_at'] is not None for row in await pending(first))
finally:
release.set()
if task is not None and not task.done():
task.cancel()
await asyncio.gather(task,return_exceptions=True)
await first.aclose()
await second.aclose()
asyncio.run(run())
def test_outbox_failure_rolls_back_real_postgres_transaction(database):
url,_ = database
async def run():
store = store_for(url)
def fail(connection,cursor,statement,parameters,context,many):
if statement.startswith('INSERT INTO runtime_outcome_outbox'):
raise RuntimeError('injected outbox insert failure')
sa.event.listen(store.engine.sync_engine,'before_cursor_execute',fail)
try:
with pytest.raises(RuntimeError,match='injected outbox'):
await commit_message(store)
async with store.sessions() as session:
for table in (runtime_messages,runtime_audit_ledger,runtime_outcome_outbox):
assert (await session.execute(sa.select(sa.func.count()).select_from(table))).scalar_one() == 0
finally:
await store.aclose()
asyncio.run(run())
def test_worker_process_exit_after_receipt_replays_same_envelope(database,tmp_path):
url,_ = database
async def prepare():
store = store_for(url)
try:
await commit_message(store)
finally:
await store.aclose()
asyncio.run(prepare())
receipt = tmp_path/'accepted.json'
script = '''
import asyncio,json,os
from hub_core.runtime.postgres_store import PostgresPortStore
class Receiver:
async def append_outcome(self,event):
with open(os.environ['HUB_TEST_RECEIPT'],'w') as stream:
json.dump(event,stream)
stream.flush()
os.fsync(stream.fileno())
os._exit(73) # Receiver accepted; local DB delivery mark never committed.
async def main():
store=PostgresPortStore.from_url(os.environ['HUB_TEST_DATABASE'])
await store.deliver_outcomes(Receiver(),limit=1)
asyncio.run(main())
'''
env = {**os.environ,'HUB_TEST_RECEIPT':str(receipt),
'HUB_TEST_DATABASE':url.set(drivername='postgresql+asyncpg').render_as_string(hide_password=False)}
result = subprocess.run([sys.executable,'-c',script],env=env,capture_output=True,timeout=15)
assert result.returncode == 73
accepted = json.loads(receipt.read_text())
async def recover():
store = store_for(url)
class DuplicateReceiver:
async def append_outcome(self,event):
assert event == accepted
try:
row, = await pending(store)
assert row['attempts'] == 0 and row['delivered_at'] is None
assert await asyncio.wait_for(store.deliver_outcomes(DuplicateReceiver(),limit=1),3) == 1
assert await store.deliver_outcomes(DuplicateReceiver()) == 0
finally:
await store.aclose()
asyncio.run(recover())
def test_receiver_failure_persists_backoff_across_connection_reopen(database):
url,_ = database
async def run():
store = store_for(url)
class FailingReceiver:
async def append_outcome(self,event):
raise OSError('fixture receiver unavailable')
try:
await commit_message(store)
assert await store.deliver_outcomes(FailingReceiver()) == 0
finally:
await store.aclose()
reopened = store_for(url)
try:
row, = await pending(reopened)
assert row['attempts'] == 1 and row['next_attempt'] >= row['created_at'] + 1
assert row['delivered_at'] is None
finally:
await reopened.aclose()
asyncio.run(run())
def test_compatibility_key_transaction_and_outcome_rollback(database):
import httpx
from hub_core.runtime.app import create_app
from hub_core.runtime.config import RuntimeSettings
from hub_core.runtime.tables import compat_api_consumers, compat_api_keys
from test_access_boundary import HEADERS
url,_ = database
async def run():
store = store_for(url)
owners = Owners()
groups = frozenset({'credentials'})
app = create_app(settings=RuntimeSettings(environment='test',access_mode='enforce',
backend='postgresql',v2_groups=groups,v2_write_groups=groups),
port_store=store,access_controller=owners.controller())
async def snapshot(table):
async with store.sessions() as session:
return [dict(row) for row in (await session.execute(sa.select(table))).mappings()]
def fail(connection,cursor,statement,parameters,context,many):
if statement.startswith('INSERT INTO runtime_outcome_outbox'):
raise RuntimeError('fixture outbox failure')
try:
async with httpx.AsyncClient(transport=httpx.ASGITransport(app=app,raise_app_exceptions=False),
base_url='http://test') as client:
created = await client.post('/api/v2/api-consumers',headers=HEADERS,json={'name':'Fixture'})
assert created.status_code == 201
path = '/api-consumers/'+created.json()['id']+'/api-keys'
before = await snapshot(compat_api_consumers)
sa.event.listen(store.engine.sync_engine,'before_cursor_execute',fail)
try:
failed = await client.post(path,headers=HEADERS,json={})
assert failed.status_code == 500
finally:
sa.event.remove(store.engine.sync_engine,'before_cursor_execute',fail)
assert await snapshot(compat_api_consumers) == before
assert not await snapshot(compat_api_keys)
assert len(await pending(store)) == 1
issued = await client.post(path,headers=HEADERS,json={})
assert issued.status_code == 201
assert len(await snapshot(compat_api_keys)) == 1
assert await snapshot(compat_api_consumers) != before
outcomes = await pending(store)
assert len(outcomes) == 2
assert {row['envelope']['data']['operation'] for row in outcomes} == {
'compat.consumer.created','compat.api_key.created'}
assert all(row['envelope']['data']['authorization']['subject'] == 'immutable-root'
for row in outcomes)
assert issued.json()['fullKey'] not in json.dumps(outcomes)
assert len(await snapshot(runtime_audit_ledger)) == 2
finally:
await store.aclose()
asyncio.run(run())