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()