345 lines
19 KiB
Python
345 lines
19 KiB
Python
#!/usr/bin/env python3
|
|
"""Actual ROS2 + native executor regression; run inside a built Humble overlay.
|
|
|
|
Uses only /sim/robot_01 fixture servers, never robot endpoints. Includes durable
|
|
CDR evidence checks, executor restart receipt recovery, and cancellation. This
|
|
is transport/orchestration evidence, not physical NAV/VLA acceptance.
|
|
"""
|
|
import argparse
|
|
import copy
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import shlex
|
|
import signal
|
|
import subprocess
|
|
import sys
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
|
|
|
|
def wait(future, seconds=45):
|
|
end = time.monotonic() + seconds
|
|
while not future.done() and time.monotonic() < end:
|
|
time.sleep(0.01)
|
|
if not future.done():
|
|
raise TimeoutError('ROS future deadline exceeded')
|
|
return future.result()
|
|
|
|
|
|
def eventually(predicate, seconds=8):
|
|
end = time.monotonic() + seconds
|
|
while time.monotonic() < end:
|
|
if predicate():
|
|
return
|
|
time.sleep(0.02)
|
|
raise AssertionError('expected native state transition did not occur')
|
|
|
|
|
|
def journal_records(path):
|
|
"""Independent decoder for durable registry v4 (does not call production)."""
|
|
latest = {}
|
|
for line in path.read_text().splitlines():
|
|
row = shlex.split(line)
|
|
assert row[0] in ('3', '4') and len(row) == (25 if row[0] == '4' else 23), 'missing durable evidence fields'
|
|
decode = lambda index: bytes.fromhex(row[index]).decode('utf-8')
|
|
latest[row[1]] = dict(goal_id=row[1], task_id=row[3], run_id=row[5],
|
|
skill=int(row[10]), state=int(row[11]), sequence=int(row[14]),
|
|
type=decode(18), wire=decode(19), feedback=decode(20),
|
|
error_code=decode(21), execution_record_ref=decode(22), result_type=decode(23) if row[0] == '4' else '',
|
|
result_wire=decode(24) if row[0] == '4' else '')
|
|
return list(latest.values())
|
|
|
|
|
|
def assert_manipulation_evidence(records, run_id, action_type, deserialize):
|
|
selected = [row for row in records if row['run_id'] == run_id and
|
|
'ExecuteManipulation' in row['type']]
|
|
assert len(selected) == 2, 'one pick and one place wire snapshot required'
|
|
decoded = []
|
|
for row in selected:
|
|
goal = deserialize(bytes.fromhex(row['wire']), action_type.Goal)
|
|
assert goal.skill in ('pick', 'place')
|
|
assert goal.trace.run_id == run_id and goal.trace.attempt == 1
|
|
assert goal.target.object_ref == 'water' and goal.target.description
|
|
assert goal.instruction == goal.skill + ' registered target water'
|
|
assert goal.timeout.sec == 5 and goal.timeout.nanosec == 0
|
|
if goal.skill == 'pick':
|
|
assert goal.destination.region_ref == goal.destination.description == ''
|
|
else:
|
|
assert goal.destination.region_ref == 'tote_A' and goal.destination.description
|
|
payload = json.loads(row['feedback'])
|
|
feedback = deserialize(bytes.fromhex(payload['cdr_hex']), action_type.Feedback)
|
|
assert feedback.sequence == row['sequence'] == payload['sequence'] > 0
|
|
assert feedback.phase == payload['phase'] and feedback.message == payload['message']
|
|
assert feedback.progress_valid == payload['progress_valid']
|
|
assert feedback.message and payload['elapsed_time_ns'] >= 0
|
|
assert row['execution_record_ref'].startswith('sim://execute_manipulation/')
|
|
assert row['state'] == 4
|
|
assert 'ExecuteManipulation' in row['result_type'] and row['result_wire']
|
|
result = deserialize(bytes.fromhex(row['result_wire']), action_type.Result)
|
|
assert result.result.status == 0 and result.result.stop_state == 1
|
|
assert result.result.stop_evidence_ref and result.result.stopped_at.sec > 0
|
|
assert result.execution_record_ref == row['execution_record_ref']
|
|
decoded.append(dict(skill=goal.skill, sequence=feedback.sequence,
|
|
phase=feedback.phase, execution_record_ref=row['execution_record_ref']))
|
|
assert {row['skill'] for row in decoded} == {'pick', 'place'}
|
|
return decoded
|
|
|
|
|
|
def run_route(root, binary, output, suffix, recovery_enabled):
|
|
import rclpy
|
|
from rclpy.action import ActionClient
|
|
from rclpy.executors import MultiThreadedExecutor
|
|
from rclpy.node import Node
|
|
from rclpy.serialization import deserialize_message
|
|
from bt_skill_interfaces.action import ExecuteTask, ExecuteManipulation
|
|
from navigation_interfaces.action import NavigateToPose
|
|
from bt_skill_interfaces.msg import RobotState
|
|
from bt_skill_interfaces.srv import ReconcileTask
|
|
from robot_bt_coordinator.plan_v2 import item_plan
|
|
from bt_mock_servers.mock_skills import MockSkills, ACTION_TYPES
|
|
|
|
site_path = root / 'config' / ('sim_site_' + suffix + '.json')
|
|
site = json.loads(site_path.read_text())
|
|
plan = item_plan(json.loads((root / 'config' / ('demo_plan_' + suffix + '.json')).read_text()), 0)
|
|
stops = ['shelf_A_stop', 'tote_A_stop'] if suffix == 'object_table' else ['observe_A', 'shelf_A_stop', 'tote_A_stop']
|
|
scenarios = {'verify_state': [{'kind': 'passed'}],
|
|
'locate_shelf_column': [{'kind': 'normal', 'shelf_id': 'shelf_A', 'side_id': 'FRONT', 'column_id': '1', 'tier_id': '2'}],
|
|
'navigate': [{'kind': 'normal', 'final_pose': site['locations'][name]} for name in stops]}
|
|
actions = [name for name in ACTION_TYPES if name != 'execute_task']
|
|
rclpy.init(args=['--ros-args', '-p', 'initial_holding_state:=EMPTY',
|
|
'-p', 'scenarios_json:=' + json.dumps(json.dumps(scenarios)),
|
|
'-p', 'enabled_actions:=' + json.dumps(actions), '-p',
|
|
'enabled_topics:=' + json.dumps(['robot_state', 'safety_state', 'visual_observation', 'dense_progress'])])
|
|
server = MockSkills()
|
|
node = Node('native_workflow_regression', namespace='/sim/robot_01')
|
|
executor = MultiThreadedExecutor(num_threads=8)
|
|
executor.add_node(server)
|
|
executor.add_node(node)
|
|
thread = threading.Thread(target=executor.spin, daemon=True)
|
|
thread.start()
|
|
client = ActionClient(node, ExecuteTask, 'tasks/execute')
|
|
service = node.create_client(ReconcileTask, 'tasks/reconcile')
|
|
feedback = []
|
|
row = dict(route=plan['route'], scope='native ROS2 executor with simulated physical skills')
|
|
process = None
|
|
log = (output / ('native-' + suffix + '.log')).open('w')
|
|
token = 'native-regression-recovery-token'
|
|
operator = 'native-regression-operator'
|
|
with tempfile.TemporaryDirectory(prefix='native-bt-regression-') as journal:
|
|
def stop():
|
|
nonlocal process
|
|
if process is not None and process.poll() is None:
|
|
os.killpg(process.pid, signal.SIGINT)
|
|
try:
|
|
process.wait(timeout=5)
|
|
except subprocess.TimeoutExpired:
|
|
os.killpg(process.pid, signal.SIGKILL)
|
|
process.wait()
|
|
process = None
|
|
|
|
def start(expect_failure=False):
|
|
nonlocal process
|
|
args = [str(binary), '--ros-args', '-r', '__ns:=/sim/robot_01',
|
|
'-p', 'robot_id:=robot_01', '-p', 'allowed_robots:=[robot_01]',
|
|
'-p', 'site_config_file:=' + str(site_path), '-p', 'journal_directory:=' + journal,
|
|
'-p', 'execution_enabled:=true', '-p', 'skill_timeout_ms:=5000',
|
|
'-p', 'recovery_token:=' + token, '-p', 'recovery_operator_id:=' + operator]
|
|
process = subprocess.Popen(args, stdout=log, stderr=subprocess.STDOUT, start_new_session=True)
|
|
if expect_failure:
|
|
code = process.wait(timeout=5)
|
|
assert code != 0, 'corrupt task journal accepted at startup'
|
|
return
|
|
assert client.wait_for_server(timeout_sec=12), 'native executor not discovered'
|
|
time.sleep(0.6)
|
|
assert process.poll() is None, 'native executor exited during startup'
|
|
|
|
def task_goal(tag):
|
|
goal = ExecuteTask.Goal()
|
|
goal.trace.task_id = 'native-' + suffix + '-' + tag
|
|
goal.trace.run_id = 'run-' + suffix + '-' + tag
|
|
goal.trace.subtask_id = 'execute_task'
|
|
goal.trace.attempt = goal.trace.task_revision = goal.trace.plan_version = goal.trace.execution_generation = 1
|
|
goal.approved_plan_json = json.dumps(plan)
|
|
goal.context_json = json.dumps(dict(schema_version=1, registry_version=1, robot_id='robot_01', target_id='water', source_shelf='shelf_A', destination_id='tote_A', observe_location='observe_A', destination_location='tote_A_stop', task_revision=1, route=plan['route'], item_index=0))
|
|
goal.timeout.sec = 30
|
|
return goal
|
|
|
|
def reconcile(goal, authorization=token, stale=False, resolution='resume_task'):
|
|
assert service.wait_for_service(timeout_sec=10), 'recovery service not discovered'
|
|
request = ReconcileTask.Request()
|
|
request.trace = copy.deepcopy(goal.trace)
|
|
if stale:
|
|
request.trace.execution_generation += 1
|
|
request.operator_id, request.authorization = operator, authorization
|
|
request.evidence_ref = 'regression/operator-observed-stop'
|
|
request.resolution = resolution
|
|
return wait(service.call_async(request), 15)
|
|
|
|
def motion_counts():
|
|
with server.lock:
|
|
return {name: server.counts[name] for name in ('navigate', 'execute_manipulation', 'execute_posture')}
|
|
|
|
try:
|
|
start()
|
|
goal = task_goal('complete')
|
|
handle = wait(client.send_goal_async(goal, feedback_callback=lambda msg: feedback.append(msg.feedback.stage)), 8)
|
|
assert handle.accepted, 'valid stage-one route was rejected'
|
|
wrapped = wait(handle.get_result_async())
|
|
result = wrapped.result
|
|
row.update(native_status=wrapped.status, status=result.result.status,
|
|
completed_quantity=result.completed_quantity, evidence=json.loads(result.evidence_json))
|
|
assert wrapped.status == 4 and result.result.status == 0, result.result.message
|
|
assert result.completed_quantity == 1 and result.result.stop_state == 1
|
|
assert row['evidence']['safe_to_release'] is True
|
|
row['decoded_manipulation'] = assert_manipulation_evidence(
|
|
journal_records(Path(journal) / 'goal_registry.log'), goal.trace.run_id,
|
|
ExecuteManipulation, deserialize_message)
|
|
navigation_rows = [entry for entry in journal_records(Path(journal) / 'goal_registry.log')
|
|
if entry['run_id'] == goal.trace.run_id and entry['skill'] == 0]
|
|
assert len(navigation_rows) == len(stops)
|
|
expected_poses = {(site['locations'][name]['x'], site['locations'][name]['y']) for name in stops}
|
|
actual_poses = set()
|
|
for entry in navigation_rows:
|
|
assert entry['type'] == 'navigation_interfaces/action/NavigateToPose_Goal'
|
|
navigation_goal = deserialize_message(bytes.fromhex(entry['wire']), NavigateToPose.Goal)
|
|
assert navigation_goal.task_id == goal.trace.task_id and navigation_goal.subtask_id
|
|
assert navigation_goal.target_pose.header.frame_id == 'map'
|
|
actual_poses.add((navigation_goal.target_pose.pose.position.x, navigation_goal.target_pose.pose.position.y))
|
|
assert entry['result_type'] == 'navigation_interfaces/action/NavigateToPose_Result'
|
|
navigation_result = deserialize_message(bytes.fromhex(entry['result_wire']), NavigateToPose.Result)
|
|
assert navigation_result.status == NavigateToPose.Result.SUCCEEDED
|
|
assert navigation_result.stop_state == NavigateToPose.Result.STOP_CONFIRMED
|
|
assert actual_poses == expected_poses
|
|
row['direct_navigation_goals'] = len(navigation_rows)
|
|
assert motion_counts()['execute_manipulation'] == 2
|
|
if recovery_enabled:
|
|
stop()
|
|
time.sleep(0.8)
|
|
start()
|
|
before = motion_counts()
|
|
wrong = reconcile(goal, authorization='wrong-token')
|
|
assert not wrong.accepted, 'wrong authorization accepted'
|
|
stale = reconcile(goal, stale=True)
|
|
assert not stale.accepted, 'stale execution generation accepted'
|
|
with server.lock:
|
|
server.scenarios['verify_state'] = [{'kind': 'unknown'}]
|
|
server.holding_state = RobotState.HOLDING_UNKNOWN
|
|
time.sleep(0.25)
|
|
unknown = reconcile(goal)
|
|
assert not unknown.accepted, 'unknown holding was treated as safe recovery'
|
|
with server.lock:
|
|
server.scenarios['verify_state'] = [{'kind': 'passed'}]
|
|
server.holding_state = RobotState.EMPTY
|
|
time.sleep(0.25)
|
|
recovered = reconcile(goal)
|
|
assert recovered.accepted, recovered.error_code + ': ' + recovered.message
|
|
state = json.loads(recovered.state_json)
|
|
assert state['verified'] is True and state['stop_confirmed'] is True
|
|
assert state['holding_state'] == 'EMPTY'
|
|
receipts = state['receipts']
|
|
assert len(receipts) == 1 and receipts[0]['task_id'] == goal.trace.task_id
|
|
assert receipts[0]['completed_quantity'] == 1
|
|
assert motion_counts() == before, 'receipt recovery dispatched motion'
|
|
row['recovery'] = dict(wrong_token=wrong.error_code, stale_trace=stale.error_code,
|
|
unknown_holding=unknown.error_code, recovered=state)
|
|
if suffix == 'object_table':
|
|
with server.lock:
|
|
server.scenarios['navigate'] = [{'kind': 'normal', 'duration_seconds': 3., 'final_pose': site['locations']['shelf_A_stop']}]
|
|
canceled_goal = task_goal('cancel')
|
|
canceled_handle = wait(client.send_goal_async(canceled_goal), 8)
|
|
assert canceled_handle.accepted
|
|
eventually(lambda: motion_counts()['navigate'] > before['navigate'])
|
|
wait(canceled_handle.cancel_goal_async(), 5)
|
|
canceled_result = wait(canceled_handle.get_result_async(), 15)
|
|
assert canceled_result.status == 5 and canceled_result.result.result.stop_state == 1
|
|
assert motion_counts()['execute_manipulation'] == before['execute_manipulation']
|
|
stop()
|
|
time.sleep(0.8)
|
|
start()
|
|
# Absence of physical-goal history cannot prove no dispatch.
|
|
stop()
|
|
goal_log = Path(journal) / 'goal_registry.log'
|
|
intact_goals = goal_log.read_bytes()
|
|
goal_log.write_bytes(b'')
|
|
start()
|
|
incomplete = reconcile(canceled_goal, resolution='resume_task')
|
|
assert not incomplete.accepted, 'missing dispatch history authorized an item replay'
|
|
stop()
|
|
goal_log.write_bytes(intact_goals)
|
|
start()
|
|
canceled_recovery = reconcile(canceled_goal, resolution='cancel_task')
|
|
assert canceled_recovery.accepted, canceled_recovery.error_code
|
|
row['cancel_recovery'] = json.loads(canceled_recovery.state_json)
|
|
row['missing_history_rejected'] = incomplete.error_code
|
|
stop()
|
|
task_log = Path(journal) / 'task_runs.jsonl'
|
|
intact_tasks = task_log.read_bytes()
|
|
assert intact_tasks.endswith(b'\n')
|
|
task_log.write_bytes(intact_tasks[:-1])
|
|
start(expect_failure=True)
|
|
row['unterminated_task_journal_rejected'] = True
|
|
stop()
|
|
task_log.write_bytes(intact_tasks)
|
|
if suffix == 'shelf_cell':
|
|
with server.lock:
|
|
server.scenarios['locate_shelf_column'] = [{'kind': 'ambiguous'}]
|
|
server.scenarios['navigate'] = [{'kind': 'normal', 'final_pose': site['locations']['observe_A']}]
|
|
before_ask = motion_counts()
|
|
ask_handle = wait(client.send_goal_async(task_goal('ambiguous')), 8)
|
|
assert ask_handle.accepted
|
|
ask_wrapped = wait(ask_handle.get_result_async(), 20)
|
|
ask = ask_wrapped.result
|
|
ask_evidence = json.loads(ask.evidence_json)
|
|
assert ask_wrapped.status == 6 and ask.completed_quantity == 0
|
|
assert ask_evidence['safe_to_release'] is True
|
|
assert ask_evidence['needs_clarification'] is True
|
|
assert ask_evidence['safe_to_retry'] is True
|
|
assert motion_counts()['execute_manipulation'] == before_ask['execute_manipulation']
|
|
row['runtime_clarification'] = ask_evidence
|
|
row['mock_calls'] = dict(server.counts)
|
|
row['feedback_stages'] = feedback
|
|
row['passed'] = True
|
|
except Exception as exc:
|
|
row['passed'] = False
|
|
row['failure'] = repr(exc)
|
|
finally:
|
|
stop()
|
|
for path in Path(journal).glob('*'):
|
|
if path.is_file():
|
|
(output / ('native-' + suffix + '-' + path.name)).write_bytes(path.read_bytes())
|
|
client.destroy()
|
|
executor.shutdown(timeout_sec=2)
|
|
thread.join(timeout=2)
|
|
server.destroy_node()
|
|
node.destroy_node()
|
|
rclpy.shutdown()
|
|
log.close()
|
|
return row
|
|
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser(description=__doc__)
|
|
parser.add_argument('--repo-root', type=Path, default=Path(__file__).resolve().parents[2])
|
|
parser.add_argument('--executor', type=Path)
|
|
parser.add_argument('--output-dir', type=Path, required=True)
|
|
parser.add_argument('--skip-recovery', action='store_true', help='Run workflow/evidence subset before recovery service is built')
|
|
args = parser.parse_args()
|
|
root = args.repo_root.resolve()
|
|
sys.path[:0] = [str(root / 'coordinator'), str(root / 'ros2/bt_mock_servers')]
|
|
binary = args.executor or root / 'install/bt_executor/lib/bt_executor/bt_executor_node'
|
|
args.output_dir.mkdir(parents=True, exist_ok=True)
|
|
rows = []
|
|
for suffix in ('object_table', 'shelf_cell'):
|
|
row = run_route(root, binary, args.output_dir, suffix, not args.skip_recovery)
|
|
rows.append(row)
|
|
(args.output_dir / 'native-workflow-regression.json').write_text(json.dumps(rows, indent=2))
|
|
print(json.dumps({k: row[k] for k in ('route', 'passed', 'failure') if k in row}), flush=True)
|
|
if not all(row['passed'] for row in rows):
|
|
raise SystemExit(1)
|
|
|
|
|
|
if __name__ == '__main__':
|
|
main()
|