MetaMetaMeta / Uebergabe /test_operator_loop.py
smlflg's picture
Initial public upload from Projekte/MetaMetaMeta
3bc7cb3 verified
Raw
History Blame Contribute Delete
5.54 kB
import json
import unittest
import tempfile
from pathlib import Path
from unittest.mock import patch
import board_worker as w
import operator_loop as op
from hermes_cli import kanban_db as kb
import test_board_lifecycle as lifecycle_tests
class OperatorTests(unittest.TestCase):
setUp=lifecycle_tests.Lifecycle.setUp
tearDown=lifecycle_tests.Lifecycle.tearDown
authority=lifecycle_tests.Lifecycle.authority
def setup_card(self,grant=None):
self.manifest.update(operator_session='lease',routes={'wggesucht':{
'description':'WG-Gesucht specialist','grant':grant}})
with w.db('test') as c:
return kb.create_task(c,title='WG-Gesucht Browser fehlt',created_by='wggesucht')
def test_missing_grant_routes_but_never_starts(self):
tid=self.setup_card()
with patch.object(op,'choose',return_value={'profile':'wggesucht','reason':'Browser blocker'}),patch.object(w,'run_one') as run:
result=op.process_one(self.manifest)
self.assertEqual(result['status'],'blocked_no_grant')
run.assert_not_called()
self.assertIsNone(op.process_one(self.manifest))
with w.db('test') as c:
task=kb.get_task(c,tid)
self.assertEqual(task.assignee,'hai:wggesucht')
self.assertEqual(task.status,'blocked')
self.assertEqual(c.execute("SELECT count(*) FROM task_events WHERE kind='operator_routed'").fetchone()[0],1)
self.assertEqual(c.execute('SELECT count(*) FROM kanban_notify_subs').fetchone()[0],0)
def test_granted_specialist_is_bound_to_card_and_result_followed(self):
tid=self.setup_card(self.grant)
with patch.object(op,'choose',return_value={'profile':'wggesucht','reason':'Match'}),patch.object(w,'run_one',return_value={'task':tid,'status':'review','evidence':'result.json'}) as run:
result=op.process_one(self.manifest)
sent=run.call_args.args[0]['tasks'][tid]
self.assertEqual(sent['worker_profile'],'wggesucht')
self.assertIn('card_sha256',sent)
self.assertEqual(result['profile'],'wggesucht')
with w.db('test') as c:
self.assertIn('meldet review',c.execute('SELECT body FROM task_comments ORDER BY id DESC LIMIT 1').fetchone()[0])
def test_foreign_assignment_is_untouched(self):
tid=self.setup_card()
with w.db('test') as c: kb.assign_task(c,tid,'hai:foreign-worker')
with patch.object(op,'choose') as choose:
self.assertIsNone(op.process_one(self.manifest));choose.assert_not_called()
def test_invalid_operator_lease_prevents_any_assignment(self):
tid=self.setup_card()
with patch.object(w,'rpc',side_effect=RuntimeError('expired')):
with self.assertRaises(RuntimeError):op.process_one(self.manifest)
with w.db('test') as c:self.assertIsNone(kb.get_task(c,tid).assignee)
def test_routing_failure_is_a_visible_blocker(self):
tid=self.setup_card()
with patch.object(op,'choose',side_effect=RuntimeError('provider unavailable')):
with self.assertRaises(RuntimeError):op.process_one(self.manifest)
with w.db('test') as c:self.assertEqual(kb.get_task(c,tid).status,'blocked')
def test_worker_failure_is_followed_up_on_ticket(self):
tid=self.setup_card(self.grant)
with patch.object(op,'choose',return_value={'profile':'wggesucht','reason':'Match'}),patch.object(w,'run_one',side_effect=RuntimeError('HAI scope denied')):
result=op.process_one(self.manifest)
self.assertEqual(result['status'],'blocked')
with w.db('test') as c:
self.assertEqual(kb.get_task(c,tid).status,'blocked')
self.assertIn('meldet blocked',c.execute('SELECT body FROM task_comments ORDER BY id DESC LIMIT 1').fetchone()[0])
class ObserverTests(unittest.TestCase):
def run_loop(self, authority_error=None, snapshots=None):
with tempfile.TemporaryDirectory() as tmp:
manifest=Path(tmp)/'manifest.json'
manifest.write_text(json.dumps({'board':'test','routes':{}}))
status=Path(tmp)/'status.json'
frames=snapshots or [{'open_tasks':[], 'ready_tasks':['new']}]*3
with patch.object(op,'observe',side_effect=frames) as observe, \
patch.object(op,'authority',side_effect=authority_error) as authority, \
patch.object(op,'process_one',return_value=None) as process, \
patch.object(op.time,'sleep'), \
patch.object(op.time,'monotonic',return_value=1):
op.serve(manifest,status_path=status,interval=0,stop=lambda:observe.call_count>=len(frames))
return json.loads(status.read_text()),authority.call_count,process.call_count
def test_expired_lease_keeps_observing_without_execution_or_retry_storm(self):
state,checks,runs=self.run_loop(PermissionError('lease_expired'))
self.assertEqual(state['phase'],'waiting_for_hai')
self.assertEqual(checks,1)
self.assertEqual(runs,0)
def test_new_card_after_idle_triggers_authorized_processing(self):
state,checks,runs=self.run_loop(snapshots=[
{'open_tasks':[], 'ready_tasks':[]},
{'open_tasks':[{'id':'new'}], 'ready_tasks':['new']},
{'open_tasks':[], 'ready_tasks':[]}])
self.assertEqual(checks,1)
self.assertEqual(runs,1)
self.assertEqual(state['phase'],'watching')
if __name__=='__main__':unittest.main()