2026-09-28 14:00:26 +02:00
|
|
|
import asyncio
|
|
|
|
|
import json
|
|
|
|
|
import time
|
|
|
|
|
from dataclasses import replace
|
|
|
|
|
from uuid import uuid4
|
|
|
|
|
|
|
|
|
|
import httpx
|
|
|
|
|
import pytest
|
|
|
|
|
import sqlalchemy as sa
|
|
|
|
|
from fastapi.testclient import TestClient
|
|
|
|
|
|
|
|
|
|
from hub_core.runtime.app import create_app
|
|
|
|
|
from hub_core.runtime.config import RuntimeSettings
|
|
|
|
|
from hub_core.runtime.postgres_store import PostgresPortStore
|
|
|
|
|
from hub_core.runtime.tables import runtime_metadata, runtime_messages, runtime_audit_ledger, runtime_outcome_outbox
|
|
|
|
|
from hub_core.security.audit import AuditCoreSink
|
|
|
|
|
from hub_core.security.context import current_authorization
|
|
|
|
|
from test_access_boundary import Owners, HEADERS
|
|
|
|
|
from test_access_audit import READY
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@pytest.fixture
|
|
|
|
|
def target(tmp_path):
|
|
|
|
|
url = f'sqlite+aiosqlite:///{tmp_path / "outcomes.db"}'
|
|
|
|
|
store = PostgresPortStore.from_url(url)
|
|
|
|
|
async def schema():
|
|
|
|
|
async with store.engine.begin() as connection:
|
|
|
|
|
await connection.run_sync(runtime_metadata.create_all)
|
|
|
|
|
asyncio.run(schema())
|
|
|
|
|
owners = Owners()
|
|
|
|
|
async def facts(actor, resource):
|
|
|
|
|
return replace(owners.facts, checked_at=time.time(), subject=actor.subject)
|
|
|
|
|
owners.resolve = facts
|
|
|
|
|
app = create_app(settings=RuntimeSettings(environment='test',access_mode='enforce',backend='postgresql',database_url=url),
|
|
|
|
|
port_store=store,access_controller=owners.controller())
|
|
|
|
|
with TestClient(app,raise_server_exceptions=False) as client:
|
|
|
|
|
yield client, store, owners, url
|
|
|
|
|
asyncio.run(store.aclose())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def message():
|
|
|
|
|
return {'schema_version':'0.1.0','correlation_id':str(uuid4()),'from_address':'agent:root',
|
|
|
|
|
'to_addresses':['agent:reader'],'body':'private body never archived'}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def rows(store, table):
|
|
|
|
|
async def read():
|
|
|
|
|
async with store.sessions() as session:
|
|
|
|
|
return [dict(r) for r in (await session.execute(sa.select(table))).mappings()]
|
|
|
|
|
return asyncio.run(read())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_committed_outcome_is_atomic_and_joins_authorization(target):
|
|
|
|
|
client,store,owners,_ = target
|
|
|
|
|
body = message()
|
|
|
|
|
response = client.post('/ports/messaging/messages',headers=HEADERS,json=body)
|
|
|
|
|
assert response.status_code == 202
|
|
|
|
|
ledger, = rows(store,runtime_audit_ledger)
|
|
|
|
|
pending, = rows(store,runtime_outcome_outbox)
|
|
|
|
|
envelope = pending['envelope']
|
|
|
|
|
authorization = envelope['data']['authorization']
|
|
|
|
|
assert envelope['id'] == ledger['id']
|
|
|
|
|
assert envelope['correlation_id'] == owners.records[0]['correlation_id'] == response.headers['x-correlation-id']
|
|
|
|
|
assert authorization == ledger['detail']['authorization']
|
|
|
|
|
assert authorization['decision_id'] == owners.records[0]['decision_id']
|
|
|
|
|
assert authorization['subject'] == 'immutable-root'
|
|
|
|
|
assert envelope['data']['business_correlation_id'] == body['correlation_id']
|
|
|
|
|
assert envelope['data']['subject_id'] == response.json()['id']
|
|
|
|
|
assert 'private body' not in json.dumps(envelope)
|
|
|
|
|
assert 'verified-root' not in json.dumps(envelope)
|
|
|
|
|
assert current_authorization.get() is None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_outbox_insert_failure_rolls_back_business_and_ledger(target):
|
|
|
|
|
client,store,owners,_ = target
|
|
|
|
|
def reject(connection,cursor,statement,parameters,context,many):
|
|
|
|
|
if statement.startswith('INSERT INTO runtime_outcome_outbox'):
|
|
|
|
|
raise RuntimeError('simulated storage failure')
|
|
|
|
|
sa.event.listen(store.engine.sync_engine,'before_cursor_execute',reject)
|
|
|
|
|
try:
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 500
|
|
|
|
|
finally:
|
|
|
|
|
sa.event.remove(store.engine.sync_engine,'before_cursor_execute',reject)
|
|
|
|
|
assert rows(store,runtime_messages) == rows(store,runtime_audit_ledger) == rows(store,runtime_outcome_outbox) == []
|
|
|
|
|
assert owners.records[0]['outcome'] == 'authorized' # Never claims commit.
|
|
|
|
|
assert current_authorization.get() is None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_denial_and_invalid_body_never_queue_committed_outcome(target):
|
|
|
|
|
client,store,owners,_ = target
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json={}).status_code == 422
|
|
|
|
|
owners.allow = False
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 403
|
|
|
|
|
assert not rows(store,runtime_outcome_outbox)
|
|
|
|
|
assert not rows(store,runtime_messages)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_lost_receipt_retry_reopen_and_duplicate_delivery(target,tmp_path):
|
|
|
|
|
client,store,owners,url = target
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202
|
|
|
|
|
credential = tmp_path/'sender'
|
|
|
|
|
credential.write_text('fixture-only')
|
|
|
|
|
archived, requests = {}, []
|
|
|
|
|
def receiver(request):
|
|
|
|
|
if request.url.path == '/readyz':
|
|
|
|
|
return httpx.Response(200,json=READY)
|
|
|
|
|
envelope = json.loads(request.content)
|
|
|
|
|
requests.append(bytes(request.content))
|
|
|
|
|
assert request.headers['idempotency-key'] == envelope['id']
|
|
|
|
|
if envelope['id'] not in archived:
|
|
|
|
|
archived[envelope['id']] = envelope
|
|
|
|
|
raise httpx.ReadTimeout('lost receipt after custody')
|
|
|
|
|
assert archived[envelope['id']] == envelope
|
|
|
|
|
return httpx.Response(200,json={'status':'duplicate','reference':'audit:'+envelope['id']})
|
|
|
|
|
async def run():
|
|
|
|
|
async with httpx.AsyncClient(transport=httpx.MockTransport(receiver)) as http:
|
|
|
|
|
sink = AuditCoreSink(base_url='https://audit.example',token_file=credential,client=http)
|
|
|
|
|
assert await store.deliver_outcomes(sink) == 0
|
|
|
|
|
assert await store.deliver_outcomes(sink) == 0 # Backoff, no immediate retry.
|
|
|
|
|
assert len(requests) == 1
|
|
|
|
|
await store.aclose()
|
|
|
|
|
reopened = PostgresPortStore.from_url(url)
|
|
|
|
|
try:
|
|
|
|
|
async with reopened.sessions.begin() as session:
|
|
|
|
|
await session.execute(runtime_outcome_outbox.update().values(next_attempt=0))
|
|
|
|
|
assert await reopened.deliver_outcomes(sink) == 1
|
|
|
|
|
assert await reopened.deliver_outcomes(sink) == 0
|
|
|
|
|
finally:
|
|
|
|
|
await reopened.aclose()
|
|
|
|
|
asyncio.run(run())
|
|
|
|
|
assert len(archived) == 1 and len(requests) == 2 and requests[0] == requests[1]
|
|
|
|
|
pending, = rows(store,runtime_outcome_outbox)
|
|
|
|
|
assert pending['attempts'] == 2 and pending['delivered_at'] is not None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_cancellation_after_receipt_leaves_replayable_row(target):
|
|
|
|
|
client,store,_,_ = target
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202
|
|
|
|
|
class Crash:
|
|
|
|
|
async def append_outcome(self,event):
|
|
|
|
|
raise asyncio.CancelledError()
|
|
|
|
|
async def run():
|
|
|
|
|
with pytest.raises(asyncio.CancelledError):
|
|
|
|
|
await store.deliver_outcomes(Crash())
|
|
|
|
|
asyncio.run(run())
|
|
|
|
|
pending, = rows(store,runtime_outcome_outbox)
|
|
|
|
|
assert pending['delivered_at'] is None and pending['attempts'] == 0
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_stale_backlog_is_visible_and_clears_after_receipt(target):
|
|
|
|
|
client,store,_,_ = target
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202
|
|
|
|
|
class Receipt:
|
|
|
|
|
async def append_outcome(self,event):
|
|
|
|
|
pass
|
|
|
|
|
async def run():
|
|
|
|
|
assert await store.outcome_readiness() == 'ok'
|
|
|
|
|
async with store.sessions.begin() as session:
|
|
|
|
|
await session.execute(runtime_outcome_outbox.update().values(created_at=time.time()-61))
|
|
|
|
|
assert await store.outcome_readiness() == 'stale'
|
|
|
|
|
await store.deliver_outcomes(Receipt())
|
|
|
|
|
assert await store.outcome_readiness() == 'ok'
|
|
|
|
|
asyncio.run(run())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_all_native_mutation_families_keep_attributed_outcomes(target):
|
|
|
|
|
from hub_core.conformance import ConformanceHarness
|
|
|
|
|
client,store,owners,_ = target
|
|
|
|
|
owners.facts = replace(owners.facts,producer_addresses=frozenset({'hub:ops-hub'}))
|
2026-09-28 14:32:36 +02:00
|
|
|
# The harness must recognize the missing dispatcher as a degraded dependency.
|
2026-09-28 14:00:26 +02:00
|
|
|
harness = ConformanceHarness(client)
|
|
|
|
|
client.headers.update(HEADERS)
|
|
|
|
|
report = harness.run()
|
2026-09-28 14:32:36 +02:00
|
|
|
assert report.passed, report.to_dict()
|
2026-09-28 14:00:26 +02:00
|
|
|
assert client.get('/readyz').json()['checks']['outcome_delivery'] == 'unavailable'
|
|
|
|
|
pending = rows(store,runtime_outcome_outbox)
|
|
|
|
|
operations = {r['envelope']['data']['operation'] for r in pending}
|
|
|
|
|
assert {'message.accepted','event.progress.accepted','event.interaction.accepted'} <= operations
|
|
|
|
|
assert any(op.startswith('registry.') for op in operations)
|
|
|
|
|
assert all(r['envelope']['data']['authorization']['subject'] == 'immutable-root' for r in pending)
|
|
|
|
|
assert len(pending) == len(rows(store,runtime_audit_ledger))
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_concurrent_requests_keep_separate_authorization_contexts(target):
|
|
|
|
|
from hub_core.security.identity import AccessFailure
|
|
|
|
|
client,store,owners,_ = target
|
|
|
|
|
async def authenticate(token):
|
|
|
|
|
if token not in {'workload:a','workload:b'}:
|
|
|
|
|
raise AccessFailure(401,'bad_fixture')
|
|
|
|
|
await asyncio.sleep(0)
|
|
|
|
|
return replace(owners.actor,subject=token,principal_type='service')
|
|
|
|
|
owners.authenticate = authenticate
|
|
|
|
|
async def run():
|
|
|
|
|
async with httpx.AsyncClient(transport=httpx.ASGITransport(app=client.app),base_url='https://hub.example') as http:
|
|
|
|
|
async def send(subject):
|
|
|
|
|
result = await http.post('/ports/messaging/messages',headers={'Authorization':'Bearer '+subject},json=message())
|
|
|
|
|
assert result.status_code == 202
|
|
|
|
|
return result.json()['id'], subject
|
|
|
|
|
return dict(await asyncio.gather(send('workload:a'),send('workload:b')))
|
|
|
|
|
subjects = asyncio.run(run())
|
|
|
|
|
for row in rows(store,runtime_outcome_outbox):
|
|
|
|
|
data = row['envelope']['data']
|
|
|
|
|
assert data['authorization']['subject'] == subjects[data['subject_id']]
|
|
|
|
|
assert current_authorization.get() is None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_outcome_migration_creates_index_and_foreign_key(tmp_path):
|
|
|
|
|
import importlib
|
|
|
|
|
from alembic.migration import MigrationContext
|
|
|
|
|
from alembic.operations import Operations
|
|
|
|
|
migration = importlib.import_module('hub_core.migrations.versions.0006_outcome_outbox')
|
|
|
|
|
engine = sa.create_engine('sqlite:///' + str(tmp_path/'migration.db'))
|
|
|
|
|
try:
|
|
|
|
|
runtime_audit_ledger.create(engine)
|
|
|
|
|
with engine.begin() as connection:
|
|
|
|
|
with Operations.context(MigrationContext.configure(connection)):
|
|
|
|
|
migration.upgrade()
|
|
|
|
|
inspect = sa.inspect(connection)
|
|
|
|
|
assert {x['name'] for x in inspect.get_columns('runtime_outcome_outbox')} == set(runtime_outcome_outbox.c.keys())
|
|
|
|
|
assert inspect.get_foreign_keys('runtime_outcome_outbox')[0]['referred_table'] == 'runtime_audit_ledger'
|
|
|
|
|
assert inspect.get_indexes('runtime_outcome_outbox')[0]['name'] == 'ix_runtime_outcome_pending'
|
|
|
|
|
connection.execute(runtime_outcome_outbox.insert().values(id='fixture',envelope={},created_at=0,attempts=0,next_attempt=0))
|
|
|
|
|
with Operations.context(MigrationContext.configure(connection)):
|
|
|
|
|
with pytest.raises(RuntimeError,match='drained outbox'):
|
|
|
|
|
migration.downgrade()
|
|
|
|
|
connection.execute(runtime_outcome_outbox.delete())
|
|
|
|
|
with Operations.context(MigrationContext.configure(connection)):
|
|
|
|
|
migration.downgrade()
|
|
|
|
|
assert 'runtime_outcome_outbox' not in sa.inspect(connection).get_table_names()
|
|
|
|
|
finally:
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_runtime_dispatcher_drains_pending_outcome_and_stops(target):
|
|
|
|
|
from hub_core.runtime.app import _deliver_outcomes
|
|
|
|
|
client,store,_,_ = target
|
|
|
|
|
assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202
|
|
|
|
|
async def run():
|
|
|
|
|
class Receipt:
|
|
|
|
|
async def append_outcome(self,envelope):
|
|
|
|
|
pass
|
|
|
|
|
task = asyncio.create_task(_deliver_outcomes(store,Receipt()))
|
|
|
|
|
try:
|
|
|
|
|
async with asyncio.timeout(3):
|
|
|
|
|
while True:
|
|
|
|
|
async with store.sessions() as session:
|
|
|
|
|
count = (await session.execute(sa.select(sa.func.count()).select_from(runtime_outcome_outbox)
|
|
|
|
|
.where(runtime_outcome_outbox.c.delivered_at.is_not(None)))).scalar_one()
|
|
|
|
|
if count == 1:
|
|
|
|
|
break
|
|
|
|
|
await asyncio.sleep(.01)
|
|
|
|
|
finally:
|
|
|
|
|
task.cancel()
|
|
|
|
|
with pytest.raises(asyncio.CancelledError):
|
|
|
|
|
await task
|
|
|
|
|
asyncio.run(run())
|