File size: 5,542 Bytes
3bc7cb3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
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()