import unittest,tempfile,sys,copy from pathlib import Path sys.path.insert(0,str(Path(__file__).resolve().parents[1]/'coordinator')) from robot_bt_coordinator.service import Coordinator,demo_site from robot_bt_coordinator.backends import ManualBackend from robot_bt_coordinator.replay import replay_events from test_plan_v2 import sample class MultiItemTests(unittest.TestCase): def setUp(self): self.tmp=tempfile.TemporaryDirectory();self.backend=ManualBackend();self.site=demo_site();self.site['execution_route']='OBJECT_TABLE' self.c=Coordinator(self.tmp.name+'/tasks.db',self.backend,['r'],self.site) def tearDown(self):self.c.close();self.tmp.cleanup() def start(self): try:t=self.c.submit(dict(client_request_id='q',robot_id='r',instruction='two waters then one doll',known_info=sample()['slots'])) except Exception as ex:self.fail(str(ex)) self.tid=t['task_id'];self.c.tick();t=self.c.get(self.tid) self.backend.emit(dict(type='plan',task_id=self.tid,task_revision=t['task_revision'],planning_generation=t['planning_generation'],status='PLAN_READY',plan=sample()));self.c.tick() self.assertEqual(self.c.get(self.tid)['status'],'EXECUTING') def complete(self,**changes): t=self.c.get(self.tid);target=t['context']['target_id'];e=dict(type='execution_result',task_id=self.tid,run_id=t['run_id'],status='SUCCEEDED',completed_quantity=1,stop_confirmed=True,evidence=dict(passed=True,empty_hand=True,in_destination=True,valid=True,safe_to_release=True,evidence_id='ev'+t['run_id'],target_ref=target,destination_ref='tote_A'));e.update(changes) self.backend.emit(e);self.c.tick();return e def test_three_runs_and_idempotent_delivery(self): self.start();runs=[] for n in range(3): runs.append(self.c.get(self.tid)['run_id']);e=self.complete();self.backend.emit(e);self.c.tick() self.assertEqual(self.c.get(self.tid)['completed_quantity'],n+1) t=self.c.get(self.tid);self.assertEqual(t['status'],'SUCCEEDED');self.assertEqual(len(set(runs)),3) self.assertEqual([x['context']['target_id'] for x in self.backend.executions],['water','water','doll']) self.assertEqual(replay_events(self.c.events(self.tid,limit=500))['completed_quantity'],3) def test_partial_failure_retains_quantity_and_blocks_queue(self): self.start();self.complete();self.complete(status='INTERVENTION_REQUIRED',stop_confirmed=False,completed_quantity=0) t=self.c.get(self.tid);self.assertEqual(t['status'],'INTERVENTION_REQUIRED');self.assertEqual(t['completed_quantity'],1);self.assertEqual(len(self.backend.executions),2) def test_cancel_racing_success_does_not_start_next_item(self): self.start();self.c.control(self.tid,'cancel');self.complete() self.assertEqual(self.c.get(self.tid)['status'],'CANCELED');self.assertEqual(len(self.backend.executions),1) def test_restart_after_partial_delivery_requires_reconciliation(self): self.start();self.complete();self.c.close();self.backend=ManualBackend() self.c=Coordinator(self.tmp.name+'/tasks.db',self.backend,['r'],self.site) self.assertEqual(self.c.get(self.tid)['status'],'INTERVENTION_REQUIRED');self.assertEqual(self.c.get(self.tid)['completed_quantity'],1) self.c.tick();self.assertEqual(self.backend.executions,[]) def test_advisory_progress_cannot_commit_delivery(self): self.start();t=self.c.get(self.tid) self.backend.emit(dict(type='advisory',task_id=self.tid,run_id=t['run_id'],state='RUNNING',progress=1.));self.c.tick() self.assertEqual(self.c.get(self.tid)['status'],'EXECUTING');self.assertEqual(self.c.get(self.tid)['completed_quantity'],0)