import asyncio import pytest from tests.test_host_registry_control import _runtime, _input, _preparation, _admission, _Sink from EvoScientist.llm.contracts import WebHostContext from EvoScientist.llm.host_execution_registry import SQLiteHostRegistry from EvoScientist.llm.runtime import _construction_owner @pytest.mark.asyncio async def test_owner_observation(tmp_path, monkeypatch, caplog): runtime, authority = _runtime(tmp_path, monkeypatch) runtime.host_registry = SQLiteHostRegistry(tmp_path / 'host.db', host_id='h', boot_id='b') release = asyncio.Event() owners = [] class Client: async def aclose(self): await release.wait() raise ValueError('SECRET cleanup detail') def broken(*args): run = _construction_owner.get() owners.append(run) client = Client() run._owned_clients[id(client)] = client raise ValueError('construction') runtime.agent_factory = broken value = _input() sink = _Sink() quote = await runtime.prepare_model_run( _preparation(authority, value, title_policy='disabled'), value, WebHostContext('/tmp', '/tmp', object(), object(), runtime_event_sink=sink)) with pytest.raises(TimeoutError): await runtime.start_web_run(_admission(authority, quote)) run = owners[0] release.set() await asyncio.wait({run._terminal_task}, timeout=2) await asyncio.sleep(0) try: state = runtime.inspect(run.run_id, owner_epoch=1, boot_id='b') assert state['terminal_error_code'] == 'RUN_TERMINAL_CLEANUP_UNCONFIRMED' assert run._terminal_task._log_traceback is False assert state['status'] == 'unknown' assert state['resources_confirmed_exited'] is False assert 'SECRET' not in caplog.text assert 'RUN_TERMINAL_CLEANUP_UNCONFIRMED' in caplog.text assert not sink.events finally: run._terminal_task.exception() @pytest.mark.asyncio async def test_durable_recovery(tmp_path, monkeypatch): import hashlib import json import sqlite3 from dataclasses import asdict from EvoScientist.llm.contracts import EvoRuntimeError, canonical_json_v1 from tests.test_host_registry_control import live_run runtime, registry, run = await live_run(tmp_path, monkeypatch) run._agent_task.cancel() await asyncio.gather(run._agent_task, return_exceptions=True) sink_path = tmp_path / 'sink.db' class DurableSink: uncertain = True def __init__(self): with sqlite3.connect(sink_path) as db: db.execute('CREATE TABLE IF NOT EXISTS events (id TEXT PRIMARY KEY, digest TEXT, body TEXT)') async def commit(self, event): body = canonical_json_v1(asdict(event)) digest = hashlib.sha256(body).hexdigest() with sqlite3.connect(sink_path) as db: prior = db.execute('SELECT digest FROM events WHERE id=?', (event.event_id,)).fetchone() if prior and prior[0] != digest: return 'conflict' db.execute('INSERT OR IGNORE INTO events VALUES (?, ?, ?)', (event.event_id, digest, body.decode())) if self.uncertain: raise OSError('lost commit response') return 'duplicate' if prior else 'committed' async def confirm(self, event_id, payload_digest): if self.uncertain: raise OSError('confirm unavailable') with sqlite3.connect(sink_path) as db: row = db.execute('SELECT digest FROM events WHERE id=?', (event_id,)).fetchone() return 'absent' if row is None else 'committed' if row[0] == payload_digest else 'conflict' sink = DurableSink() object.__setattr__(run._host, 'runtime_event_sink', sink) with pytest.raises(EvoRuntimeError, match='EVENT_COMMIT_INDETERMINATE'): await run._terminal_locked('awaiting_input', checkpoint_details={'checkpoint_id': 'cp'}) await asyncio.sleep(0) assert registry.inspect(run.run_id)['status'] == 'unknown' intent = registry.terminal_intent(run.run_id) assert intent['cleanup_confirmed'] is True assert intent['phase'] == 'prepared' assert intent['event']['payload']['outcome'] == 'awaiting_input' restarted = SQLiteHostRegistry(registry.path, host_id='host', boot_id='new') # A new boot cannot create resource evidence for an old execution. with pytest.raises(EvoRuntimeError, match='EXECUTION_BOOT_MISMATCH'): restarted.prepare_terminal(run.run_id, event=intent['event']) registry.bind(execution_id='unconfirmed', grant_id='unconfirmed', digest='d', thread_id='other', turn_id='t') runtime.host_registry = restarted sink = DurableSink() sink.uncertain = False assert await runtime.recover_terminal(run.run_id, sink=sink) == 'awaiting_input' assert await runtime.recover_terminal(run.run_id, sink=sink) == 'awaiting_input' assert restarted.inspect(run.run_id)['status'] == 'awaiting_input' assert restarted.terminal_intent(run.run_id)['phase'] == 'registry_finished' assert await run.wait_stopped(timeout=0.1) == 'awaiting_input' with pytest.raises(EvoRuntimeError, match='TERMINAL_CLEANUP_EVIDENCE_REQUIRED'): await runtime.recover_terminal('unconfirmed', sink=sink) with sqlite3.connect(sink_path) as db: rows = db.execute('SELECT body FROM events').fetchall() assert len(rows) == 1 assert json.loads(rows[0][0]) == intent['event'] with sqlite3.connect(registry.path) as db: assert db.execute('SELECT checkpoint_id FROM pending_continuations WHERE execution_id=?', (run.run_id,)).fetchone()[0] == 'cp' assert db.execute('SELECT execution_id FROM active_claims').fetchall() == [('unconfirmed',)]