import copy import json from pathlib import Path import sqlite3 import sys import tempfile import unittest from unittest.mock import patch ROOT=Path(__file__).resolve().parents[1] sys.path[:0]=[str(ROOT/'coordinator'),str(ROOT/'robobrain')] from robot_bt_coordinator.service import Coordinator,demo_site from robot_bt_coordinator.backends import ManualBackend,demo_plan from robot_bt_coordinator.errors import ApiError from robot_bt_coordinator.replay import replay_events from robot_bt_coordinator.store import Store from robot_bt_coordinator.plan_v2 import make_plan class TrustedBackend(ManualBackend): def reconcile(self,task,request): return dict(verified=True,stop_confirmed=True,holding_state='EMPTY',run_id=task['run_id'],evidence_ref=request['evidence_ref'],receipts=copy.deepcopy(self.receipts),safe_to_retry=self.retry) class RecoveryTest(unittest.TestCase): def setUp(self): self.tmp=tempfile.TemporaryDirectory();self.backend=TrustedBackend();self.backend.receipts=[];self.backend.retry=False self.site=demo_site();self.site['object_aliases']={'water':['water'],'doll':['doll']} 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): task=self.c.submit(dict(client_request_id='q',robot_id='r',instruction='move item',known_info={'target_name':'water','source_location':'shelf_A','destination':'tote_A'})) self.c.tick();task=self.c.get(task['task_id']) self.backend.emit(dict(type='plan',task_id=task['task_id'],task_revision=task['task_revision'],planning_generation=task['planning_generation'],status='PLAN_READY',plan=demo_plan(task['request']['known_info']))) self.c.tick();return self.c.get(task['task_id']) def quarantine(self,task): self.backend.emit(dict(type='execution_result',task_id=task['task_id'],run_id=task['run_id'],status='INTERVENTION_REQUIRED',completed_quantity=0,stop_confirmed=False)) self.c.tick() def request(self,task,resolution='resume_task'): return dict(run_id=task['run_id'],evidence_ref='trusted-proof',resolution=resolution) def receipt(self,task): return dict(task_id=task['task_id'],item_index=0,completed_quantity=1,evidence=dict(evidence_id='delivery-proof',target_ref='water',destination_ref='tote_A',passed=True,empty_hand=True,in_destination=True,valid=True)) def test_unsafe_resume_cannot_replay_an_uncertain_manipulation(self): task=self.start();self.quarantine(task) with self.assertRaises(ApiError) as error:self.c.intervene(task['task_id'],self.request(task)) self.assertEqual(error.exception.code,'RECONCILIATION_REJECTED') self.assertEqual(len(self.backend.executions),1) def test_trusted_receipt_recovers_lost_result_without_redispatch(self): task=self.start();self.quarantine(task);self.backend.receipts=[self.receipt(task)] out=self.c.intervene(task['task_id'],self.request(task)) self.assertEqual(out['status'],'SUCCEEDED');self.assertEqual(out['completed_quantity'],1) self.assertEqual(len(self.backend.executions),1) replay=replay_events(self.c.events(task['task_id'],limit=500)) self.assertEqual(replay['completed_quantity'],1) def test_empty_hand_alone_cannot_credit_or_resume(self): task=self.start();self.quarantine(task) with self.assertRaises(ApiError):self.c.intervene(task['task_id'],self.request(task)) out=self.c.intervene(task['task_id'],self.request(task,'cancel_task')) self.assertEqual(out['status'],'CANCELED');self.assertEqual(out['completed_quantity'],0) def clarification_result(self,task,**changes): evidence=dict(needs_clarification=True,safe_to_retry=True,empty_hand=True,valid=True,safe_to_release=True,questions=['Which registered shelf should be used?']) evidence.update(changes) self.backend.emit(dict(type='execution_result',task_id=task['task_id'],run_id=task['run_id'],status='FAILED',completed_quantity=0,stop_confirmed=True,evidence=evidence)) self.c.tick();return self.c.get(task['task_id']) def test_verified_runtime_question_returns_to_planning_after_answer(self): task=self.start();out=self.clarification_result(task) self.assertEqual(out['status'],'NEEDS_CLARIFICATION');self.assertFalse(out['motion_dispatched']) self.assertEqual(out['question']['source_run_id'],task['run_id']) out=self.c.clarify(task['task_id'],dict(question_id=out['question']['question_id'],task_revision=out['task_revision'],known_info={'source_location':'shelf_A'})) self.c.tick();self.assertEqual(self.c.get(task['task_id'])['status'],'PLANNING') self.assertTrue(self.backend.plans[-1]['clarification_confirmed']) self.assertEqual(len(self.backend.executions),1) def test_runtime_question_without_no_manipulation_proof_is_quarantined(self): task=self.start();out=self.clarification_result(task,safe_to_retry=False) self.assertEqual(out['status'],'INTERVENTION_REQUIRED') self.assertIsNone(out['question']) def test_ordinary_failure_does_not_fabricate_question(self): task=self.start();out=self.clarification_result(task,needs_clarification=False) self.assertEqual(out['status'],'FAILED');self.assertIsNone(out['question']) def test_operator_replan_requires_no_manipulation_and_no_delivery(self): task=self.start();self.quarantine(task) with self.assertRaises(ApiError):self.c.intervene(task['task_id'],self.request(task,'replan_task')) self.backend.retry=True out=self.c.intervene(task['task_id'],self.request(task,'replan_task')) self.assertEqual(out['status'],'NEEDS_CLARIFICATION');self.assertFalse(out['motion_dispatched']) self.assertEqual(len(self.backend.executions),1) def test_operator_replan_cannot_rewrite_already_delivered_scope(self): task=self.start();self.quarantine(task);self.backend.retry=True;self.backend.receipts=[self.receipt(task)] with self.assertRaises(ApiError):self.c.intervene(task['task_id'],self.request(task,'replan_task')) self.assertEqual(self.c.get(task['task_id'])['status'],'INTERVENTION_REQUIRED') def test_question_cannot_hide_a_verified_physical_delivery(self): task=self.start();proof=self.receipt(task)['evidence'] proof.update(needs_clarification=True,safe_to_retry=True,safe_to_release=True,questions=['Which shelf?']) self.backend.emit(dict(type='execution_result',task_id=task['task_id'],run_id=task['run_id'],status='FAILED',completed_quantity=1,stop_confirmed=True,evidence=proof));self.c.tick() out=self.c.get(task['task_id']) self.assertEqual(out['completed_quantity'],1) self.assertEqual(out['status'],'INTERVENTION_REQUIRED') def test_proven_pre_manipulation_resume_uses_new_run(self): task=self.start();self.quarantine(task);self.backend.retry=True out=self.c.intervene(task['task_id'],self.request(task)) self.assertEqual(out['status'],'EXECUTING');self.assertNotEqual(out['run_id'],task['run_id']) self.assertEqual(len(self.backend.executions),2) def test_paused_execution_accepts_authenticated_reconciliation(self): task=self.start();self.c.control(task['task_id'],'pause') self.backend.emit(dict(type='execution_result',task_id=task['task_id'],run_id=task['run_id'],status='CANCELED',completed_quantity=0,stop_confirmed=True)) self.c.tick();self.assertEqual(self.c.get(task['task_id'])['status'],'PAUSED') self.backend.retry=True self.assertEqual(self.c.intervene(task['task_id'],self.request(task))['status'],'EXECUTING') def test_unverified_resume_preserves_paused_state(self): task=self.start();self.c.control(task['task_id'],'pause') self.backend.emit(dict(type='execution_result',task_id=task['task_id'],run_id=task['run_id'],status='CANCELED',completed_quantity=0,stop_confirmed=True)) self.c.tick() with self.assertRaises(ApiError) as error:self.c.control(task['task_id'],'resume') self.assertEqual(error.exception.code,'RECONCILIATION_REQUIRED') self.assertEqual(self.c.get(task['task_id'])['status'],'PAUSED') self.assertEqual(len(self.backend.executions),1) def test_completed_planning_releases_timer_state(self): self.start();self.assertEqual(self.c._planning_started,{}) def test_duplicate_recovered_receipts_credit_once_then_start_next_item(self): self.site['execution_route']='OBJECT_TABLE' slots=dict(items=[dict(target_name='water',quantity=2,source_location='shelf_A')],destination='tote_A') task=self.c.submit(dict(client_request_id='multi',robot_id='r',instruction='move two',known_info=slots)) self.c.tick();task=self.c.get(task['task_id']) self.backend.emit(dict(type='plan',task_id=task['task_id'],task_revision=1,planning_generation=task['planning_generation'],status='PLAN_READY',plan=make_plan('move two',slots,'OBJECT_TABLE'))) self.c.tick();task=self.c.get(task['task_id']);self.quarantine(task) receipt=self.receipt(task);self.backend.receipts=[receipt,copy.deepcopy(receipt)] out=self.c.intervene(task['task_id'],self.request(task)) self.assertEqual(out['status'],'EXECUTING');self.assertEqual(out['completed_quantity'],1) self.assertEqual(out['active_item_index'],1);self.assertNotEqual(out['run_id'],task['run_id']) self.assertEqual([e['item_index'] for e in self.c.events(task['task_id'],limit=500) if e['kind']=='delivery_committed'],[0]) def test_empty_clarification_cannot_mark_intent_as_confirmed(self): task=self.c.submit(dict(client_request_id='intent',robot_id='r',instruction='1 water shelf_A tote_A',known_info=dict(target_name='doll'))) with self.assertRaises(ApiError): self.c.clarify(task['task_id'],dict(question_id=task['question']['question_id'],task_revision=task['task_revision'],known_info={})) self.assertFalse(self.c.get(task['task_id'])['clarification_confirmed']) def test_invalid_receipts_are_atomic_and_do_not_release_quarantine(self): for mutate in [lambda r:r.update(task_id='other'),lambda r:r.update(item_index=True),lambda r:r.update(completed_quantity=2),lambda r:r['evidence'].update(destination_ref='wrong')]: task=self.start() if not hasattr(self,'current') else self.current self.current=task;self.quarantine(task) receipt=self.receipt(task);mutate(receipt);self.backend.receipts=[receipt] with self.assertRaises(ApiError):self.c.intervene(task['task_id'],self.request(task,'cancel_task')) self.assertEqual(self.c.get(task['task_id'])['completed_quantity'],0) self.assertEqual(self.c.get(task['task_id'])['status'],'INTERVENTION_REQUIRED') def test_tick_avoids_historical_task_scan(self): for n in range(30): t=self.c.submit(dict(client_request_id=str(n),robot_id='r',instruction='move')) self.c.control(t['task_id'],'cancel') with patch.object(self.c.store,'all',side_effect=AssertionError('historical scan')): self.c.tick() def test_progress_replay_preserves_stage_sequence_and_detail(self): task=self.start() self.backend.emit(dict(type='progress',task_id=task['task_id'],run_id=task['run_id'],sequence=1,stage='NavigateSource',detail='waiting obstacle')) self.c.tick() events=self.c.events(task['task_id'],limit=500) progress=next(e for e in events if e['kind']=='progress') self.assertEqual(progress['stage'],'NavigateSource');self.assertEqual(progress['sequence'],1) self.assertEqual(progress['detail'],'waiting obstacle') self.assertEqual(replay_events(events)['progress'][-1]['stage'],'NavigateSource') def test_supplied_slots_conflicting_with_instruction_require_clarification(self): task=self.c.submit(dict(client_request_id='intent',robot_id='r',instruction='1 water shelf_A tote_A',known_info=dict(target_name='doll',quantity=1,source_location='shelf_A',destination='tote_A'))) self.assertEqual(task['status'],'NEEDS_CLARIFICATION');self.assertFalse(task['clarification_confirmed']) self.c.tick();self.assertEqual(self.backend.plans,[]) task=self.c.clarify(task['task_id'],dict(question_id=task['question']['question_id'],task_revision=task['task_revision'],known_info=dict(target_name='doll'))) self.assertTrue(task['clarification_confirmed']) self.c.close();self.c=Coordinator(self.tmp.name+'/tasks.db',self.backend,['r'],self.site) self.assertTrue(self.c.get(task['task_id'])['clarification_confirmed']) self.c.tick();self.assertTrue(self.backend.plans[-1]['clarification_confirmed']) def test_approved_plan_is_independently_checked_against_instruction(self): task=self.c.submit(dict(client_request_id='intent',robot_id='r',instruction='1 water shelf_A tote_A',known_info={})) self.c.tick();task=self.c.get(task['task_id']) plan=demo_plan(dict(target_name='doll',quantity=1,source_location='shelf_A',destination='tote_A')) self.backend.emit(dict(type='plan',task_id=task['task_id'],task_revision=1,planning_generation=task['planning_generation'],status='PLAN_READY',plan=plan));self.c.tick() self.assertEqual(self.c.get(task['task_id'])['status'],'NEEDS_CLARIFICATION');self.assertEqual(self.backend.executions,[]) class MigrationTest(unittest.TestCase): def test_legacy_database_migrates_active_status_without_losing_history(self): with tempfile.TemporaryDirectory() as folder: path=folder+'/old.db';db=sqlite3.connect(path) db.execute('CREATE TABLE tasks(seq INTEGER PRIMARY KEY AUTOINCREMENT,task_id TEXT UNIQUE NOT NULL,robot_id TEXT NOT NULL,request_id TEXT NOT NULL,request_hash TEXT NOT NULL,data TEXT NOT NULL,UNIQUE(robot_id,request_id))') for n,status in enumerate(('SUCCEEDED','QUEUED','PAUSED')): db.execute('INSERT INTO tasks(task_id,robot_id,request_id,request_hash,data) VALUES(?,?,?,?,?)',(str(n),'r',str(n),'hash',json.dumps(dict(task_id=str(n),status=status)))) db.commit();db.close();store=Store(path) try: self.assertEqual([t['status'] for t in store.active()],['QUEUED','PAUSED']) self.assertEqual(len(store.all()),3) finally:store.close() if __name__=='__main__':unittest.main()