test: verify outcome migrations and worker recovery on PostgreSQL
Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929
This commit is contained in:
parent
f0eff0ac92
commit
7f0dc78607
6 changed files with 315 additions and 8 deletions
258
tests/test_postgres_integration.py
Normal file
258
tests/test_postgres_integration.py
Normal file
|
|
@ -0,0 +1,258 @@
|
|||
"""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())
|
||||
Loading…
Add table
Add a link
Reference in a new issue