d4b53bfb08
Docker / build (push) Has been cancelled
Test / pytest (ubuntu-latest, 3.12) (push) Has been cancelled
Test / pytest (windows-latest, 3.11) (push) Has been cancelled
Test / pytest (windows-latest, 3.12) (push) Has been cancelled
Lint / ruff (push) Has been cancelled
Test / pytest (ubuntu-latest, 3.11) (push) Has been cancelled
Build / build (push) Has been cancelled
128 lines
5.7 KiB
Python
128 lines
5.7 KiB
Python
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',)] |