From 931530eefd0b9f92b7830b12f03c64082c169fed Mon Sep 17 00:00:00 2001 From: zogoo Date: Sun, 31 Oct 2021 18:44:01 +0900 Subject: [PATCH 01/14] Internal communication between camera and motors --- Camera.py | 60 +++---- Excavator.py | 8 + camera_node.py | 114 +++++++++++++ client_message.py | 208 ++++++++++++++++++++++++ motor_node.py | 79 +++++++++ server_message.py | 212 +++++++++++++++++++++++++ detect_picamera.py => test_picamera.py | 0 7 files changed, 652 insertions(+), 29 deletions(-) create mode 100644 camera_node.py create mode 100644 client_message.py create mode 100644 motor_node.py create mode 100644 server_message.py rename detect_picamera.py => test_picamera.py (100%) diff --git a/Camera.py b/Camera.py index 92847f3..d49c2ed 100644 --- a/Camera.py +++ b/Camera.py @@ -121,7 +121,7 @@ def annotate_objects(self, annotator, results, labels): annotator.text([xmin, ymin], '%s\n%.2f' % (labels[obj['class_id']], obj['score'])) - def detect_size(self, results, labels, obj_label): + def detect_sizes(self, results, labels): sizes = [] for obj in results: ymin, xmin, ymax, xmax = obj['bounding_box'] @@ -129,19 +129,19 @@ def detect_size(self, results, labels, obj_label): xmax = int(xmax * self.CAMERA_WIDTH) ymin = int(ymin * self.CAMERA_HEIGHT) ymax = int(ymax * self.CAMERA_HEIGHT) - if labels[obj['class_id']] == obj_label: - obj = {} - obj['height'] = xmax - xmin - obj['width'] = ymax - ymin - # 55mm width, 80mm height - obj['pixel_metric'] = (obj['width'] / 55 + obj['height'] / 80) / 2 - print("Pixel metrics: " + - str(round(obj['pixel_metric'], 1)) + "\n") - sizes.append(obj) + size_obj = {} + size_obj['name'] = labels[obj['class_id']] + size_obj['height'] = xmax - xmin + size_obj['width'] = ymax - ymin + # 55mm width, 80mm height + size_obj['pixel_metric'] = (size_obj['width'] / 55 + size_obj['height'] / 80) / 2 + print("Pixel metrics: " + + str(round(size_obj['pixel_metric'], 1)) + "\n") + sizes.append(size_obj) return sizes - def detect_distance(self, results, labels, obj_label): + def detect_distances(self, results, labels): distances = [] for obj in results: ymin, xmin, ymax, xmax = obj['bounding_box'] @@ -149,16 +149,16 @@ def detect_distance(self, results, labels, obj_label): xmax = int(xmax * self.CAMERA_WIDTH) ymin = int(ymin * self.CAMERA_HEIGHT) ymax = int(ymax * self.CAMERA_HEIGHT) - if labels[obj['class_id']] == obj_label: - obj = {} - obj['height'] = xmax - xmin - obj['width'] = ymax - ymin - # When pixel metric 2.1 distance will 155mm - obj['focal_distance'] = ( - (obj['width'] * 155) / 55 + obj['height'] * 155 / 80) / 2 - print("Focal distance: " + - str(round(obj['focal_distance'], 1)) + "\n") - distances.append(obj) + dist_obj = {} + dist_obj['name'] = labels[obj['class_id']] + dist_obj['height'] = xmax - xmin + dist_obj['width'] = ymax - ymin + # When pixel metric 2.1 distance will 155mm + dist_obj['focal_distance'] = ( + (dist_obj['width'] * 155) / 55 + dist_obj['height'] * 155 / 80) / 2 + print("Focal distance: " + + str(round(dist_obj['focal_distance'], 1)) + "\n") + distances.append(dist_obj) return distances def print_objects(self, results, labels): @@ -189,18 +189,20 @@ def execute_command(self): for interpreter in self.interpreters: image = image.resize((interpreter['shape'][1], interpreter['shape'][2]), Image.ANTIALIAS) - result = self.detect_objects(interpreter['interpreter'], image, 0.5) + results = self.detect_objects(interpreter['interpreter'], image, 0.5) # Annotate objects in terminal - self.print_objects(result, interpreter['labels']) + self.print_objects(results, interpreter['labels']) # Annotate object in view # self.annotate_objects(annotator, result, interpreter['labels']) - # Detect size and distance TODO: improve with contanstant object - size = self.detect_size( - result, interpreter['labels'], interpreter['name']) - distance = self.detect_distance( - result, interpreter['labels'], interpreter['name']) + # Detect size and distance + # TODO: improve with physical object with 1cm length + sizes = self.detect_sizes( + results, interpreter['labels'], interpreter['name']) + distances = self.detect_distances( + results, interpreter['labels'], interpreter['name']) if bool(interpreter.get('function')): - interpreter['function'](result, interpreter['labels'], size, distance) + interpreter['function']( + results, interpreter['labels'], sizes, distances) elapsed_ms = (time.monotonic() - start_time) * 1000 diff --git a/Excavator.py b/Excavator.py index f928c84..4f6d092 100644 --- a/Excavator.py +++ b/Excavator.py @@ -57,6 +57,14 @@ def backward_right_chain(self, speed=100): clockwise=False, speed=speed) + def move_forward(self, speed=100): + self.forward_left_chain(speed) + self.forward_right_chain(speed) + + def move_backward(self, speed=100): + self.backward_left_chain(speed) + self.backward_right_chain(speed) + def turn_left_body(self, speed=100): self.motors_memo.append(self.BODY_MOTOR) self.motors.run_dc_motor(self.BODY_MOTOR, clockwise=True, speed=speed) diff --git a/camera_node.py b/camera_node.py new file mode 100644 index 0000000..0bd6069 --- /dev/null +++ b/camera_node.py @@ -0,0 +1,114 @@ +#!/usr/bin/env python3 + +import sys +import socket +import selectors +import traceback +from Camera import Camera + +from client_message import Message + +class CameraNode: + def __init__(self) -> None: + self.sel = selectors.DefaultSelector() + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + try: + self.sel.close() + except RuntimeWarning: + return True + + def create_request(self, action, value, encode): + if encode == "bin": + return dict( + type="binary/custom-client-binary-type", + encoding="binary", + content=bytes(action + value, encoding="utf-8"), + ) + else: + return dict( + type="text/json", + encoding="utf-8", + content=dict(action=action, value=value), + ) + + def start_connection(self): + addr = ('127.0.0.1', '65432') + print("starting connection to", addr) + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.setblocking(False) + sock.connect_ex(addr) + return sock, addr + + def send_instruction(self, addr, sock, request): + events = selectors.EVENT_READ | selectors.EVENT_WRITE + message = Message(self.sel, sock, addr, request) + self.sel.register(sock, events, data=message) + try: + while True: + events = self.sel.select(timeout=1) + for key, mask in events: + message = key.data + try: + message.process_events(mask) + except Exception: + print( + "main: error: exception for", + f"{message.addr}:\n{traceback.format_exc()}", + ) + message.close() + # Check for a socket being monitored to continue. + if not self.sel.get_map(): + break + except KeyboardInterrupt: + print("caught keyboard interrupt, exiting") + finally: + self.sel.close() + + +cnode = CameraNode() +sock, addr = cnode.start_connection() + +def find_object(obj_name, results, labels, sizes, distances): + score = 0 + + for obj in results: + if labels[obj['class_id']] == obj_name: + score = obj['score'] + + obj_size = next((size for size in sizes if size["name"] == obj_name), None) + obj_dist = next((dist for dist in distances if dist["name"] == obj_name), None) + + while score < 0.5: + request = cnode.create_request("turn-left", "360") + cnode.send_instruction(sock, addr, request) + + request = cnode.create_request("stop", "all") + cnode.send_instruction(sock, addr, request) + + while obj_dist > 100: + request = cnode.create_request("move-forward", "1") + cnode.send_instruction(sock, addr, request) + + request = cnode.create_request("stop", "all") + cnode.send_instruction(sock, addr, request) + +tl_models = [ + { + 'name': 'shovel', + 'model_path': './trained_model/shovel_model/model.tflite', + 'label_path': './trained_model/shovel_model/model-dict.txt', + 'function': None + }, + { + 'name': 'apple', + 'model_path': './trained_model/object/detect.tflite', + 'label_path': './trained_model/object/coco_labels.txt', + 'function': find_object + } +] +camera = Camera(tl_models) +camera.execute_command() diff --git a/client_message.py b/client_message.py new file mode 100644 index 0000000..9262417 --- /dev/null +++ b/client_message.py @@ -0,0 +1,208 @@ +import sys +import selectors +import json +import io +import struct + + +class Message: + def __init__(self, selector, sock, addr, request): + self.selector = selector + self.sock = sock + self.addr = addr + self.request = request + self._recv_buffer = b"" + self._send_buffer = b"" + self._request_queued = False + self._jsonheader_len = None + self.jsonheader = None + self.response = None + + def _set_selector_events_mask(self, mode): + """Set selector to listen for events: mode is 'r', 'w', or 'rw'.""" + if mode == "r": + events = selectors.EVENT_READ + elif mode == "w": + events = selectors.EVENT_WRITE + elif mode == "rw": + events = selectors.EVENT_READ | selectors.EVENT_WRITE + else: + raise ValueError(f"Invalid events mask mode {repr(mode)}.") + self.selector.modify(self.sock, events, data=self) + + def _read(self): + try: + # Should be ready to read + data = self.sock.recv(4096) + except BlockingIOError: + # Resource temporarily unavailable (errno EWOULDBLOCK) + pass + else: + if data: + self._recv_buffer += data + else: + raise RuntimeError("Peer closed.") + + def _write(self): + if self._send_buffer: + print("sending", repr(self._send_buffer), "to", self.addr) + try: + # Should be ready to write + sent = self.sock.send(self._send_buffer) + except BlockingIOError: + # Resource temporarily unavailable (errno EWOULDBLOCK) + pass + else: + self._send_buffer = self._send_buffer[sent:] + + def _json_encode(self, obj, encoding): + return json.dumps(obj, ensure_ascii=False).encode(encoding) + + def _json_decode(self, json_bytes, encoding): + tiow = io.TextIOWrapper( + io.BytesIO(json_bytes), encoding=encoding, newline="" + ) + obj = json.load(tiow) + tiow.close() + return obj + + def _create_message( + self, *, content_bytes, content_type, content_encoding + ): + jsonheader = { + "byteorder": sys.byteorder, + "content-type": content_type, + "content-encoding": content_encoding, + "content-length": len(content_bytes), + } + jsonheader_bytes = self._json_encode(jsonheader, "utf-8") + message_hdr = struct.pack(">H", len(jsonheader_bytes)) + message = message_hdr + jsonheader_bytes + content_bytes + return message + + def _process_response_json_content(self): + content = self.response + result = content.get("result") + print(f"got result: {result}") + + def _process_response_binary_content(self): + content = self.response + print(f"got response: {repr(content)}") + + def process_events(self, mask): + if mask & selectors.EVENT_READ: + self.read() + if mask & selectors.EVENT_WRITE: + self.write() + + def read(self): + self._read() + + if self._jsonheader_len is None: + self.process_protoheader() + + if self._jsonheader_len is not None: + if self.jsonheader is None: + self.process_jsonheader() + + if self.jsonheader: + if self.response is None: + self.process_response() + + def write(self): + if not self._request_queued: + self.queue_request() + + self._write() + + if self._request_queued: + if not self._send_buffer: + # Set selector to listen for read events, we're done writing. + self._set_selector_events_mask("r") + + def close(self): + print("closing connection to", self.addr) + try: + self.selector.unregister(self.sock) + except Exception as e: + print( + "error: selector.unregister() exception for", + f"{self.addr}: {repr(e)}", + ) + + try: + self.sock.close() + except OSError as e: + print( + "error: socket.close() exception for", + f"{self.addr}: {repr(e)}", + ) + finally: + # Delete reference to socket object for garbage collection + self.sock = None + + def queue_request(self): + content = self.request["content"] + content_type = self.request["type"] + content_encoding = self.request["encoding"] + if content_type == "text/json": + req = { + "content_bytes": self._json_encode(content, content_encoding), + "content_type": content_type, + "content_encoding": content_encoding, + } + else: + req = { + "content_bytes": content, + "content_type": content_type, + "content_encoding": content_encoding, + } + message = self._create_message(**req) + self._send_buffer += message + self._request_queued = True + + def process_protoheader(self): + hdrlen = 2 + if len(self._recv_buffer) >= hdrlen: + self._jsonheader_len = struct.unpack( + ">H", self._recv_buffer[:hdrlen] + )[0] + self._recv_buffer = self._recv_buffer[hdrlen:] + + def process_jsonheader(self): + hdrlen = self._jsonheader_len + if len(self._recv_buffer) >= hdrlen: + self.jsonheader = self._json_decode( + self._recv_buffer[:hdrlen], "utf-8" + ) + self._recv_buffer = self._recv_buffer[hdrlen:] + for reqhdr in ( + "byteorder", + "content-length", + "content-type", + "content-encoding", + ): + if reqhdr not in self.jsonheader: + raise ValueError(f'Missing required header "{reqhdr}".') + + def process_response(self): + content_len = self.jsonheader["content-length"] + if not len(self._recv_buffer) >= content_len: + return + data = self._recv_buffer[:content_len] + self._recv_buffer = self._recv_buffer[content_len:] + if self.jsonheader["content-type"] == "text/json": + encoding = self.jsonheader["content-encoding"] + self.response = self._json_decode(data, encoding) + print("received response", repr(self.response), "from", self.addr) + self._process_response_json_content() + else: + # Binary or unknown content-type + self.response = data + print( + f'received {self.jsonheader["content-type"]} response from', + self.addr, + ) + self._process_response_binary_content() + # Close when response has been processed + self.close() diff --git a/motor_node.py b/motor_node.py new file mode 100644 index 0000000..bd87295 --- /dev/null +++ b/motor_node.py @@ -0,0 +1,79 @@ +#!/usr/bin/env python3 + +import sys +import socket +import selectors +import traceback + +from server_message import Message +from Excavator import Excavator + + +excavator = Excavator() + +instructions = { + "forward": {'cmd': excavator.move_forward, 'fire': excavator.execute, 'period': 1}, + "backward": {'cmd': excavator.move_forward, 'fire': excavator.execute, 'period': 1}, + "left": {'cmd': excavator.forward_left_chain, 'fire': excavator.execute, 'period': 1}, + "right": {'cmd': excavator.forward_right_chain, 'fire': excavator.execute, 'period': 1}, + "shovel-left": {'cmd': excavator.turn_left_body, 'fire': excavator.execute, 'period': 1}, + "shovel-right": {'cmd': excavator.turn_right_body, 'fire': excavator.execute, 'period': 1}, + "shovel-up": {'cmd': excavator.move_up_shovel, 'fire': excavator.execute, 'period': 1}, + "shovel-down": {'cmd': excavator.move_down_shovel, 'fire': excavator.execute, 'period': 1}, +} + +class MotorNode: + def __init__(self) -> None: + self.sel = selectors.DefaultSelector() + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + try: + self.sel.close() + except RuntimeWarning: + return True + + def accept_wrapper(self, sock): + conn, addr = sock.accept() # Should be ready to read + print("accepted connection from", addr) + conn.setblocking(False) + message = Message(self.sel, conn, addr, instructions) + self.sel.register(conn, selectors.EVENT_READ, data=message) + + host = '127.0.0.1' + port = '65432' + lsock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + # Avoid bind() exception: OSError: [Errno 48] Address already in use + lsock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + lsock.bind((host, port)) + lsock.listen() + print("listening on", (host, port)) + lsock.setblocking(False) + self.sel.register(lsock, selectors.EVENT_READ, data=None) + + def follow_instructions(self): + try: + while True: + events = self.sel.select(timeout=None) + for key, mask in events: + if key.data is None: + self.accept_wrapper(key.fileobj) + else: + message = key.data + try: + message.process_events(mask) + except Exception: + print( + "main: error: exception for", + f"{message.addr}:\n{traceback.format_exc()}", + ) + message.close() + except KeyboardInterrupt: + print("caught keyboard interrupt, exiting") + finally: + self.sel.close() + + +MotorNode().follow_instructions() \ No newline at end of file diff --git a/server_message.py b/server_message.py new file mode 100644 index 0000000..3d40569 --- /dev/null +++ b/server_message.py @@ -0,0 +1,212 @@ +import sys +import selectors +import json +import io +import struct + +class Message: + def __init__(self, selector, sock, addr, instructions): + self.selector = selector + self.sock = sock + self.addr = addr + self._recv_buffer = b"" + self._send_buffer = b"" + self._jsonheader_len = None + self.jsonheader = None + self.request = None + self.response_created = False + self.instructions = instructions + + def _set_selector_events_mask(self, mode): + """Set selector to listen for events: mode is 'r', 'w', or 'rw'.""" + if mode == "r": + events = selectors.EVENT_READ + elif mode == "w": + events = selectors.EVENT_WRITE + elif mode == "rw": + events = selectors.EVENT_READ | selectors.EVENT_WRITE + else: + raise ValueError(f"Invalid events mask mode {repr(mode)}.") + self.selector.modify(self.sock, events, data=self) + + def _read(self): + try: + # Should be ready to read + data = self.sock.recv(4096) + except BlockingIOError: + # Resource temporarily unavailable (errno EWOULDBLOCK) + pass + else: + if data: + self._recv_buffer += data + else: + raise RuntimeError("Peer closed.") + + def _write(self): + if self._send_buffer: + print("sending", repr(self._send_buffer), "to", self.addr) + try: + # Should be ready to write + sent = self.sock.send(self._send_buffer) + except BlockingIOError: + # Resource temporarily unavailable (errno EWOULDBLOCK) + pass + else: + self._send_buffer = self._send_buffer[sent:] + # Close when the buffer is drained. The response has been sent. + if sent and not self._send_buffer: + self.close() + + def _json_encode(self, obj, encoding): + return json.dumps(obj, ensure_ascii=False).encode(encoding) + + def _json_decode(self, json_bytes, encoding): + tiow = io.TextIOWrapper( + io.BytesIO(json_bytes), encoding=encoding, newline="" + ) + obj = json.load(tiow) + tiow.close() + return obj + + def _create_message( + self, *, content_bytes, content_type, content_encoding + ): + jsonheader = { + "byteorder": sys.byteorder, + "content-type": content_type, + "content-encoding": content_encoding, + "content-length": len(content_bytes), + } + jsonheader_bytes = self._json_encode(jsonheader, "utf-8") + message_hdr = struct.pack(">H", len(jsonheader_bytes)) + message = message_hdr + jsonheader_bytes + content_bytes + return message + + def _create_response_json_content(self): + action = self.request.get("action") + query = self.request.get("value") + if bool(self.instructions.get(query)): + self.instructions.get(query)['cmd']() + self.instructions.get(query)['fire']( + int(self.instructions.get(query)['period'])) + content = {"result": action} + else: + content = {"result": f'Error: invalid action "{action}".'} + content_encoding = "utf-8" + response = { + "content_bytes": self._json_encode(content, content_encoding), + "content_type": "text/json", + "content_encoding": content_encoding, + } + return response + + def _create_response_binary_content(self): + response = { + "content_bytes": b"First 10 bytes of request: " + + self.request[:10], + "content_type": "binary/custom-server-binary-type", + "content_encoding": "binary", + } + return response + + def process_events(self, mask): + if mask & selectors.EVENT_READ: + self.read() + if mask & selectors.EVENT_WRITE: + self.write() + + def read(self): + self._read() + + if self._jsonheader_len is None: + self.process_protoheader() + + if self._jsonheader_len is not None: + if self.jsonheader is None: + self.process_jsonheader() + + if self.jsonheader: + if self.request is None: + self.process_request() + + def write(self): + if self.request: + if not self.response_created: + self.create_response() + + self._write() + + def close(self): + print("closing connection to", self.addr) + try: + self.selector.unregister(self.sock) + except Exception as e: + print( + "error: selector.unregister() exception for", + f"{self.addr}: {repr(e)}", + ) + + try: + self.sock.close() + except OSError as e: + print( + "error: socket.close() exception for", + f"{self.addr}: {repr(e)}", + ) + finally: + # Delete reference to socket object for garbage collection + self.sock = None + + def process_protoheader(self): + hdrlen = 2 + if len(self._recv_buffer) >= hdrlen: + self._jsonheader_len = struct.unpack( + ">H", self._recv_buffer[:hdrlen] + )[0] + self._recv_buffer = self._recv_buffer[hdrlen:] + + def process_jsonheader(self): + hdrlen = self._jsonheader_len + if len(self._recv_buffer) >= hdrlen: + self.jsonheader = self._json_decode( + self._recv_buffer[:hdrlen], "utf-8" + ) + self._recv_buffer = self._recv_buffer[hdrlen:] + for reqhdr in ( + "byteorder", + "content-length", + "content-type", + "content-encoding", + ): + if reqhdr not in self.jsonheader: + raise ValueError(f'Missing required header "{reqhdr}".') + + def process_request(self): + content_len = self.jsonheader["content-length"] + if not len(self._recv_buffer) >= content_len: + return + data = self._recv_buffer[:content_len] + self._recv_buffer = self._recv_buffer[content_len:] + if self.jsonheader["content-type"] == "text/json": + encoding = self.jsonheader["content-encoding"] + self.request = self._json_decode(data, encoding) + print("received request", repr(self.request), "from", self.addr) + else: + # Binary or unknown content-type + self.request = data + print( + f'received {self.jsonheader["content-type"]} request from', + self.addr, + ) + # Set selector to listen for write events, we're done reading. + self._set_selector_events_mask("w") + + def create_response(self): + if self.jsonheader["content-type"] == "text/json": + response = self._create_response_json_content() + else: + # Binary or unknown content-type + response = self._create_response_binary_content() + message = self._create_message(**response) + self.response_created = True + self._send_buffer += message diff --git a/detect_picamera.py b/test_picamera.py similarity index 100% rename from detect_picamera.py rename to test_picamera.py From 653ac46db32d6ccc541af0ff91ebd7548cc70b8d Mon Sep 17 00:00:00 2001 From: zogoo Date: Sun, 31 Oct 2021 18:50:17 +0900 Subject: [PATCH 02/14] Cleanup --- camera_node.py | 1 + client_message.py | 1 - motor_node.py | 2 +- 3 files changed, 2 insertions(+), 2 deletions(-) diff --git a/camera_node.py b/camera_node.py index 0bd6069..6bbdb31 100644 --- a/camera_node.py +++ b/camera_node.py @@ -110,5 +110,6 @@ def find_object(obj_name, results, labels, sizes, distances): 'function': find_object } ] + camera = Camera(tl_models) camera.execute_command() diff --git a/client_message.py b/client_message.py index 9262417..db07305 100644 --- a/client_message.py +++ b/client_message.py @@ -4,7 +4,6 @@ import io import struct - class Message: def __init__(self, selector, sock, addr, request): self.selector = selector diff --git a/motor_node.py b/motor_node.py index bd87295..c2a6bbf 100644 --- a/motor_node.py +++ b/motor_node.py @@ -76,4 +76,4 @@ def follow_instructions(self): self.sel.close() -MotorNode().follow_instructions() \ No newline at end of file +MotorNode().follow_instructions() From c077ab55c61b6b816fcc478c21bb844da442bf71 Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 19:47:45 +0900 Subject: [PATCH 03/14] Update structure for motor --- Excavator.py | 6 ++++++ camera_node.py | 10 ++++++---- motor_node.py | 17 +++++++++-------- server_message.py | 7 +++---- 4 files changed, 24 insertions(+), 16 deletions(-) diff --git a/Excavator.py b/Excavator.py index 4f6d092..be5e0ab 100644 --- a/Excavator.py +++ b/Excavator.py @@ -85,6 +85,12 @@ def move_down_shovel(self, speed=100): clockwise=False, speed=speed) + def stop_all_motors(self): + for motor in [self.LEFT_CHAIN_MOTOR, self.RIGHT_CHAIN_MOTOR, self.BODY_MOTOR, self.SHOVEL_MOTOR]: + self.motors.stop_dc_motor(motor) + self.motors_memo = [] + time.sleep(1) + def test_move(self): self.forward_left_chain() self.forward_right_chain() diff --git a/camera_node.py b/camera_node.py index 6bbdb31..66c3d4a 100644 --- a/camera_node.py +++ b/camera_node.py @@ -82,15 +82,17 @@ def find_object(obj_name, results, labels, sizes, distances): obj_size = next((size for size in sizes if size["name"] == obj_name), None) obj_dist = next((dist for dist in distances if dist["name"] == obj_name), None) - while score < 0.5: - request = cnode.create_request("turn-left", "360") + while score < 0.5: + request = cnode.create_request("left", "4") cnode.send_instruction(sock, addr, request) - + request = cnode.create_request("right", "4") + cnode.send_instruction(sock, addr, request) + request = cnode.create_request("stop", "all") cnode.send_instruction(sock, addr, request) while obj_dist > 100: - request = cnode.create_request("move-forward", "1") + request = cnode.create_request("forward", "1") cnode.send_instruction(sock, addr, request) request = cnode.create_request("stop", "all") diff --git a/motor_node.py b/motor_node.py index c2a6bbf..833a0e1 100644 --- a/motor_node.py +++ b/motor_node.py @@ -12,14 +12,15 @@ excavator = Excavator() instructions = { - "forward": {'cmd': excavator.move_forward, 'fire': excavator.execute, 'period': 1}, - "backward": {'cmd': excavator.move_forward, 'fire': excavator.execute, 'period': 1}, - "left": {'cmd': excavator.forward_left_chain, 'fire': excavator.execute, 'period': 1}, - "right": {'cmd': excavator.forward_right_chain, 'fire': excavator.execute, 'period': 1}, - "shovel-left": {'cmd': excavator.turn_left_body, 'fire': excavator.execute, 'period': 1}, - "shovel-right": {'cmd': excavator.turn_right_body, 'fire': excavator.execute, 'period': 1}, - "shovel-up": {'cmd': excavator.move_up_shovel, 'fire': excavator.execute, 'period': 1}, - "shovel-down": {'cmd': excavator.move_down_shovel, 'fire': excavator.execute, 'period': 1}, + "forward": {'cmd': excavator.move_forward, 'fire': excavator.execute}, + "backward": {'cmd': excavator.move_forward, 'fire': excavator.execute}, + "left": {'cmd': excavator.forward_left_chain, 'fire': excavator.execute}, + "right": {'cmd': excavator.forward_right_chain, 'fire': excavator.execute}, + "shovel-left": {'cmd': excavator.turn_left_body, 'fire': excavator.execute}, + "shovel-right": {'cmd': excavator.turn_right_body, 'fire': excavator.execute}, + "shovel-up": {'cmd': excavator.move_up_shovel, 'fire': excavator.execute}, + "shovel-down": {'cmd': excavator.move_down_shovel, 'fire': excavator.execute}, + "stop": { 'cmd': excavator.stop_all_motors } } class MotorNode: diff --git a/server_message.py b/server_message.py index 3d40569..dda7d6d 100644 --- a/server_message.py +++ b/server_message.py @@ -85,10 +85,9 @@ def _create_message( def _create_response_json_content(self): action = self.request.get("action") query = self.request.get("value") - if bool(self.instructions.get(query)): - self.instructions.get(query)['cmd']() - self.instructions.get(query)['fire']( - int(self.instructions.get(query)['period'])) + if bool(self.instructions.get(action)): + self.instructions.get(action)['cmd']() + self.instructions.get(action)['fire'](int(query)) content = {"result": action} else: content = {"result": f'Error: invalid action "{action}".'} From bf3c779361c05950e12eb243928c551a941017ec Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 20:11:26 +0900 Subject: [PATCH 04/14] Port must be integer --- Excavator.py | 3 +-- camera_node.py | 2 +- motor_node.py | 2 +- 3 files changed, 3 insertions(+), 4 deletions(-) diff --git a/Excavator.py b/Excavator.py index be5e0ab..af57fc9 100644 --- a/Excavator.py +++ b/Excavator.py @@ -86,8 +86,7 @@ def move_down_shovel(self, speed=100): speed=speed) def stop_all_motors(self): - for motor in [self.LEFT_CHAIN_MOTOR, self.RIGHT_CHAIN_MOTOR, self.BODY_MOTOR, self.SHOVEL_MOTOR]: - self.motors.stop_dc_motor(motor) + self.motors.stop_dc_motors([self.LEFT_CHAIN_MOTOR, self.RIGHT_CHAIN_MOTOR, self.BODY_MOTOR, self.SHOVEL_MOTOR]) self.motors_memo = [] time.sleep(1) diff --git a/camera_node.py b/camera_node.py index 66c3d4a..b14228f 100644 --- a/camera_node.py +++ b/camera_node.py @@ -36,7 +36,7 @@ def create_request(self, action, value, encode): ) def start_connection(self): - addr = ('127.0.0.1', '65432') + addr = ('127.0.0.1', 65432) print("starting connection to", addr) sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.setblocking(False) diff --git a/motor_node.py b/motor_node.py index 833a0e1..2b87806 100644 --- a/motor_node.py +++ b/motor_node.py @@ -44,7 +44,7 @@ def accept_wrapper(self, sock): self.sel.register(conn, selectors.EVENT_READ, data=message) host = '127.0.0.1' - port = '65432' + port = 65432 lsock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) # Avoid bind() exception: OSError: [Errno 48] Address already in use lsock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) From 6f1c1a0276ebc93a4ef9c2be04b864c0fec12876 Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 20:15:14 +0900 Subject: [PATCH 05/14] Argument error --- Camera.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/Camera.py b/Camera.py index d49c2ed..d643dbb 100644 --- a/Camera.py +++ b/Camera.py @@ -196,10 +196,8 @@ def execute_command(self): # self.annotate_objects(annotator, result, interpreter['labels']) # Detect size and distance # TODO: improve with physical object with 1cm length - sizes = self.detect_sizes( - results, interpreter['labels'], interpreter['name']) - distances = self.detect_distances( - results, interpreter['labels'], interpreter['name']) + sizes = self.detect_sizes(results, interpreter['labels']) + distances = self.detect_distances(results, interpreter['labels']) if bool(interpreter.get('function')): interpreter['function']( results, interpreter['labels'], sizes, distances) From 5fb2f58be9a7965c418eaec3c70355d65d5a8a71 Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 20:21:51 +0900 Subject: [PATCH 06/14] Argument errors --- Camera.py | 2 +- camera_node.py | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/Camera.py b/Camera.py index d643dbb..24e181a 100644 --- a/Camera.py +++ b/Camera.py @@ -200,7 +200,7 @@ def execute_command(self): distances = self.detect_distances(results, interpreter['labels']) if bool(interpreter.get('function')): interpreter['function']( - results, interpreter['labels'], sizes, distances) + results, interpreter['labels'], sizes, distances, interpreter['name']) elapsed_ms = (time.monotonic() - start_time) * 1000 diff --git a/camera_node.py b/camera_node.py index b14228f..996708c 100644 --- a/camera_node.py +++ b/camera_node.py @@ -72,7 +72,8 @@ def send_instruction(self, addr, sock, request): cnode = CameraNode() sock, addr = cnode.start_connection() -def find_object(obj_name, results, labels, sizes, distances): + +def find_object(results, labels, sizes, distances, obj_name): score = 0 for obj in results: From 60e8990c91c55fd39f9b899b8c5aa3477194348a Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 20:23:42 +0900 Subject: [PATCH 07/14] Another argument error --- camera_node.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/camera_node.py b/camera_node.py index 996708c..60516e2 100644 --- a/camera_node.py +++ b/camera_node.py @@ -21,7 +21,7 @@ def __exit__(self, exc_type, exc_val, exc_tb): except RuntimeWarning: return True - def create_request(self, action, value, encode): + def create_request(self, action, value, encode="utf-8"): if encode == "bin": return dict( type="binary/custom-client-binary-type", From 5bade97b19e94c60e053d44496a5cdc0f238a415 Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 21:26:09 +0900 Subject: [PATCH 08/14] Another argument issue --- camera_node.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/camera_node.py b/camera_node.py index 60516e2..3694bbc 100644 --- a/camera_node.py +++ b/camera_node.py @@ -43,7 +43,7 @@ def start_connection(self): sock.connect_ex(addr) return sock, addr - def send_instruction(self, addr, sock, request): + def send_instruction(self, sock, addr, request): events = selectors.EVENT_READ | selectors.EVENT_WRITE message = Message(self.sel, sock, addr, request) self.sel.register(sock, events, data=message) From 055ac51ae2722600ecbf3cff17e9443cfe45a703 Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 22:31:07 +0900 Subject: [PATCH 09/14] No need miltiple registration for multiplexer --- camera_node.py | 26 ++++++++++++++++---------- client_message.py | 6 +++++- server_message.py | 4 +++- 3 files changed, 24 insertions(+), 12 deletions(-) diff --git a/camera_node.py b/camera_node.py index 3694bbc..cd3f61b 100644 --- a/camera_node.py +++ b/camera_node.py @@ -41,18 +41,24 @@ def start_connection(self): sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.setblocking(False) sock.connect_ex(addr) - return sock, addr - - def send_instruction(self, sock, addr, request): events = selectors.EVENT_READ | selectors.EVENT_WRITE - message = Message(self.sel, sock, addr, request) + handshake_req = dict( + type="text/json", + encoding="utf-8", + content=dict(action="hello", value="motor"), + ) + message = Message(self.sel, sock, addr, handshake_req) self.sel.register(sock, events, data=message) + return sock, addr + + def send_instruction(self, request): try: while True: events = self.sel.select(timeout=1) for key, mask in events: message = key.data try: + message.set_request(request) message.process_events(mask) except Exception: print( @@ -70,7 +76,7 @@ def send_instruction(self, sock, addr, request): cnode = CameraNode() -sock, addr = cnode.start_connection() +cnode.start_connection() def find_object(results, labels, sizes, distances, obj_name): @@ -85,19 +91,19 @@ def find_object(results, labels, sizes, distances, obj_name): while score < 0.5: request = cnode.create_request("left", "4") - cnode.send_instruction(sock, addr, request) + cnode.send_instruction(request) request = cnode.create_request("right", "4") - cnode.send_instruction(sock, addr, request) + cnode.send_instruction(request) request = cnode.create_request("stop", "all") - cnode.send_instruction(sock, addr, request) + cnode.send_instruction(request) while obj_dist > 100: request = cnode.create_request("forward", "1") - cnode.send_instruction(sock, addr, request) + cnode.send_instruction(request) request = cnode.create_request("stop", "all") - cnode.send_instruction(sock, addr, request) + cnode.send_instruction(request) tl_models = [ { diff --git a/client_message.py b/client_message.py index db07305..a28d6f2 100644 --- a/client_message.py +++ b/client_message.py @@ -9,9 +9,9 @@ def __init__(self, selector, sock, addr, request): self.selector = selector self.sock = sock self.addr = addr - self.request = request self._recv_buffer = b"" self._send_buffer = b"" + self.request = request self._request_queued = False self._jsonheader_len = None self.jsonheader = None @@ -140,6 +140,10 @@ def close(self): # Delete reference to socket object for garbage collection self.sock = None + def set_request(self, request): + self.request = request + self._request_queued = False + def queue_request(self): content = self.request["content"] content_type = self.request["type"] diff --git a/server_message.py b/server_message.py index dda7d6d..f4604ba 100644 --- a/server_message.py +++ b/server_message.py @@ -85,7 +85,9 @@ def _create_message( def _create_response_json_content(self): action = self.request.get("action") query = self.request.get("value") - if bool(self.instructions.get(action)): + if action == "hello": + content = {"result": "hello camera"} + elif bool(self.instructions.get(action)): self.instructions.get(action)['cmd']() self.instructions.get(action)['fire'](int(query)) content = {"result": action} From 4961c8bfcacc877fc0af7c2b825fd20f23e5c575 Mon Sep 17 00:00:00 2001 From: zogoo Date: Mon, 1 Nov 2021 22:41:15 +0900 Subject: [PATCH 10/14] Server started not being executed --- motor_node.py | 41 +++++++++++++++++++++-------------------- 1 file changed, 21 insertions(+), 20 deletions(-) diff --git a/motor_node.py b/motor_node.py index 2b87806..76a9abf 100644 --- a/motor_node.py +++ b/motor_node.py @@ -8,21 +8,6 @@ from server_message import Message from Excavator import Excavator - -excavator = Excavator() - -instructions = { - "forward": {'cmd': excavator.move_forward, 'fire': excavator.execute}, - "backward": {'cmd': excavator.move_forward, 'fire': excavator.execute}, - "left": {'cmd': excavator.forward_left_chain, 'fire': excavator.execute}, - "right": {'cmd': excavator.forward_right_chain, 'fire': excavator.execute}, - "shovel-left": {'cmd': excavator.turn_left_body, 'fire': excavator.execute}, - "shovel-right": {'cmd': excavator.turn_right_body, 'fire': excavator.execute}, - "shovel-up": {'cmd': excavator.move_up_shovel, 'fire': excavator.execute}, - "shovel-down": {'cmd': excavator.move_down_shovel, 'fire': excavator.execute}, - "stop": { 'cmd': excavator.stop_all_motors } -} - class MotorNode: def __init__(self) -> None: self.sel = selectors.DefaultSelector() @@ -43,14 +28,14 @@ def accept_wrapper(self, sock): message = Message(self.sel, conn, addr, instructions) self.sel.register(conn, selectors.EVENT_READ, data=message) - host = '127.0.0.1' - port = 65432 + def init_listener(self): + addr = ('127.0.0.1', 65432) lsock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) # Avoid bind() exception: OSError: [Errno 48] Address already in use lsock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - lsock.bind((host, port)) + lsock.bind(addr) lsock.listen() - print("listening on", (host, port)) + print("listening on", addr) lsock.setblocking(False) self.sel.register(lsock, selectors.EVENT_READ, data=None) @@ -60,6 +45,7 @@ def follow_instructions(self): events = self.sel.select(timeout=None) for key, mask in events: if key.data is None: + print("Accept wrapper is executed") self.accept_wrapper(key.fileobj) else: message = key.data @@ -77,4 +63,19 @@ def follow_instructions(self): self.sel.close() -MotorNode().follow_instructions() +excavator = Excavator() +instructions = { + "forward": {'cmd': excavator.move_forward, 'fire': excavator.execute}, + "backward": {'cmd': excavator.move_forward, 'fire': excavator.execute}, + "left": {'cmd': excavator.forward_left_chain, 'fire': excavator.execute}, + "right": {'cmd': excavator.forward_right_chain, 'fire': excavator.execute}, + "shovel-left": {'cmd': excavator.turn_left_body, 'fire': excavator.execute}, + "shovel-right": {'cmd': excavator.turn_right_body, 'fire': excavator.execute}, + "shovel-up": {'cmd': excavator.move_up_shovel, 'fire': excavator.execute}, + "shovel-down": {'cmd': excavator.move_down_shovel, 'fire': excavator.execute}, + "stop": {'cmd': excavator.stop_all_motors} +} + +motor_node = MotorNode() +motor_node.init_listener() +motor_node.follow_instructions() From 32e26790c3efbb0131ba2740f8ac36c54b6cae94 Mon Sep 17 00:00:00 2001 From: zogoo Date: Tue, 2 Nov 2021 15:56:37 +0900 Subject: [PATCH 11/14] To see result, let's do it with non efficient way --- camera_node.py | 33 ++++++++++++++------------------- 1 file changed, 14 insertions(+), 19 deletions(-) diff --git a/camera_node.py b/camera_node.py index cd3f61b..2905aba 100644 --- a/camera_node.py +++ b/camera_node.py @@ -1,5 +1,6 @@ #!/usr/bin/env python3 +from _typeshed import Self import sys import socket import selectors @@ -35,30 +36,24 @@ def create_request(self, action, value, encode="utf-8"): content=dict(action=action, value=value), ) - def start_connection(self): + def start_connection(self, request): addr = ('127.0.0.1', 65432) print("starting connection to", addr) sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.setblocking(False) sock.connect_ex(addr) events = selectors.EVENT_READ | selectors.EVENT_WRITE - handshake_req = dict( - type="text/json", - encoding="utf-8", - content=dict(action="hello", value="motor"), - ) - message = Message(self.sel, sock, addr, handshake_req) + message = Message(self.sel, sock, addr, request) self.sel.register(sock, events, data=message) return sock, addr - def send_instruction(self, request): + def send_instruction(self): try: while True: events = self.sel.select(timeout=1) for key, mask in events: message = key.data try: - message.set_request(request) message.process_events(mask) except Exception: print( @@ -74,6 +69,11 @@ def send_instruction(self, request): finally: self.sel.close() + def send_request(self, action, value): + request = self.create_request(action, value) + self.start_connection(request) + self.send_instruction() + cnode = CameraNode() cnode.start_connection() @@ -90,20 +90,15 @@ def find_object(results, labels, sizes, distances, obj_name): obj_dist = next((dist for dist in distances if dist["name"] == obj_name), None) while score < 0.5: - request = cnode.create_request("left", "4") - cnode.send_instruction(request) - request = cnode.create_request("right", "4") - cnode.send_instruction(request) + cnode.send_request("left", "4") + cnode.send_request("right", "4") - request = cnode.create_request("stop", "all") - cnode.send_instruction(request) + cnode.send_request("stop", "all") while obj_dist > 100: - request = cnode.create_request("forward", "1") - cnode.send_instruction(request) + cnode.send_request("forward", "1") - request = cnode.create_request("stop", "all") - cnode.send_instruction(request) + cnode.send_request("stop", "all") tl_models = [ { From 5bd48e5251e64518836b8f9a6be2599036c6e9be Mon Sep 17 00:00:00 2001 From: zogoo Date: Tue, 2 Nov 2021 16:06:13 +0900 Subject: [PATCH 12/14] Auto-complete feature miss --- camera_node.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/camera_node.py b/camera_node.py index 2905aba..fff3489 100644 --- a/camera_node.py +++ b/camera_node.py @@ -1,6 +1,5 @@ #!/usr/bin/env python3 -from _typeshed import Self import sys import socket import selectors @@ -76,7 +75,6 @@ def send_request(self, action, value): cnode = CameraNode() -cnode.start_connection() def find_object(results, labels, sizes, distances, obj_name): From d6ef2d76bd17ab47be0d47c5765994a4d3df66b4 Mon Sep 17 00:00:00 2001 From: zogoo Date: Tue, 2 Nov 2021 16:29:34 +0900 Subject: [PATCH 13/14] We need to register object every time when we send request --- camera_node.py | 25 +++++++++++-------------- client_message.py | 2 +- 2 files changed, 12 insertions(+), 15 deletions(-) diff --git a/camera_node.py b/camera_node.py index fff3489..4a4cb30 100644 --- a/camera_node.py +++ b/camera_node.py @@ -35,18 +35,18 @@ def create_request(self, action, value, encode="utf-8"): content=dict(action=action, value=value), ) - def start_connection(self, request): - addr = ('127.0.0.1', 65432) - print("starting connection to", addr) - sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - sock.setblocking(False) - sock.connect_ex(addr) - events = selectors.EVENT_READ | selectors.EVENT_WRITE - message = Message(self.sel, sock, addr, request) - self.sel.register(sock, events, data=message) - return sock, addr + def start_connection(self): + self.addr = ('127.0.0.1', 65432) + print("starting connection to", self.addr) + self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self.sock.setblocking(False) + self.sock.connect_ex(self.addr) def send_instruction(self): + events = selectors.EVENT_READ | selectors.EVENT_WRITE + message = Message(self.sel, self.sock, self.addr) + self.sel.register(self.sock, events, data=message) + try: while True: events = self.sel.select(timeout=1) @@ -73,11 +73,8 @@ def send_request(self, action, value): self.start_connection(request) self.send_instruction() - -cnode = CameraNode() - - def find_object(results, labels, sizes, distances, obj_name): + cnode = CameraNode() score = 0 for obj in results: diff --git a/client_message.py b/client_message.py index a28d6f2..0cc2ed8 100644 --- a/client_message.py +++ b/client_message.py @@ -208,4 +208,4 @@ def process_response(self): ) self._process_response_binary_content() # Close when response has been processed - self.close() + # self.close() From c7faa5453eecb3babf87afb0cbc405aff01ec8d1 Mon Sep 17 00:00:00 2001 From: zogoo Date: Tue, 2 Nov 2021 16:34:11 +0900 Subject: [PATCH 14/14] Missing arguments --- camera_node.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/camera_node.py b/camera_node.py index 4a4cb30..c309f78 100644 --- a/camera_node.py +++ b/camera_node.py @@ -42,9 +42,9 @@ def start_connection(self): self.sock.setblocking(False) self.sock.connect_ex(self.addr) - def send_instruction(self): + def send_instruction(self, request): events = selectors.EVENT_READ | selectors.EVENT_WRITE - message = Message(self.sel, self.sock, self.addr) + message = Message(self.sel, self.sock, self.addr, request) self.sel.register(self.sock, events, data=message) try: @@ -70,8 +70,8 @@ def send_instruction(self): def send_request(self, action, value): request = self.create_request(action, value) - self.start_connection(request) - self.send_instruction() + self.start_connection() + self.send_instruction(request) def find_object(results, labels, sizes, distances, obj_name): cnode = CameraNode()