259 lines
10 KiB
Python
259 lines
10 KiB
Python
|
|
"""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())
|